多智能体系统的核心不是同时调用几个模型,而是明确谁可以执行、何时停止,以及失败后哪些节点不应再运行。本文用一个不依赖智能体框架的 TypeScript 任务图,演示状态契约、依赖调度、并行汇合、预算限制和图级测试。

先把协作流程表示成图

把多个角色放进群聊,让它们轮流发言,是一种容易实现但难以约束的编排方式。只要提示词没有稳定地产生“结束”信号,智能体就可能继续讨论;上下文变化后,同一任务还可能被重复执行。

任务图采用更明确的模型:节点负责执行,边表示依赖。只有全部依赖成功的节点才能进入运行态,失败节点的下游会被跳过,所有节点进入终态后整张图结束。

例如,一个资料分析流程可以写成:

节点依赖职责
question规范化用户问题
searchquestion检索资料
outlinequestion生成分析提纲
answersearchoutline汇合结果并生成答案

searchoutline 可以并行,但 answer 必须等待两者成功。这种汇合条件由调度器判断,不交给模型自行协商。

节点状态至少需要覆盖以下集合:

type NodeStatus =
  | 'pending'
  | 'running'
  | 'succeeded'
  | 'failed'
  | 'skipped';

其中 succeededfailedskipped 是终态。状态转换也应受到限制:正常路径是 pending -> running -> succeeded | failed;依赖失败或预算耗尽时,则是 pending -> skipped。把契约固定下来后,调度器才能可靠地判断终止,而不是猜测模型是否已经完成任务。

实现可终止的任务图

下面的实现只使用 Node.js 标准 API。它在启动前检查未知依赖和环,按批次并行执行就绪节点,并提供节点次数、总时长、单节点超时和最大并发数四项限制。

import assert from 'node:assert/strict';

type NodeStatus = 'pending' | 'running' | 'succeeded' | 'failed' | 'skipped';
type StopReason = 'settled' | 'node-budget' | 'time-budget';

type TaskContext = {
  signal: AbortSignal;
  output<T>(id: string): T;
};

type TaskNode = {
  id: string;
  dependencies: string[];
  run(context: TaskContext): Promise<unknown>;
};

type NodeState = {
  status: NodeStatus;
  attempts: number;
  output?: unknown;
  error?: string;
};

type TraceEvent = {
  step: number;
  atMs: number;
  event: string;
  nodeId?: string;
  detail?: string;
};

type RunOptions = {
  maxNodeRuns: number;
  maxDurationMs: number;
  nodeTimeoutMs: number;
  maxParallel: number;
};

type GraphResult = {
  reason: StopReason;
  states: Record<string, NodeState>;
  trace: TraceEvent[];
};

class TaskGraph {
  private readonly nodes = new Map<string, TaskNode>();

  constructor(nodes: TaskNode[]) {
    for (const node of nodes) {
      if (this.nodes.has(node.id)) throw new Error(`duplicate node: ${node.id}`);
      this.nodes.set(node.id, node);
    }
    this.validate();
  }

