#编排
本页展示 A3S Code 中的可编程编排原语:用于扇出的 session.parallel、用于按条目执行多阶段链的 session.pipeline,以及用于在崩溃后仍可恢复的带日志运行的 session.parallelResumable。parallel 还接受一个可选的 token 预算,让整个扇出共享同一个账本。当你有多个相互独立的子代理任务时使用 parallel;当你需要让每个输入流经一组有序阶段时使用 pipeline。
关于这些原语背后的概念模型,请参阅编排。
#用 session.parallel 扇出
parallel 接收一个 AgentStepSpec 数组并发执行它们,按输入顺序(而非完成顺序)为每个 spec 返回一个 StepOutcome。每个 spec 路由到一个具名子代理(explore、plan、review、verification、general 等)。在 spec 上设置 outputSchema / output_schema 即可拿到经过 schema 校验的 structured 结果。
use a3s_code_core::{Agent, AgentStepSpec};#[tokio::main]async fn main() -> a3s_code_core::Result<()> {let agent = Agent::new("agent.acl").await?;let session = agent.session_builder(".").build().await?;let workflow = session.workflow();let outcomes = workflow.parallel(vec![AgentStepSpec::new("langs","general","列出语言","列出三种系统编程语言。",).with_max_steps(2),AgentStepSpec::new("verdict","general","分类","Rust 是否无需 GC 就能保证内存安全?只回答是或否。",).with_max_steps(2),]).await;for outcome in outcomes {println!("[parallel] {}: success={}",outcome.task_id, outcome.success);}session.close().await;agent.close().await;Ok(())}
import { Agent } from '@a3s-lab/code';const agent = await Agent.create('agent.acl');const session = agent.session('.', {});// 各步骤互不依赖;结果按输入顺序返回,而不是按完成顺序返回。const outcomes = await session.parallel([{taskId: 'langs',agent: 'general',description: 'list',prompt: 'Name three systems languages.',maxSteps: 2,},{taskId: 'safe',agent: 'general',description: 'classify',prompt: 'Is Rust memory-safe without a GC? yes/no.',maxSteps: 2,},]);for (const o of outcomes) {console.log(`[parallel] ${o.taskId}: success=${o.success}`);}session.close();
from a3s_code import Agent, SessionOptionsagent = Agent.create("agent.acl")session = agent.session(".", SessionOptions())# 各步骤互不依赖;结果按输入顺序返回,而不是按完成顺序返回。outcomes = session.parallel([{"task_id": "langs","agent": "general","description": "list languages","prompt": "Name three systems programming languages, comma-separated.","max_steps": 2,},{"task_id": "verdict","agent": "general","description": "classify","prompt": "Is Rust memory-safe without a GC? Answer yes or no.","max_steps": 2,# 这个步骤会返回经过模式校验的结构化输出。"output_schema": {"type": "object","properties": {"memory_safe": {"type": "boolean"}},"required": ["memory_safe"],},},])for o in outcomes:print(f"[parallel] {o['task_id']}: success={o['success']} structured={o.get('structured')}")session.close()
package mainimport ("context""encoding/json""fmt""log"code "github.com/A3S-Lab/Code/sdk/go/v6")func main() {ctx := context.Background()agent, err := code.Create(ctx, "agent.acl")if err != nil {log.Fatal(err)}defer agent.Close(ctx)session, err := agent.Session(ctx, ".", nil)if err != nil {log.Fatal(err)}defer session.Close(ctx)maxSteps := uint(2)result, err := session.Parallel(ctx, []code.AgentStepSpec{{TaskID: "langs",Description: "列出语言",Agent: "general",Prompt: "列出三种系统编程语言,用逗号分隔。",MaxSteps: &maxSteps,},{TaskID: "verdict",Description: "分类",Agent: "general",Prompt: "Rust 是否无需 GC 就能保证内存安全?只回答是或否。",MaxSteps: &maxSteps,OutputSchema: json.RawMessage(`{"type":"object","properties":{"memory_safe":{"type":"boolean"}},"required":["memory_safe"]}`),},}, nil)if err != nil {log.Fatal(err)}for _, outcome := range result.Outcomes {fmt.Printf("[parallel] %s: success=%t structured=%s\n",outcome.TaskID,outcome.Success,outcome.Structured,)}}
结果在 Python 中是字典,在 Node 中是对象,在 Go 中是 StepOutcome 值。会话选项
maxParallelTasks / max_parallel_tasks / MaxParallelTasks 限制并发量;多出的
spec 会排队,而返回的结果数组仍然完整且保持顺序。
#用 session.pipeline 构建按条目执行的链
pipeline 接收一个输入 items 列表和一个有序的 stages 列表。每个条目独立地流经各个阶段——阶段之间没有屏障,因此一个较快的条目可以在一个较慢的条目仍处于阶段 1 时就到达阶段 2。阶段回调接收一个 ctx:第一个阶段看到 ctx.item,后续阶段看到 ctx.previous(上一个 StepOutcome,你可以基于其 .output 继续构建)。返回下一个 spec 以继续,或返回 null / None 以提前停止该条目的链。
use std::sync::Arc;use a3s_code_core::{Agent, AgentStepSpec, PipelineStage};#[tokio::main]async fn main() -> a3s_code_core::Result<()> {let agent = Agent::new("agent.acl").await?;let session = agent.session_builder(".").build().await?;let workflow = session.workflow();let stages: Vec<PipelineStage<String>> = vec![Arc::new(|_, item| {Some(AgentStepSpec::new("summarize","general","总结",format!("用一句话说明什么是{item}。"),).with_max_steps(2),)}),Arc::new(|previous, _| {previous.map(|outcome| {AgentStepSpec::new("classify","general","分类",format!("只回答是或否:下面描述的是编程语言吗?\n\n{}",outcome.output),).with_max_steps(2)})}),];let results = workflow.pipeline(vec!["Rust 编程语言".to_string()], stages).await;for result in results.into_iter().flatten() {println!("[pipeline] {}", result.output);}session.close().await;agent.close().await;Ok(())}
import { Agent } from '@a3s-lab/code';const agent = await Agent.create('agent.acl');const session = agent.session('.', {});// 第二阶段使用第一阶段的输出。阶段回调不能抛出异常;// 返回 null 可停止当前条目的处理链。const results = await session.pipeline(['the Rust programming language'],[(ctx) => ({taskId: 'sum',agent: 'general',description: 'summarize',prompt: `In one sentence, what is ${ctx.item}?`,maxSteps: 2,}),(ctx) => ({taskId: 'cls',agent: 'general',description: 'classify',prompt: `Reply YES or NO: does this describe a programming language?\n\n${ctx.previous.output}`,maxSteps: 2,}),],);for (const r of results) {console.log(`[pipeline] final=${r === null ? null : JSON.stringify(r.output.slice(0, 60))}`,);}session.close();
from a3s_code import Agent, SessionOptionsagent = Agent.create("agent.acl")session = agent.session(".", SessionOptions())# 每个条目依次经过各阶段,第二阶段使用第一阶段的输出。# 阶段返回 None(抛出的异常也会被视为 None)可停止当前条目的处理链。results = session.pipeline(["the Rust programming language"],[lambda ctx: {"task_id": "summarize","agent": "general","description": "summarize","prompt": f"In one sentence, what is {ctx['item']}?","max_steps": 2,},lambda ctx: {"task_id": "classify","agent": "general","description": "classify","prompt": "Reply with one word YES or NO: does this describe a "f"programming language?\n\n{ctx['previous']['output']}","max_steps": 2,},],)for r in results:print(f"[pipeline] final={None if r is None else r['output'][:60]!r}")session.close()
package mainimport ("context""fmt""log"code "github.com/A3S-Lab/Code/sdk/go/v6")func main() {ctx := context.Background()agent, err := code.Create(ctx, "agent.acl")if err != nil {log.Fatal(err)}defer agent.Close(ctx)session, err := agent.Session(ctx, ".", nil)if err != nil {log.Fatal(err)}defer session.Close(ctx)maxSteps := uint(2)results, err := session.Pipeline(ctx,[]any{"Rust 编程语言"},[]code.PipelineStage{func(_ context.Context,stage code.PipelineContext,) (*code.AgentStepSpec, error) {return &code.AgentStepSpec{TaskID: "summarize",Agent: "general",Description: "总结",Prompt: fmt.Sprintf("用一句话说明什么是 %v。", stage.Item),MaxSteps: &maxSteps,}, nil},func(_ context.Context,stage code.PipelineContext,) (*code.AgentStepSpec, error) {return &code.AgentStepSpec{TaskID: "classify",Agent: "general",Description: "分类",Prompt: "只回答是或否:下面描述的是编程语言吗?\n\n" + stage.Previous.Output,MaxSteps: &maxSteps,}, nil},},)if err != nil {log.Fatal(err)}for _, outcome := range results {if outcome != nil {fmt.Println(outcome.Output)}}}
与 parallel 的关键区别:阶段是有序且相互依赖的,但各条目在阶段之间不会
彼此等待。Node 的阶段回调绝不能抛出异常——出错时返回 null;Python 的阶段可以
抛出并被视为 None;Go 阶段返回 (*AgentStepSpec, error)。
#用 session.parallelResumable 恢复运行
parallelResumable 就是带日志的 parallel。它的第一个参数是 specs,第二个参数是
稳定的 workflowId;每个步骤的结果都会被记录到会话的 store,因此如果进程在
运行中途崩溃,你可以用同一个 workflowId 再次调用它,已完成步骤会从日志中重放。
它需要一个会话 store——打开会话时传入 sessionStore、session_store 或
Go 的 FileSessionStoreDir。
Rust 通过具名 workflow.phase 暴露相同的检查点行为。配置会话存储后,每个 phase
都是可恢复的屏障。
use a3s_code_core::{Agent, AgentStepSpec, SessionOptions};#[tokio::main]async fn main() -> a3s_code_core::Result<()> {let agent = Agent::new("agent.acl").await?;let session = agent.session_builder(".").options(SessionOptions::new().with_file_session_store("./.a3s/sessions").with_session_id("nightly-audit-session"),).build().await?;let workflow = session.workflow();let outcomes = workflow.phase("nightly-audit",vec![AgentStepSpec::new("deps","general","审计依赖","检查清单中的过时依赖。",).with_max_steps(2),AgentStepSpec::new("tests","verification","运行测试","运行测试套件并总结失败。",).with_max_steps(2),],).await;for outcome in outcomes {println!("{}:{}", outcome.task_id, outcome.success);}session.close().await;agent.close().await;Ok(())}
import { Agent, FileSessionStore } from '@a3s-lab/code';const agent = await Agent.create('agent.acl');// parallelResumable 会把日志写入会话存储;没有存储时会抛出异常。const session = agent.session('.', {sessionStore: new FileSessionStore('./.a3s/sessions'),});// 签名是 (specs, workflowId):先传规格,再传稳定的工作流标识。const outcomes = await session.parallelResumable([{taskId: 'deps',agent: 'general',description: 'audit deps',prompt: 'Check manifests for outdated dependencies.',maxSteps: 2,},{taskId: 'tests',agent: 'verification',description: 'run tests',prompt: 'Run the test suite and summarize failures.',maxSteps: 2,},],'nightly-audit',);// 中断后使用相同 workflowId 重启。已完成步骤从日志加载;// 完全成功的运行会删除检查点。console.log(outcomes.map((o) => `${o.taskId}:${o.success}`).join(' '));session.close();
from a3s_code import Agent, SessionOptions, FileSessionStoreagent = Agent.create("agent.acl")# parallel_resumable 会把日志写入会话存储;没有存储时会抛出异常。opts = SessionOptions()opts.session_store = FileSessionStore("./.a3s/sessions")session = agent.session(".", opts)# 签名是 (specs, workflow_id):先传规格,再传稳定的工作流标识。outcomes = session.parallel_resumable([{"task_id": "deps", "agent": "general", "description": "audit deps", "prompt": "Check manifests for outdated dependencies.", "max_steps": 2},{"task_id": "tests", "agent": "verification", "description": "run tests", "prompt": "Run the test suite and summarize failures.", "max_steps": 2},],"nightly-audit",)# 中断后使用相同 workflow_id 重启。已完成步骤从日志加载;# 完全成功的运行会删除检查点。print(" ".join(f"{o['task_id']}:{o['success']}" for o in outcomes))session.close()
package mainimport ("context""fmt""log"code "github.com/A3S-Lab/Code/sdk/go/v6")func main() {ctx := context.Background()agent, err := code.Create(ctx, "agent.acl")if err != nil {log.Fatal(err)}defer agent.Close(ctx)options := &code.SessionOptions{SessionID: "nightly-audit-session",FileSessionStoreDir: "./.a3s/sessions",}session, err := agent.Session(ctx, ".", options)if err != nil {log.Fatal(err)}defer session.Close(ctx)maxSteps := uint(2)outcomes, err := session.ParallelResumable(ctx, []code.AgentStepSpec{{TaskID: "deps",Agent: "general",Description: "审计依赖",Prompt: "检查 manifest 中是否存在过期依赖。",MaxSteps: &maxSteps,},{TaskID: "tests",Agent: "verification",Description: "运行测试",Prompt: "运行测试套件并总结失败。",MaxSteps: &maxSteps,},}, "nightly-audit")if err != nil {log.Fatal(err)}for _, outcome := range outcomes {fmt.Printf("%s:%t ", outcome.TaskID, outcome.Success)}}
#用 session.parallel 做预算受限的扇出
传入 token 预算后,所有子代理就会汇入同一个账本。传入预算时,
parallel 解析为 { outcomes, budget }(账本快照)而非原来的结果数组;一旦达到
上限,之后启动的 step 会被拒绝(success: false)。它是软上限——宽扇出可能冲过
上限几个在飞回合;在飞工作绝不会被强杀。
use a3s_code_core::{Agent, AgentStepSpec};#[tokio::main]async fn main() -> a3s_code_core::Result<()> {let agent = Agent::new("agent.acl").await?;let session = agent.session_builder(".").build().await?;let workflow = session.workflow_with_token_budget(Some(50_000));let outcomes = workflow.parallel(vec![AgentStepSpec::new("a", "general", "问题一", "回答:准备。").with_max_steps(2),AgentStepSpec::new("b", "general", "问题二", "回答:开始。").with_max_steps(2),]).await;for outcome in outcomes {println!("{}: success={}", outcome.task_id, outcome.success);}if let Some(budget) = workflow.budget_snapshot() {println!("spent {} / {} tokens",budget.consumed_tokens,budget.limit_tokens.unwrap_or_default());}session.close().await;agent.close().await;Ok(())}
import { Agent } from '@a3s-lab/code';const agent = await Agent.create('agent.acl');const session = agent.session('.', {});const specs = [{taskId: 'a',agent: 'general',description: 'q1',prompt: 'Reply with one word: ready.',maxSteps: 2,},{taskId: 'b',agent: 'general',description: 'q2',prompt: 'Reply with one word: go.',maxSteps: 2,},];// 传入预算时,parallel() 解析为 { outcomes, budget } —— 所有子代理共享一个账本。// (不传时,parallel(specs) 返回原来的数组。)const { outcomes, budget } = await session.parallel(specs, 50_000);for (const o of outcomes)console.log(`[budget] ${o.taskId}: success=${o.success}`);console.log(`spent ${budget.consumedTokens} / ${budget.limitTokens} tokens`);session.close();
from a3s_code import Agent, SessionOptionsagent = Agent.create("agent.acl")session = agent.session(".", SessionOptions())specs = [{"task_id": "a", "agent": "general", "description": "q1", "prompt": "Reply with one word: ready.", "max_steps": 2},{"task_id": "b", "agent": "general", "description": "q2", "prompt": "Reply with one word: go.", "max_steps": 2},]# 传入预算时,parallel() 返回 {"outcomes", "budget"}(共享账本)。res = session.parallel(specs, budget_tokens=50_000)for o in res["outcomes"]:print(f"[budget] {o['task_id']}: success={o['success']}")print(f"spent {res['budget']['consumed_tokens']} / {res['budget']['limit_tokens']} tokens")session.close()
package mainimport ("context""fmt""log"code "github.com/A3S-Lab/Code/sdk/go/v6")func main() {ctx := context.Background()agent, err := code.Create(ctx, "agent.acl")if err != nil {log.Fatal(err)}defer agent.Close(context.Background())session, err := agent.Session(ctx, ".", nil)if err != nil {log.Fatal(err)}defer session.Close(context.Background())maxSteps := uint(2)budgetTokens := uint64(50_000)result, err := session.Parallel(ctx, []code.AgentStepSpec{{TaskID: "a", Agent: "general", Description: "问题 1", Prompt: "只回答一个词:ready。", MaxSteps: &maxSteps},{TaskID: "b", Agent: "general", Description: "问题 2", Prompt: "只回答一个词:go。", MaxSteps: &maxSteps},}, &budgetTokens)if err != nil {log.Fatal(err)}for _, outcome := range result.Outcomes {fmt.Printf("[budget] %s: success=%t\n", outcome.TaskID, outcome.Success)}if result.Budget != nil && result.Budget.LimitTokens != nil {fmt.Printf("spent %d / %d tokens\n",result.Budget.ConsumedTokens,*result.Budget.LimitTokens,)}}
说明:
- 三个原语都返回按输入顺序对齐的结果。Go 使用
StepOutcome;Node 使用对象; Python 使用字典。 - 在 spec 上设置
outputSchema/output_schema/OutputSchema,可在structured/Structured中拿到解析后的结果。 maxSteps/max_steps/MaxSteps限制每个子代理的步数;会话选项maxParallelTasks/max_parallel_tasks/MaxParallelTasks限制扇出并发量。- 给
parallel传入 token 预算即可让整个扇出对同一个账本计数。Go 把*uint64作为Parallel的第三个参数,并从ParallelResult.Budget读取账本。它是软上限。 - Node 的 pipeline 阶段回调绝不能抛出异常——出错时返回
null。Python 的阶段可以 抛出并被视为None;Go 阶段返回error。
可运行的 Node.js 和 Python 版本见
sdk/node/examples/orchestration/parallel-pipeline.mjs 和
sdk/python/examples/orchestration_workflow.py;上面的 Rust 与 Go 页签均为自包含示例。