编排

编排是高级宿主编排面——模型驱动委派的可编程姊妹能力。日常编码 Agent 扇出请优先使用统一的 task。在 Tasks 与 Teams 中,由 LLM 在运行时决定是否调用含一个或多个 tasks[] 项的 task——扇出的形状取决于模型的选择。编排把这个决策搬进你的代码:你 用一套语法表达扇出、流水线与验证面板,因此形状是可复现、可测试、受预算约束 且可恢复的——与模型的临场选择无关。

当工作的结构是宿主提前已知的(并行跑这三个 reviewer;让每个候选项依次经过 explore → verify → review;崩溃后恢复这一批),就用编排。当你希望由模型决定是否 以及如何委派时,就用 Tasks/Teams。

框架 / 宿主边界(接缝)

本层的一切都围绕单一接缝 AgentExecutor 编写:"运行这个 step,把结果给我。"这 条接缝把职责切得很干净:

  • 框架拥有语法——有哪些 step、如何组合、并发提示,以及可序列化契约 AgentStepSpec / StepOutcome。
  • 宿主拥有放置——传输、调度,以及 step 实际在哪里运行。

内置的默认 executor(TaskExecutor)在本地、进程内、基于 tokio 运行每个 step。 宿主可以替换为自己的 AgentExecutor,把 step 放置到集群各处;组合子从不观察 step 在哪里运行,因此同一套编排无需改动即可从单进程扩展到集群。

concurrency_hint() 是建议性的,而非硬性本地上限——正是它让编排能扩展到单进 程之外(由调度器支撑的宿主会返回其集群级目标,而不是本地上限)。

SDK 已为你接好这一切:AgentSession::agent_executor() 返回由 session 支撑的 executor(它把每个 step 作为子 agent 在本节点运行,继承 session 的 agent registry、 LLM client、workspace 与 MCP 工具),session_store() 返回 session 的 store。 parallel / pipeline / parallelResumable 方法会替你调用它们。

步骤契约

一个 step 由 AgentStepSpec 描述,并解析为 StepOutcome。两者都刻意保持可序列 化:宿主可以把 spec 发送到另一个节点,可恢复组合子也会把 outcome 持久化进 checkpoint。

AgentStepSpec 字段:

  • task_id——step 的稳定 id(由你指定);会流入生命周期事件与 checkpoint。
  • agent——要运行的 agent 的 registry key(例如 explore、review)。
  • description——用于展示/追踪的简短人类标签。
  • prompt——交给子 agent 的指令。
  • max_steps(可选)——每个 step 的 tool-round 上限。
  • parent_session_id(可选)——用于事件关联的父 session id。
  • output_schema(可选)——设置后,步骤必须返回符合此 JSON Schema 的值(见 强制模式约束的步骤输出)。

StepOutcome 字段:

  • task_id——产生该结果的 step id。
  • session_id——子运行的 session id(失败的 step 仍可寻址)。
  • agent——运行的 agent。
  • output——step 的文本输出。
  • success——失败或 panic 的 step 为 false(绝不会丢弃兄弟 step)。
  • structured(可选)——经模式校验的对象,仅当任务规格携带 output_schema 时存在。
  • source_anchors——成功的子研究工具调用所观察到的来源(tool 加上规范化的 URL 或工作区相对路径)。

key 的大小写风格因 SDK 而异:

概念Node(camelCase)Python(snake_case)Go(struct 字段)
step idtaskIdtask_idTaskID
tool-round 上限maxStepsmax_stepsMaxSteps
父 sessionparentSessionIdparent_session_idParentSessionID
强制模式约束outputSchemaoutput_schemaOutputSchema
子 sessionsessionIdsession_idSessionID
来源锚点sourceAnchorssource_anchorsSourceAnchors

agent、description、prompt、output、success 与 structured 在 Node 和 Python 中拼写相同;Go 使用导出字段 Agent、Description、Prompt、Output、 Success 与 Structured。

