跳转到正文

fan-out-fan-in.ts ​

源文件: examples/control-flow/parallel/fan-out-fan-in.ts · 下载原文件 · 示例使用说明

ts
import { buildParallelGraph, type Branch } from "./concurrency-limit.ts";
import { runCli, isMain } from "./cli.ts";
import { accepted, deliveryInput, savedOrder, toolData, validateBatch, validateOptions, type BatchInput, type ExecutionOptions, type Runner } from "./shared.ts";

export function buildSummaryGraph(branches: readonly Branch[], joinId = "summary") {
  const dependencies = branches.map(branch => `save:${branch.taskId}`);
  return buildParallelGraph(branches, "parallel-summary")
    .node(`join:${joinId}`, "INTERACTION.ACT.TOOL", dependencies, (input, outputs) => {
      for (const id of dependencies) savedOrder(outputs[id]);
      return { call: { id: `${input.id}:report`, name: "build_order_report", arguments: {
        batchId: input.id, successIds: branches.map(branch => branch.sourceId), failures: [],
      } } };
    })
    .node("delivered", "INTERACTION.OUTPUT", [`join:${joinId}`], (input, outputs) => deliveryInput(input.id, outputs[`join:${joinId}`]));
}
export async function runSummary(runtime: Runner, input: BatchInput, options: ExecutionOptions = {}) {
  validateBatch(input); validateOptions(options);
  const branches = input.sources.map(source => ({ taskId: source.id, sourceId: source.id }));
  const output = await runtime.run(buildSummaryGraph(branches), input, { concurrency: options.concurrency ?? 3, ...(options.signal ? { signal: options.signal } : {}) });
  accepted(output.delivered);
  return { report: toolData(output["join:summary"]), receipt: output.delivered, output };
}
if (isMain(import.meta.url)) await runCli((runtime, input, options) => runSummary(runtime, input, options));

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