多智能体系统的核心不是同时调用几个模型,而是明确谁可以执行、何时停止,以及失败后哪些节点不应再运行。本文用一个不依赖智能体框架的 TypeScript 任务图,演示状态契约、依赖调度、并行汇合、预算限制和图级测试。
先把协作流程表示成图
把多个角色放进群聊,让它们轮流发言,是一种容易实现但难以约束的编排方式。只要提示词没有稳定地产生“结束”信号,智能体就可能继续讨论;上下文变化后,同一任务还可能被重复执行。
任务图采用更明确的模型:节点负责执行,边表示依赖。只有全部依赖成功的节点才能进入运行态,失败节点的下游会被跳过,所有节点进入终态后整张图结束。
例如,一个资料分析流程可以写成:
| 节点 | 依赖 | 职责 |
|---|---|---|
question | 无 | 规范化用户问题 |
search | question | 检索资料 |
outline | question | 生成分析提纲 |
answer | search、outline | 汇合结果并生成答案 |
search 与 outline 可以并行,但 answer 必须等待两者成功。这种汇合条件由调度器判断,不交给模型自行协商。
节点状态至少需要覆盖以下集合:
type NodeStatus =
| 'pending'
| 'running'
| 'succeeded'
| 'failed'
| 'skipped';
其中 succeeded、failed 和 skipped 是终态。状态转换也应受到限制:正常路径是 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 install 和 npm start 即可运行。测试不仅检查最终答案,还验证两个前置节点确实并行,以及汇合节点一定在依赖成功之后启动。这比只比较最终文本更能定位协作流程中的错误。
轨迹与预算如何阻止失控
图级轨迹记录的是调度事实,而不是模型自述。示例中的每条事件都有递增的 step、相对时间、事件类型、节点 ID 和可选详情。由此可以回答:节点是否重复启动、失败发生在哪一步、下游为何跳过,以及图因何停止。
建议至少监控以下不变量:
| 不变量 | 发现的问题 |
|---|---|
| 每个节点默认只启动一次 | 重复执行或幂等性缺失 |
| 节点启动时依赖均成功 | 越过依赖读取不完整结果 |
running 最终进入终态 | 超时失效或 Promise 长期悬挂 |
图结束后不存在 pending | 终止清理不完整 |
| 实际启动次数不超过预算 | 循环、重试或动态扩图失控 |
maxNodeRuns 限制工作总量,maxDurationMs 限制图的墙钟时间,nodeTimeoutMs 限制单次执行,maxParallel 控制外部服务压力。四者解决的问题不同,不应只保留一个总超时。
需要注意,AbortSignal 只是协作式取消。节点内部调用 fetch 时应把 signal 传进去;如果第三方 SDK 不接受取消信号,Promise 超时只能让调度器停止等待,不能保证远端计算立即停止。因此,有副作用的节点还需要幂等键、请求 ID 或持久化去重记录。
从示例走向生产环境
这个实现刻意保持了较小的边界。接入模型时,模型调用应位于节点的 run 函数内部,调度器不需要知道 OpenAI、Anthropic 或本地模型的具体 API。这样可以用假节点测试图,再对单个模型节点做集成测试。
生产环境通常还需要补充持久化、重试策略和恢复机制。重试应由节点策略明确控制,并计入执行预算;不能把失败节点重新改成 pending 后无限尝试。轨迹可写入数据库或可观测平台,但应过滤提示词中的密钥和个人信息。
动态生成子任务时也要克制。允许模型无限添加节点,会重新引入不可终止问题。更稳妥的做法是限制新增节点数量、允许的节点类型和图深度,并在扩图后再次执行依赖与环检查。
此外,示例的内存状态适合单进程流程。如果任务跨分钟甚至跨小时,应持久化节点状态和输出,并通过租约或原子状态更新避免多个工作进程重复领取同一节点。此时状态契约仍然适用,只是 Map 会被数据库中的条件更新替代。
总结
多智能体工程首先是调度与状态管理问题,其次才是提示词问题。一个可控的任务图应具备以下能力:
- 用有限状态集合约束节点生命周期,并把失败下游显式标记为
skipped。 - 在执行前检查未知依赖和环,运行时只调度依赖全部成功的节点。
- 让独立节点并行,让汇合节点等待全部输入完成。
- 同时限制节点次数、图时长、节点超时和并发量,并记录明确的终止原因。
- 保存图级事件轨迹,用确定性节点验证顺序、并行度、失败传播和预算行为。
最终答案只能说明某次运行看起来可用;状态与轨迹才能说明协作过程是否按设计发生。先把任务图测清楚,再替换成真实模型节点,通常更容易得到可重复、可诊断的多智能体系统。
评论