跳转到正文

graph-loop.ts ​

源文件: docs/worker-api/examples/runtime/graph-loop.ts · 下载原文件 · 示例使用说明

ts
import { createDitto, createInteractionWorker, createRuntimeServices, graph, loadRuntimeConfigFile, loop } from "@codesoul-co/ditto";

// Graphs describe dependencies and data. Worker placement is supplied at execution.
interface Input { value: string; iteration: number }
const makeGraph = (id: string, tool: string) => graph<Input>(id)
  .node("act", "INTERACTION.ACT.TOOL", [], input => ({
    call: { id: `${id}:${input.iteration}`, name: tool, arguments: { value: input.value } },
  }))
  .node("observe", "INTERACTION.OBSERVE", ["act"], (_input, { act }) => {
    if (act.status !== "success") throw new Error(act.error?.message ?? "Tool failed");
    return { result: act };
  });
const read = makeGraph("read", "read_text");
const count = makeGraph("count", "count_text");
const config = loadRuntimeConfigFile("ditto.yaml", {});
const runtime = createDitto({ config });
const reader = runtime.register(createInteractionWorker({ tools: [{
  name: "read_text", inputSchema: { type: "object" }, effects: ["read"], requiresApproval: false,
  validate(args) { if (typeof args.value !== "string") throw new Error("value must be a path"); },
  async execute(args, ctx) {
    return { status: "success", content: await ctx.services.sandbox.readText(args.value as string) };
  },
}] }), { id: "reader", services: createRuntimeServices({ config, sandbox: { tools: ["read_text"], read: true } }) });
const counter = runtime.register(createInteractionWorker({ tools: [{
  name: "count_text", inputSchema: { type: "object" }, effects: [], requiresApproval: false,
  validate(args) { if (typeof args.value !== "string") throw new Error("value must be text"); },
  async execute(args, ctx) {
    const content = String((args.value as string).length);
    await ctx.emit({ type: "counted", payload: { worker: ctx.worker.workerId, characters: Number(content) } });
    return { status: "success", content };
  },
}] }), { id: "counter", services: createRuntimeServices({ config, sandbox: { tools: ["count_text"] } }) });
runtime.subscribe("counted", event => { console.log("event", event.payload); });

interface State { value: string; iteration: number }
const workflow = loop({
  // read once, then run the same count graph twice.
  graph: (state: State) => state.iteration === 0 ? read : count,
  maxIterations: 3,
  bind: (state: State) => state,
  update: (state: State, output) => ({ value: String(output.observe.message.content), iteration: state.iteration + 1 }),
  done: (state: State) => state.iteration === 3,
});
try {
  console.log("result", await runtime.loop(workflow, { value: "README.md", iteration: 0 }, {
    concurrency: 1,
    workers: { read: { act: reader.address.workerId, observe: reader.address.workerId },
      count: { act: counter.address.workerId, observe: counter.address.workerId } },
  }));
  const failures = await runtime.drainEvents();
  if (failures.length) throw new AggregateError(failures.map(f => f.error), "Event consumer failed");
} finally { await runtime.close(); }

Ditto · @codesoul-co/ditto · Node.js 24+