parallel:带屏障的扇出 [#parallel--barrier-fan-out]

session.parallel(specs) 把每个任务规格作为扇出分支运行,并按输入顺序解析出每 个 spec 对应的一个 StepOutcome。它映射到核心组合子 execute_steps_parallel。它 是一个屏障:在返回前会等待每个 step。

每个分支彼此隔离——失败或 panic 的 step 会变成 success: false,绝不会丢弃兄弟 step。并发受 executor 的并发提示约束(默认即 session 配置的并行度)。 在 TypeScript 中返回类型是 StepOutcomeObject[] | WorkflowParallelResult;不传预算时 请转换为 StepOutcomeObject[]。

TypeScript
const outcomes = await session.parallel([
{
taskId: 'explore',
agent: 'explore',
description: 'Risky changes',
prompt: 'Find risky changed files in this diff.',
},
{
taskId: 'verify',
agent: 'verification',
description: 'Test gaps',
prompt: 'Identify missing or weak verification.',
},
{
taskId: 'review',
agent: 'review',
description: 'Correctness',
prompt: 'Review the diff for correctness risks.',
},
]);
for (const outcome of outcomes) {
if (outcome.success) {
console.log(outcome.taskId, outcome.output);
} else {
console.warn('failed:', outcome.taskId, outcome.output);
}
}
Python
outcomes = session.parallel([
{"task_id": "explore", "agent": "explore", "description": "Risky changes",
"prompt": "Find risky changed files in this diff."},
{"task_id": "verify", "agent": "verification", "description": "Test gaps",
"prompt": "Identify missing or weak verification."},
{"task_id": "review", "agent": "review", "description": "Correctness",
"prompt": "Review the diff for correctness risks."},
])
for outcome in outcomes:
if outcome["success"]:
print(outcome["task_id"], outcome["output"])
else:
print("failed:", outcome["task_id"], outcome["output"])
Go
result, err := session.Parallel(ctx, []code.AgentStepSpec{
{TaskID: "explore", Agent: "explore", Description: "Risky changes",
Prompt: "Find risky changed files in this diff."},
{TaskID: "verify", Agent: "verification", Description: "Test gaps",
Prompt: "Identify missing or weak verification."},
{TaskID: "review", Agent: "review", Description: "Correctness",
Prompt: "Review the diff for correctness risks."},
}, nil)
if err != nil {
return err
}
for _, outcome := range result.Outcomes {
fmt.Println(outcome.TaskID, outcome.Success, outcome.Output)
}

pipeline:阶段之间无屏障

session.pipeline(items, stages) 让每个 item 独立地流经一连串阶段——阶段之间 没有屏障,因此 item A 可以处于 stage 3,而 item B 还在 stage 1。墙钟时间是最慢的 单条链,而不是逐阶段屏障会带来的"每阶段最慢之和"。

阶段是 spec 构造器,而非 spec:每个阶段接收上一个 outcome 和原始 item,返回要运行 的下一个 step,或返回 null / None 提前停止该 item 的链。失败的 step 同样会停止 链(后续阶段只会基于失败结果继续)。回调形状:

  • Node:(ctx) => spec | null,其中 ctx = { previous: StepOutcome | null, item }
  • Python:stage(ctx) -> spec | None,其中 ctx = {"previous": <dict|None>, "item": <item>}
  • Go:code.PipelineStage,即 func(ctx, code.PipelineContext) (*code.AgentStepSpec, error), 其中 PipelineContext{Previous *StepOutcome, Item any};返回 nil, nil 停止该链

阶段可以基于上一个 outcome 分支——例如"验证 review 阶段产出的发现"。

每次调用按输入顺序为每个 item 返回一项:该 item 的最后一个 outcome;若其第一个 阶段没有返回 spec,则为 null / None / nil。

约束:

  • Node: 阶段必须是同步函数。抛异常、返回 Promise 或格式错误的 spec、或超过 timeoutMs(第 3 个参数,默认 30000)的阶段会 fail closed——被当作 null 处理,仅停止该条链。阶段 spec 上的 outputSchema 会被忽略。
  • Python: 抛异常的阶段 callable 会被捕获并当作 None(仅停止该条链)。阶段 spec 上的 output_schema 会被忽略。
  • Go: 每个阶段的 step 都通过 WorkflowStep 运行,因此 OutputSchema 生效。 返回 error 的阶段(或 step 请求)会停止该条链,Pipeline 会连同已收集的 outcomes 一起返回第一个 error。

在 Node 或 Python 中需要模式校验的步骤请用 parallel。

TypeScript
const outcomes = await session.pipeline(
['src/auth.ts', 'src/payments.ts'],
[
(ctx) => ({
taskId: `explore-${ctx.item}`,
agent: 'explore',
description: 'Inspect file',
prompt: `Summarize the responsibilities and risks of ${ctx.item}.`,
}),
(ctx) => {
if (!ctx.previous) return null;
return {
taskId: `review-${ctx.item}`,
agent: 'review',
description: 'Review of prior finding',
prompt: `Review this summary for correctness risks:\n${ctx.previous.output}`,
};
},
],
);
Python
def explore_stage(ctx):
item = ctx["item"]
return {
"task_id": f"explore-{item}",
"agent": "explore",
"description": "Inspect file",
"prompt": f"Summarize the responsibilities and risks of {item}.",
}
def review_stage(ctx):
prev = ctx["previous"]
if prev is None:
return None
item = ctx["item"]
return {
"task_id": f"review-{item}",
"agent": "review",
"description": "Review of prior finding",
"prompt": f"Review this summary for correctness risks:\n{prev['output']}",
}
outcomes = session.pipeline(
["src/auth.ts", "src/payments.ts"],
[explore_stage, review_stage],
)
Go
outcomes, err := session.Pipeline(ctx,
[]any{"src/auth.ts", "src/payments.ts"},
[]code.PipelineStage{
func(_ context.Context, c code.PipelineContext) (*code.AgentStepSpec, error) {
return &code.AgentStepSpec{
TaskID: fmt.Sprintf("explore-%v", c.Item),
Agent: "explore",
Description: "Inspect file",
Prompt: fmt.Sprintf("Summarize the responsibilities and risks of %v.", c.Item),
}, nil
},
func(_ context.Context, c code.PipelineContext) (*code.AgentStepSpec, error) {
if c.Previous == nil {
return nil, nil
}
return &code.AgentStepSpec{
TaskID: fmt.Sprintf("review-%v", c.Item),
Agent: "review",
Description: "Review of prior finding",
Prompt: "Review this summary for correctness risks:\n" + c.Previous.Output,
}, nil
},
})

可恢复 / 可迁移工作流

session.parallelResumable(specs, workflowId)(Node)/ session.parallel_resumable(specs, workflow_id)(Python)/ session.ParallelResumable(ctx, specs, workflowID)(Go)是 parallel 加上一份 日志。它映射到 execute_steps_parallel_resumable。

它在每个 step 边界把一份 WorkflowCheckpoint 写入 session store。恢复时它跳过已完 成的 step(复用其缓存的 outcome),只重新派发其余 step。它只记录成功的 step—— 失败的 step 不入日志,因此恢复时会重试。完全成功后 checkpoint 会被删除;崩溃或存在 失败 step 时会留下一份供恢复。checkpoint 保存失败只会记录日志,不会让进行中的运行失败。

由于 checkpoint 可序列化、executor 是一个参数,宿主可以通过传入另一个节点的 executor,在另一个节点上恢复被中断的工作流(迁移)。

该组合子要求已配置 session store——各 SDK 方法在缺少 store 时都会 reject 或 raise(Node 的错误信息是 parallelResumable requires a sessionStore on the session)。

WorkflowCheckpoint 的模式字段为 schema_version / workflow_id / steps / checkpoint_ms;每条 step 记录可携带一份绑定到该 step 身份(workflow id 加 spec) 的结果回执。恢复采取 fail closed,而不是重跑外部效果不明确的工作:当 checkpoint 无法读取、由未来的 schema_version 写入(ensure_loadable)、属于另一个 workflow id,或记录的结果所对应 task_id 的 spec 已改变时,每个 step 都返回 success: false,并带有 workflow checkpoint cannot be resumed 信息。specs 中 重复出现的 task_id 对应的 step 也会被拒绝。

store 相关见 Persistence,迁移路径见 Multi-Machine。

TypeScript
import { Agent, FileSessionStore } from '@a3s-lab/code';
const agent = await Agent.create('agent.acl');
const session = agent.session('/repo', {
sessionStore: new FileSessionStore('./.a3s/sessions'),
});
// 第一次尝试可能会在中途被打断。
let outcomes = await session.parallelResumable(specs, 'release-batch-42');
// 崩溃或重启后:使用同一个 workflowId 恢复,并跳过已完成的步骤。
outcomes = await session.parallelResumable(specs, 'release-batch-42');
Python
from a3s_code import Agent, FileSessionStore, SessionOptions
agent = Agent.create("agent.acl")
opts = SessionOptions()
opts.session_store = FileSessionStore("./.a3s/sessions")
session = agent.session("/repo", opts)
# 第一次尝试可能会在中途被打断。
outcomes = session.parallel_resumable(specs, "release-batch-42")
# 崩溃或重启后:使用同一个 workflow_id 恢复,并跳过已完成的步骤。
outcomes = session.parallel_resumable(specs, "release-batch-42")
Go
session, err := agent.Session(ctx, "/repo", &code.SessionOptions{
FileSessionStoreDir: "./.a3s/sessions",
})
if err != nil {
return err
}
// 使用同一个 workflow id 重新运行会跳过已完成的步骤。
outcomes, err := session.ParallelResumable(ctx, specs, "release-batch-42")

跨扇出任务的共享预算 [#shared-budget]

默认情况下,每个子代理各自统计自己的 LLM 成本。给 parallel 传入一个 token 预算后,所有子智能体改为汇入同一个账本——为整个扇出设一个统一上限。它映射 到核心的 WorkflowBudget:一个安装到每个子运行上的、聚合型 BudgetGuard。若 session 已设置自己的预算守卫,它会被包裹在内,宿主记账仍能收到每次调用。

预算是一个可选参数:

  • 不传预算时,parallel(specs) 返回普通的结果数组。
  • 传入预算时,parallel(specs, budgetTokens) 解析为 { outcomes, budget }, 其中 budget 是账本快照(consumedTokens / limitTokens)。

一旦达到上限,每个子运行的下一次 LLM 调用和工具调用都会被拒绝,因此之后启动的 step 会以 success: false 结束,并带有预算耗尽信息。它是一个软上限:由于用量 是在每次 LLM 调用之后记账的,宽扇出可能在账本追上之前冲过上限几个在飞回合。 框架绝不强杀进行中的扇出。

TypeScript
import type { WorkflowParallelResult } from '@a3s-lab/code';
// 不传预算 → 普通结果数组。
const outcomes = await session.parallel(specs);
// 传入预算 → { outcomes, budget }。所有子代理共享一个账本。
// TS 返回类型是联合类型,因此对带预算的调用做类型收窄。
const { outcomes: out, budget } = (await session.parallel(
specs,
500_000,
)) as WorkflowParallelResult;
console.log(budget.consumedTokens, budget.limitTokens); // 例如 48213, 500000
Python
# 不传预算 → 普通列表。
outcomes = session.parallel(specs)
# 传入预算 → {"outcomes": [...], "budget": {"consumed_tokens", "limit_tokens"}}。
res = session.parallel(specs, budget_tokens=500_000)
print(res["budget"]["consumed_tokens"], res["budget"]["limit_tokens"])

Go 始终返回带有 Budget 快照的 *code.ParallelResult;未传入上限时 Budget.LimitTokens 为 nil。

Go
limit := uint64(500_000)
result, err := session.Parallel(ctx, specs, &limit)
if err == nil && result.Budget != nil {
fmt.Println(result.Budget.ConsumedTokens, *result.Budget.LimitTokens)
}

循环直到完成(execute_loop)[#looping]

对于长度未知、需要迭代至收敛的工作(循环直至没有剩余项、反复打磨直到满意),核心语法 提供了 execute_loop。每一轮都是一个屏障(execute_steps_parallel);宿主提供的 谓词看到本轮的结果,返回 LoopDecision::Continue(next_specs) 或 LoopDecision::Stop。必填的 max_iterations 是一个硬上限——一旦达到,即使谓词 还想继续也会停止,从而让 LLM 驱动的循环永远不会失控。

Rust
use a3s_code_core::orchestration::{execute_loop, AgentStepSpec, LoopDecision};
let outcomes = execute_loop(executor, initial_specs, /* max_iterations */ 5, None, |round| {
// 本轮没有新发现就停止;否则扇出后续 step。
let follow_ups = derive_follow_ups(round);
if follow_ups.is_empty() {
LoopDecision::Stop
} else {
LoopDecision::Continue(follow_ups)
}
})
.await;

从宿主 SDK 你并不需要专门的 loop 动词——直接用你自己语言里的 while/for 围绕 parallel 写循环,根据结果决定下一轮即可。execute_loop 是为 Rust 语法 而存在,并为循环提供一个单一、强制的终止守卫。

工作流外观接口(Rust / 嵌入)[#workflow-facade]

session.workflow() 返回一个可廉价克隆的 Workflow,它预先接好了会话的 executor、持久化 store、逐 step 事件流,以及一个稳定的、由会话派生的 root id。它是 把以上能力打包起来的可编程句柄;控制流就是普通 Rust——await 一个动词、查看结果、 决定下一步运行什么。

  • 动词——agent(单步)、parallel(屏障式扇出)、phase(命名的、 可恢复的屏障,并发出里程碑)、pipeline(按条目的链),以及不会失败的 log。 每个动词都只委派给一个 combinator。
  • Phase 与事件——phase(name, specs) 派生确定性 checkpoint id ({root}/{index}:{name}),在配置了 store 时走可恢复屏障,并在一个广播上发出 WorkflowEvent::PhaseStart / PhaseEnd,你可用 subscribe() 读取。log() 发出 WorkflowEvent::Log。
  • 预算——每个 session 工作流都带有一个共享的 WorkflowBudget; session.workflow_with_token_budget(Some(limit)) 为其设置上限。 budget_snapshot() 读取账本;结束时已达到上限的 phase 会发出 WorkflowEvent::BudgetExhausted。
Rust
let wf = session.workflow(); // 或 session.workflow_with_token_budget(Some(500_000))
let mut events = wf.subscribe();
// 先跑一步,再根据其结果计算出一个*可变*数量的扇出——这正是“动态”所在:
// 形状在运行时决定,而非提前声明。
let plan = wf.agent(AgentStepSpec::new("plan", "plan", "plan", goal)).await;
let specs = derive_specs(&plan); // 你的代码
let done = wf.phase("implement", specs).await; // 可恢复屏障 + 里程碑
let reviews = wf.phase("review", to_review(&done)).await; // 预算在各 phase 间共享
if let Some(b) = wf.budget_snapshot() {
println!("spent {} / {:?} tokens", b.consumed_tokens, b.limit_tokens);
}

SDK 暴露的是扁平的 parallel / pipeline / parallelResumable 动词(以及上面 parallel 的预算重载);完整的 Workflow 句柄——phases、事件订阅、loop combinator——属于 Rust / 嵌入层 API。

受模式约束的步骤输出 [#schema-forced-step-output]

携带 output_schema(Node 中为 outputSchema,Go 中为 OutputSchema)的任务规格会强制步骤返回符合该 JSON Schema 的值;经校验的对象落在 StepOutcome.structured 中。这复用了与 A3S Code 其余 部分相同的结构化输出强转 + 修复机制。强转失败会把该 step 降级为不成功 (success: false),因此调用方绝不会把未经校验的文本当作承诺的对象。

强制模式约束适用于 parallel / parallelResumable 的任务规格。Node 和 Python 的 pipeline 阶段会忽略它;Go 的 pipeline 阶段会遵循它。

TypeScript
const [outcome] = await session.parallel([
{
taskId: 'triage',
agent: 'review',
description: 'Structured triage',
prompt: 'Triage this diff.',
outputSchema: {
type: 'object',
properties: {
severity: { type: 'string', enum: ['low', 'medium', 'high'] },
summary: { type: 'string' },
},
required: ['severity', 'summary'],
},
},
]);
if (outcome.success) {
console.log(outcome.structured.severity, outcome.structured.summary);
}
Python
outcomes = session.parallel([
{
"task_id": "triage",
"agent": "review",
"description": "Structured triage",
"prompt": "Triage this diff.",
"output_schema": {
"type": "object",
"properties": {
"severity": {"type": "string", "enum": ["low", "medium", "high"]},
"summary": {"type": "string"},
},
"required": ["severity", "summary"],
},
},
])
outcome = outcomes[0]
if outcome["success"]:
print(outcome["structured"]["severity"], outcome["structured"]["summary"])

成本治理与生命周期

编排 step 走的是同一个 session,因此 session 的各项控制对它们直接生效。 session 预算守卫(Node 为 setBudgetGuard,Python 为 budget_guard / set_budget_guard,Go 为 SetBudgetGuard)会安装到每个 step 上;close() 会触发 session 取消令牌,所有派生的 step 都继承该令牌。宿主身份标签(tenant_id、 principal、agent_template_id、correlation_id)保留在父 session 上,供宿主侧 聚合与计费使用。这些控制的细节见 Sessions 与 Limits。