  async run(options: RunOptions): Promise<GraphResult> {
    if (options.maxNodeRuns < 1 || options.maxDurationMs < 1 ||
        options.nodeTimeoutMs < 1 || options.maxParallel < 1) {
      throw new Error('all run options must be positive');
    }

    const startedAt = Date.now();
    const states = new Map<string, NodeState>();
    const trace: TraceEvent[] = [];
    let step = 0;
    let runs = 0;
    let reason: StopReason = 'settled';

    const record = (event: string, nodeId?: string, detail?: string) => {
      trace.push({ step: step++, atMs: Date.now() - startedAt, event, nodeId, detail });
    };

    for (const id of this.nodes.keys()) {
      states.set(id, { status: 'pending', attempts: 0 });
    }
    record('graph_started');

    const execute = async (node: TaskNode, timeoutMs: number) => {
      const state = states.get(node.id)!;
      state.status = 'running';
      state.attempts += 1;
      runs += 1;
      record('node_started', node.id);

      const controller = new AbortController();
      let timer: ReturnType<typeof setTimeout> | undefined;
      const timeout = new Promise<never>((_, reject) => {
        timer = setTimeout(() => {
          controller.abort();
          reject(new Error(`timeout after ${timeoutMs}ms`));
        }, timeoutMs);
      });

      try {
        const context: TaskContext = {
          signal: controller.signal,
          output: <T>(id: string) => {
            const dependency = states.get(id);
            if (!dependency || dependency.status !== 'succeeded') {
              throw new Error(`output is unavailable: ${id}`);
            }
            return dependency.output as T;
          }
        };
        state.output = await Promise.race([node.run(context), timeout]);
        state.status = 'succeeded';
        record('node_succeeded', node.id);
      } catch (error) {
        state.status = 'failed';
        state.error = error instanceof Error ? error.message : String(error);
        record('node_failed', node.id, state.error);
      } finally {
        if (timer) clearTimeout(timer);
      }
    };

    while (true) {
      const elapsed = Date.now() - startedAt;
      if (elapsed >= options.maxDurationMs) {
        reason = 'time-budget';
        break;
      }

      for (const node of this.nodes.values()) {
        const state = states.get(node.id)!;
        if (state.status !== 'pending') continue;
        const blocked = node.dependencies.some((id) => {
          const status = states.get(id)!.status;
          return status === 'failed' || status === 'skipped';
        });
        if (blocked) {
          state.status = 'skipped';
          state.error = 'dependency did not succeed';
          record('node_skipped', node.id, state.error);
        }
      }

      const pending = [...this.nodes.values()].filter(
        (node) => states.get(node.id)!.status === 'pending'
      );
      if (pending.length === 0) break;

      const ready = pending.filter((node) =>
        node.dependencies.every((id) => states.get(id)!.status === 'succeeded')
      );
      if (ready.length === 0) throw new Error('graph cannot make progress');

      const remainingRuns = options.maxNodeRuns - runs;
      if (remainingRuns <= 0) {
        reason = 'node-budget';
        break;
      }

      const batch = ready.slice(0, Math.min(options.maxParallel, remainingRuns));
      const remainingMs = options.maxDurationMs - (Date.now() - startedAt);
      const timeoutMs = Math.max(1, Math.min(options.nodeTimeoutMs, remainingMs));
      await Promise.all(batch.map((node) => execute(node, timeoutMs)));
    }

    if (reason !== 'settled') {
      for (const [id, state] of states) {
        if (state.status === 'pending') {
          state.status = 'skipped';
          state.error = `graph stopped: ${reason}`;
          record('node_skipped', id, state.error);
        }
      }
    }

    record('graph_finished', undefined, reason);
    return {
      reason,
      states: Object.fromEntries(states),
      trace
    };
  }

  private validate(): void {
    const incoming = new Map<string, number>();
    const outgoing = new Map<string, string[]>();
    for (const node of this.nodes.values()) {
      incoming.set(node.id, node.dependencies.length);
      for (const dependency of node.dependencies) {
        if (!this.nodes.has(dependency)) {
          throw new Error(`unknown dependency: ${dependency}`);
        }
        outgoing.set(dependency, [...(outgoing.get(dependency) ?? []), node.id]);
      }
    }

    const queue = [...incoming].filter(([, count]) => count === 0).map(([id]) => id);
    let visited = 0;
    while (queue.length > 0) {
      const id = queue.shift()!;
      visited += 1;
      for (const next of outgoing.get(id) ?? []) {
        const count = incoming.get(next)! - 1;
        incoming.set(next, count);
        if (count === 0) queue.push(next);
      }
    }
    if (visited !== this.nodes.size) throw new Error('graph contains a cycle');
  }
}

这里采用“批次屏障”:同一批就绪节点并行运行,全部结束后再调度下一批。它不追求极致吞吐,但状态顺序更容易理解和测试。对多数包含外部模型调用的流程,这通常是合适的起点;需要更高利用率时,可以改成节点完成后立即触发下游,但轨迹并发控制也会更复杂。

用确定性节点验证并行汇合

不要一开始就用真实模型测试编排。模型输出、网络延迟和限流会同时引入变量,很难判断失败来自调度器还是模型。可以先追加以下确定性测试:

