多机器

A3S Code 把多智能体编排表达为代码中的一套语法(grammar),再把由此产生的步骤 放到你希望它运行的任何地方。这一划分沿着框架 / 宿主边界展开:

  • 框架拥有编排语法和可序列化的数据契约。它从不决定步骤在哪里运行。
  • 宿主拥有放置(placement)、传输和调度——哪个节点执行某个步骤、step 规格如何 送达、以及并发如何映射到集群。

两者之间唯一的接触点是一个 trait:AgentExecutor。

框架 / 宿主边界

框架的契约是两个可序列化类型:

  • AgentStepSpec——运行什么,与在哪里运行无关:task_id、agent、 description、prompt,以及可选的 max_steps、parent_session_id、 output_schema。
  • StepOutcome——运行一个 spec 的结果:task_id、session_id、agent、 output、success、可选的 structured,以及 source_anchors(成功的子工具 调用观察到的源码锚点)。

两者都能干净地序列化,因此宿主可以把一个 spec 送到另一个节点,并把 outcome 持久化 到 checkpoint。组合这些 spec 的 combinators 完全针对 AgentExecutor trait 编写, 从不观察步骤实际在哪里运行——所以同一套编排可以从单进程扩展到集群而无需改动。

AgentExecutor 扩展边界

AgentExecutor 是语法与宿主之间的边界:

Text
combinators (parallel / pipeline / resumable)
-> AgentExecutor::execute_step(spec, event_tx) -> StepOutcome
├─ in-box TaskExecutor: runs the step as a local child agent
└─ host executor: places the step on a remote node

execute_step 以 StepOutcome { success: false, .. } 报告失败,而不是硬错误, 因此某个分支失败时 fan-out 会继续;event_tx 存在时会接收该步骤的生命周期事件。

内置的 TaskExecutor 把每个步骤作为子 agent 在本地运行——在进程内、基于 Tokio—— 继承会话的 agent 注册表、LLM 客户端、工作区、MCP 工具和 subagent 跟踪器。诸如 集群运行时这样的宿主用自己的 executor 替换它,把步骤分布到集群;combinators 不受影响。

concurrency_hint() 是建议性的,不是本地硬上限。本地默认返回会话的 max_parallel_tasks;由调度器支撑的宿主可以返回其集群范围的目标值。因为它是 hint 而非 ceiling,编排得以扩展到单进程之外。

会话直接暴露内置 seam:

  • AgentSession::agent_executor() 返回一个由会话支撑的 AgentExecutor。
  • AgentSession::session_store() 返回会话的 store(在配置了的情况下),可恢复的 combinator 需要它来记录进度。

下面的 SDK 语法会替你调用 agent_executor();只有在实现或替换自定义 executor 时 你才需要直接使用它们。

并行:带屏障的扇出

execute_steps_parallel 把 specs 在 executor 上 fan-out 并等待全部完成(一个 屏障 / barrier)。结果保持输入顺序,panic 的分支会变成一个失败的 StepOutcome 而不会拖垮整批,并且并发受 executor 的 concurrency hint 限制。

在 Node.js 中,session.parallel(specs) 解析为 outcome 数组。传入可选的 budgetTokens 参数会限制子 LLM 花费,并把结果改为 { outcomes, budget },因此 声明的返回类型是联合类型;迭代前需要先收窄。Go 的 session.Parallel(ctx, specs, budgetTokens) 返回携带 Outcomes 与可选 Budget 的 *ParallelResult。

TypeScript
import type { StepOutcomeObject } from '@a3s-lab/code';
const outcomes = (await session.parallel([
{
taskId: 'a',
agent: 'explore',
description: 'survey',
prompt: 'Map the auth module',
},
{
taskId: 'b',
agent: 'review',
description: 'audit',
prompt: 'Review error handling',
},
])) as StepOutcomeObject[];
for (const o of outcomes) {
console.log(o.taskId, o.success, o.output);
}
Python
outcomes = session.parallel([
{"task_id": "a", "agent": "explore", "description": "survey", "prompt": "Map the auth module"},
{"task_id": "b", "agent": "review", "description": "audit", "prompt": "Review error handling"},
])
for o in outcomes:
print(o["task_id"], o["success"], o["output"])

流水线:逐项串联,阶段之间无屏障

execute_pipeline 让每个 item 独立地流经一条由 PipelineStage 组成的链路。 stage 之间没有屏障——item A 可以处于 stage 3,而 item B 仍在 stage 1——因此 墙钟时间取决于最慢的单条链路,而不是每个 stage 最慢步骤之和。

一个 stage 是 (ctx) => spec | null 回调,其中 ctx 携带上一步的 outcome 和原始 item。返回一个 spec 以运行下一步,或返回 null 提前停止该 item 的链路;当某步失败 时链路也会停止(后续 stage 只会建立在失败结果之上)。每次调用按输入顺序为每个 item 返回一项——该 item 的最后一个 outcome;若其第一个 stage 就返回 null,则为 null。

Node 与 Python 桥接按链路 fail closed。在 Node 中,抛出异常、返回 Promise、 返回格式错误的 spec,或超过 timeoutMs(第三个参数,默认 30 秒)的 stage 都被当作 null。在 Python 中,raise 的 stage 被当作 None。stage 回调必须是同步的。Go 的 session.Pipeline(ctx, items, stages) 不同:返回 error 的 stage 会使整个调用在返回 部分 outcome 的同时返回该 error。