async function main() {
  let active = 0;
  let peak = 0;

  const parallelNode = (id: string): TaskNode => ({
    id,
    dependencies: [],
    async run() {
      active += 1;
      peak = Math.max(peak, active);
      await Promise.resolve();
      active -= 1;
      return id.toUpperCase();
    }
  });

  const graph = new TaskGraph([
    parallelNode('search'),
    parallelNode('outline'),
    {
      id: 'answer',
      dependencies: ['search', 'outline'],
      async run(context) {
        return `${context.output<string>('search')}+${context.output<string>('outline')}`;
      }
    }
  ]);

  const result = await graph.run({
    maxNodeRuns: 10,
    maxDurationMs: 1_000,
    nodeTimeoutMs: 200,
    maxParallel: 2
  });

  assert.equal(result.reason, 'settled');
  assert.equal(result.states.answer.status, 'succeeded');
  assert.equal(result.states.answer.output, 'SEARCH+OUTLINE');
  assert.equal(peak, 2);

  const answerStart = result.trace.findIndex(
    (event) => event.event === 'node_started' && event.nodeId === 'answer'
  );
  const dependencyEnds = result.trace
    .map((event, index) => ({ event, index }))
    .filter(({ event }) =>
      event.event === 'node_succeeded' &&
      (event.nodeId === 'search' || event.nodeId === 'outline')
    );
  assert.equal(dependencyEnds.length, 2);
  assert.ok(dependencyEnds.every(({ index }) => index < answerStart));

  console.log(JSON.stringify(result, null, 2));
}

await main();

将两段代码放在同一个 task-graph.ts 文件中,再准备最小化的运行配置:

{
  "type": "module",
  "scripts": {
    "start": "tsx task-graph.ts"
  },
  "devDependencies": {
    "tsx": "^4.0.0",
    "typescript": "^5.0.0"
  }
}

执行 npm installnpm start 即可运行。测试不仅检查最终答案,还验证两个前置节点确实并行,以及汇合节点一定在依赖成功之后启动。这比只比较最终文本更能定位协作流程中的错误。

轨迹与预算如何阻止失控

图级轨迹记录的是调度事实,而不是模型自述。示例中的每条事件都有递增的 step、相对时间、事件类型、节点 ID 和可选详情。由此可以回答:节点是否重复启动、失败发生在哪一步、下游为何跳过,以及图因何停止。

建议至少监控以下不变量:

不变量发现的问题
每个节点默认只启动一次重复执行或幂等性缺失
节点启动时依赖均成功越过依赖读取不完整结果
running 最终进入终态超时失效或 Promise 长期悬挂
图结束后不存在 pending终止清理不完整
实际启动次数不超过预算循环、重试或动态扩图失控

maxNodeRuns 限制工作总量,maxDurationMs 限制图的墙钟时间,nodeTimeoutMs 限制单次执行,maxParallel 控制外部服务压力。四者解决的问题不同,不应只保留一个总超时。

需要注意,AbortSignal 只是协作式取消。节点内部调用 fetch 时应把 signal 传进去;如果第三方 SDK 不接受取消信号,Promise 超时只能让调度器停止等待,不能保证远端计算立即停止。因此,有副作用的节点还需要幂等键、请求 ID 或持久化去重记录。

从示例走向生产环境

这个实现刻意保持了较小的边界。接入模型时,模型调用应位于节点的 run 函数内部,调度器不需要知道 OpenAI、Anthropic 或本地模型的具体 API。这样可以用假节点测试图,再对单个模型节点做集成测试。

生产环境通常还需要补充持久化、重试策略和恢复机制。重试应由节点策略明确控制,并计入执行预算;不能把失败节点重新改成 pending 后无限尝试。轨迹可写入数据库或可观测平台,但应过滤提示词中的密钥和个人信息。

动态生成子任务时也要克制。允许模型无限添加节点,会重新引入不可终止问题。更稳妥的做法是限制新增节点数量、允许的节点类型和图深度,并在扩图后再次执行依赖与环检查。

此外,示例的内存状态适合单进程流程。如果任务跨分钟甚至跨小时,应持久化节点状态和输出,并通过租约或原子状态更新避免多个工作进程重复领取同一节点。此时状态契约仍然适用,只是 Map 会被数据库中的条件更新替代。

总结

多智能体工程首先是调度与状态管理问题,其次才是提示词问题。一个可控的任务图应具备以下能力:

  • 用有限状态集合约束节点生命周期,并把失败下游显式标记为 skipped
  • 在执行前检查未知依赖和环,运行时只调度依赖全部成功的节点。
  • 让独立节点并行,让汇合节点等待全部输入完成。
  • 同时限制节点次数、图时长、节点超时和并发量,并记录明确的终止原因。
  • 保存图级事件轨迹,用确定性节点验证顺序、并行度、失败传播和预算行为。

最终答案只能说明某次运行看起来可用;状态与轨迹才能说明协作过程是否按设计发生。先把任务图测清楚,再替换成真实模型节点,通常更容易得到可重复、可诊断的多智能体系统。