TypeScript
const outcomes = await session.pipeline(
['src/auth', 'src/api'],
[
(ctx) => ({
taskId: 's1',
agent: 'explore',
description: 'survey',
prompt: `Survey ${ctx.item}`,
}),
(ctx) =>
ctx.previous?.success
? {
taskId: 's2',
agent: 'review',
description: 'review',
prompt: `Review: ${ctx.previous.output}`,
}
: null,
],
);
Python
def survey(ctx):
return {"task_id": "s1", "agent": "explore", "description": "survey",
"prompt": f"Survey {ctx['item']}"}
def review(ctx):
prev = ctx["previous"]
if prev and prev["success"]:
return {"task_id": "s2", "agent": "review", "description": "review",
"prompt": f"Review: {prev['output']}"}
return None
outcomes = session.pipeline(["src/auth", "src/api"], [survey, review])

可恢复执行:跨节点恢复

execute_steps_parallel_resumable 是 parallel 加上一份日志。在每个步骤边界,它把 一个 WorkflowCheckpoint 以 workflowId 为键写入 SessionStore。恢复时它跳过已 完成的步骤,重新派发其余步骤。它只记录成功的步骤,因此失败的步骤会在恢复时重试 ——它的效果尚未完成。所有步骤都成功时,checkpoint 会被删除。

需要一个 SessionStore。未配置时,Node 以 parallelResumable requires a sessionStore on the session 拒绝,Python 抛出 RuntimeError: parallel_resumable requires a session_store on the session。

TypeScript
import { Agent, FileSessionStore } from '@a3s-lab/code';
const session = agent.session('/repo', {
sessionStore: new FileSessionStore('./.a3s/sessions'),
});
const outcomes = await session.parallelResumable(
[
{
taskId: 'a',
agent: 'explore',
description: 'survey',
prompt: 'Map the auth module',
},
{
taskId: 'b',
agent: 'review',
description: 'audit',
prompt: 'Review error handling',
},
],
'release-audit',
);
Python
from a3s_code import Agent, FileSessionStore, SessionOptions
opts = SessionOptions()
opts.session_store = FileSessionStore("./.a3s/sessions")
session = agent.session("/repo", opts)
outcomes = session.parallel_resumable([
{"task_id": "a", "agent": "explore", "description": "survey", "prompt": "Map the auth module"},
{"task_id": "b", "agent": "review", "description": "audit", "prompt": "Review error handling"},
], "release-audit")

Go 以 session.Parallel(ctx, specs, budgetTokens)、 session.Pipeline(ctx, items, stages) 与 session.ParallelResumable(ctx, specs, workflowID) 暴露同一套语法:

Go
outcomes, err := session.ParallelResumable(ctx, []code.AgentStepSpec{
{TaskID: "a", Agent: "explore", Description: "survey", Prompt: "Map the auth module"},
{TaskID: "b", Agent: "review", Description: "audit", Prompt: "Review error handling"},
}, "release-audit")
if err != nil {
return err
}

因为 checkpoint 是可序列化的,且 executor 是一个参数,宿主可以通过传入另一个节点的 executor,在不同的节点上恢复被中断的工作流——框架迁移运行什么,宿主提供 在哪里运行。

模式约束的步骤输出

携带 output_schema(Node 中为 outputSchema)的 spec 必须返回一个符合该 JSON Schema 的值;经校验的对象落入 StepOutcome.structured。executor 复用结构化输出的 coercion 与 repair 机制。coercion 失败会把该步骤降级为不成功,并在其输出后追加 [structured output failed: ...],因此调用方永远不会把未经校验的文本当作所承诺的 对象。Python 的 pipeline 不接受逐 stage 的 output_schema;需要模式校验的步骤请使用 parallel。

TypeScript
import type { StepOutcomeObject } from '@a3s-lab/code';
const [finding] = (await session.parallel([
{
taskId: 'classify',
agent: 'review',
description: 'classify',
prompt: 'Classify this defect',
outputSchema: {
type: 'object',
properties: { severity: { type: 'string' }, summary: { type: 'string' } },
required: ['severity', 'summary'],
},
},
])) as StepOutcomeObject[];
if (finding.success && finding.structured) {
console.log(finding.structured.severity);
}

跨机器放置编排步骤

lane queue 仍然是一种有效的传输——但它现在是AgentExecutor seam 背后的一个选项, 而不是唯一的集成点。要分布工作,宿主针对任何适合其平台的传输(HTTP、消息队列、 任务系统、lane queue、内部 RPC)实现 AgentExecutor::execute_step,把 concurrency_hint() 接到其集群目标,再把该 executor 交给 combinators。语法—— parallel、pipeline、resumable——保持不变;移动的只有放置。

结果是:coordinator 会话拥有对话、最终合成和发布决策,而它的编排步骤在宿主放置的 任何地方执行,跨节点恢复由可序列化的 checkpoint 承载。

参见 Orchestration 深入了解 combinator 语法,以及 Persistence 了解 checkpoint store。