Orchestration

Orchestration is an Advanced host-authored surface — the programmable sibling of model-driven delegation. For ordinary coding-agent fan-out, prefer unified task. With Tasks and Teams the LLM decides, at run time, to call task with one or more tasks[] items — the shape of the fan-out is whatever the model chose. Orchestration moves that decision into your code: you express fan-out, pipelines, and verification panels as a grammar, so the shape is reproducible, testable, budget-bounded, and resumable — independent of what the model picks.

Reach for orchestration when the structure of the work is known to the host ahead of time (run these three reviewers in parallel; flow each candidate through explore → verify → review; resume this batch after a crash). Reach for Tasks/Teams when you want the model to decide whether and how to delegate.

The framework / host boundary (the seam)

Everything in this layer is written against a single seam, AgentExecutor: "run this step, give me the result." That seam splits responsibilities cleanly:

  • The framework owns the grammar — which steps exist, how they compose, the concurrency hint, and the serializable contracts AgentStepSpec / StepOutcome.
  • The host owns placement — transport, scheduling, and where a step actually runs.

The in-box default executor (TaskExecutor) runs every step locally, in-process, on tokio. A host substitutes its own AgentExecutor to place steps across a cluster; the combinators never observe where a step ran, so the same orchestration scales from one process to a cluster without change.

concurrency_hint() is advisory, not a hard local bound — it is the lever that lets orchestration scale past a single process (a scheduler-backed host returns its cluster-wide target instead of a local cap).

The SDK wires this up for you: AgentSession::agent_executor() returns the session-backed executor (it runs each step as a child agent on this node, inheriting the session's agent registry, LLM client, workspace, and MCP tools), and session_store() returns the session's store. The parallel / pipeline / parallelResumable methods call these for you.

Step contracts

A step is described by an AgentStepSpec and resolves to a StepOutcome. Both are serializable on purpose: a host may ship a spec to another node, and the resumable combinator persists outcomes into checkpoints.

AgentStepSpec fields:

  • task_id — stable id for the step (you assign it); flows into lifecycle events and checkpoints.
  • agent — registry key of the agent to run (e.g. explore, review).
  • description — short human label for display/tracking.
  • prompt — the instruction handed to the child agent.
  • max_steps (optional) — per-step tool-round cap.
  • parent_session_id (optional) — parent session id for event correlation.
  • output_schema (optional) — when set, the step must return a value conforming to this JSON Schema (see Schema-forced step output).

StepOutcome fields:

  • task_id — the originating step's id.
  • session_id — the child run's session id (failed steps remain addressable).
  • agent — the agent that ran.
  • output — the step's text output.
  • success — false for a failed or panicked step (never a dropped sibling).
  • structured (optional) — schema-validated object, present only when the spec carried an output_schema.
  • source_anchors — sources (tool plus a normalized URL or workspace-relative path) that successful child research tool calls observed.

The key casing differs by SDK:

ConceptNode (camelCase)Python (snake_case)Go (struct field)
step idtaskIdtask_idTaskID
tool-round capmaxStepsmax_stepsMaxSteps
parent sessionparentSessionIdparent_session_idParentSessionID
forced schemaoutputSchemaoutput_schemaOutputSchema
child sessionsessionIdsession_idSessionID
source anchorssourceAnchorssource_anchorsSourceAnchors

agent, description, prompt, output, success, and structured are spelled the same in Node and Python; Go uses the exported Agent, Description, Prompt, Output, Success, and Structured fields.

parallel — barrier fan-out

session.parallel(specs) runs every spec as a fan-out and resolves with one StepOutcome per spec, in input order. It maps to the core execute_steps_parallel combinator. It is a barrier: it awaits every step before returning.

Each branch is isolated — a step that fails or panics becomes success: false; it never drops a sibling. Concurrency is bounded by the executor's concurrency hint (the session's configured parallelism by default). In TypeScript the return type is StepOutcomeObject[] | WorkflowParallelResult; cast to StepOutcomeObject[] when you pass no budget.

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 — no barrier between stages

session.pipeline(items, stages) flows each item through a chain of stages independently — there is no barrier between stages, so item A can be in stage 3 while item B is still in stage 1. Wall-clock time is the slowest single chain, not the sum-of-slowest-per-stage that a per-stage barrier would incur.

Stages are spec-builders, not specs: each stage receives the prior outcome and the original item and returns the next step to run, or null / None to stop that item's chain early. A failed step also stops the chain (a later stage would only build on a failed result). The callback shapes:

  • Node: (ctx) => spec | null where ctx = { previous: StepOutcome | null, item }
  • Python: stage(ctx) -> spec | None where ctx = {"previous": <dict|None>, "item": <item>}
  • Go: code.PipelineStage, func(ctx, code.PipelineContext) (*code.AgentStepSpec, error), where PipelineContext{Previous *StepOutcome, Item any}; return nil, nil to stop

A stage can branch on the prior outcome — e.g. "verify the finding the review stage produced".

Each call resolves with one entry per item, in input order: the item's last outcome, or null / None / nil when its first stage returned no spec.

Constraints:

  • Node: stages must be synchronous. A stage that throws, returns a Promise or a malformed spec, or runs past timeoutMs (the 3rd argument, default 30000) fails closed — it is treated as null, stopping only that chain. outputSchema on a stage spec is ignored.
  • Python: a stage callable that raises is caught and treated as None (stops only that chain). output_schema on a stage spec is ignored.
  • Go: each stage step runs through WorkflowStep, so OutputSchema is honored. A stage (or step request) that returns an error stops that chain, and Pipeline returns the first error together with the outcomes collected so far.

Use parallel for schema-validated steps from Node or Python.

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
},
})

Resumable / migratable workflows

session.parallelResumable(specs, workflowId) (Node) / session.parallel_resumable(specs, workflow_id) (Python) / session.ParallelResumable(ctx, specs, workflowID) (Go) is parallel plus a journal. It maps to execute_steps_parallel_resumable.

At each step boundary it writes a WorkflowCheckpoint to the session store. On resume it skips already-completed steps (reusing their cached outcomes) and re-dispatches only the rest. It records only successful steps — a failed step is not journaled, so it retries on resume. On full success the checkpoint is deleted; a crash or a failed step leaves one behind for resume. A failed checkpoint save is logged and does not fail the live run.

Because the checkpoint is serializable and the executor is a parameter, a host can resume an interrupted workflow on a different node (migration) by passing that node's executor.

This combinator requires a configured session store — the SDK methods reject or raise without one (the Node error message is parallelResumable requires a sessionStore on the session).

The WorkflowCheckpoint schema is schema_version / workflow_id / steps / checkpoint_ms; each step record can carry a result receipt bound to the step's identity (workflow id plus spec). Resume fails closed rather than re-running work whose external effect is ambiguous: when the checkpoint cannot be read, was written by a future schema_version (ensure_loadable), is keyed to another workflow id, or records a result for a task_id whose spec has changed, every step returns success: false with a workflow checkpoint cannot be resumed message. Steps whose task_id appears more than once in specs are refused as well.

See Persistence for the store and Multi-Machine for the migration path.

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'),
});
// First attempt — may be interrupted partway through.
let outcomes = await session.parallelResumable(specs, 'release-batch-42');
// After a crash/restart: same workflowId resumes, skipping completed steps.
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)
# First attempt — may be interrupted partway through.
outcomes = session.parallel_resumable(specs, "release-batch-42")
# After a crash/restart: same workflow_id resumes, skipping completed steps.
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
}
// A rerun with the same workflow id skips completed steps.
outcomes, err := session.ParallelResumable(ctx, specs, "release-batch-42")

Shared budget across a fan-out

By default each child agent counts its own LLM cost. Pass a token budget to parallel and every child instead feeds one shared ledger — a single cap for the whole fan-out. It maps to the core WorkflowBudget, an aggregating BudgetGuard installed onto each child run. It wraps the session's own budget guard, if one is set, so host accounting keeps receiving each call.

The budget is an optional argument:

  • Without a budget, parallel(specs) returns the plain outcomes array.
  • With a budget, parallel(specs, budgetTokens) resolves to { outcomes, budget }, where budget is the ledger snapshot (consumedTokens / limitTokens).

Once the cap is reached, every child's next LLM call and tool call is denied, so a step that starts afterwards ends with success: false and a budget-exhausted message. It is a soft cap: because usage is recorded after each LLM call, a wide fan-out can race a few in-flight turns past the cap before the ledger catches up. The framework never force-kills an in-flight fan-out.

TypeScript
import type { WorkflowParallelResult } from '@a3s-lab/code';
// No budget → plain outcomes array.
const outcomes = await session.parallel(specs);
// With a budget → { outcomes, budget }. All children share one ledger.
// The TS return type is a union, so narrow it for the budgeted call.
const { outcomes: out, budget } = (await session.parallel(
specs,
500_000,
)) as WorkflowParallelResult;
console.log(budget.consumedTokens, budget.limitTokens); // e.g. 48213, 500000
Python
# No budget → plain list.
outcomes = session.parallel(specs)
# With a budget → {"outcomes": [...], "budget": {"consumed_tokens", "limit_tokens"}}.
res = session.parallel(specs, budget_tokens=500_000)
print(res["budget"]["consumed_tokens"], res["budget"]["limit_tokens"])

Go always returns a *code.ParallelResult with a Budget snapshot; Budget.LimitTokens is nil when no limit was passed.

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)
}

Looping until done (execute_loop)

For unknown-length, iterate-until-converge work (loop-until-dry, refine-until-good), the core grammar adds execute_loop. Each round is a barrier (execute_steps_parallel); a host-supplied predicate sees the round's outcomes and returns LoopDecision::Continue(next_specs) or LoopDecision::Stop. A required max_iterations is a hard cap — once reached the loop stops even if the predicate would continue, so an LLM-driven loop can never run away.

Rust
use a3s_code_core::orchestration::{execute_loop, AgentStepSpec, LoopDecision};
let outcomes = execute_loop(executor, initial_specs, /* max_iterations */ 5, None, |round| {
// Stop when the round produced no new findings; otherwise fan out follow-ups.
let follow_ups = derive_follow_ups(round);
if follow_ups.is_empty() {
LoopDecision::Stop
} else {
LoopDecision::Continue(follow_ups)
}
})
.await;

From the host SDKs you don't need a dedicated loop verb — write the loop in your own language (while/for) around parallel, deciding the next round from the outcomes. execute_loop exists for the Rust grammar and to give the loop a single, enforced termination guard.

The Workflow facade (Rust / embedding)

session.workflow() returns a cheaply-clonable Workflow that pre-wires the session's executor, persistence store, per-step event stream, and a stable, session-derived root id. It is the programmable handle that bundles everything above; control flow is ordinary Rust — await a verb, inspect the outcomes, decide what runs next.

  • Verbs — agent (one step), parallel (barrier fan-out), phase (a named, resumable barrier that emits milestones), pipeline (per-item chains), and the non-failing log. Each delegates to exactly one combinator.
  • Phases & events — phase(name, specs) derives a deterministic checkpoint id ({root}/{index}:{name}), runs the resumable barrier when a store is present, and emits WorkflowEvent::PhaseStart / PhaseEnd on a broadcast you read with subscribe(). log() emits WorkflowEvent::Log.
  • Budget — every session workflow carries a shared WorkflowBudget; session.workflow_with_token_budget(Some(limit)) gives it a cap. budget_snapshot() reads the ledger, and a phase that ends with the cap reached emits WorkflowEvent::BudgetExhausted.
Rust
let wf = session.workflow(); // or session.workflow_with_token_budget(Some(500_000))
let mut events = wf.subscribe();
// One step, then a *variable* fan-out computed from its result — the "dynamic"
// part: the shape is decided at run time, not declared up front.
let plan = wf.agent(AgentStepSpec::new("plan", "plan", "plan", goal)).await;
let specs = derive_specs(&plan); // your code
let done = wf.phase("implement", specs).await; // resumable barrier + milestones
let reviews = wf.phase("review", to_review(&done)).await; // budget shared across phases
if let Some(b) = wf.budget_snapshot() {
println!("spent {} / {:?} tokens", b.consumed_tokens, b.limit_tokens);
}

The SDKs expose the flat parallel / pipeline / parallelResumable verbs (and the parallel budget overload above); the full Workflow handle — phases, event subscription, the loop combinator — is a Rust/embedding API.

Schema-forced step output

A spec carrying output_schema (outputSchema in Node, OutputSchema in Go) forces the step to return a value conforming to that JSON Schema; the validated object lands in StepOutcome.structured. This reuses the same structured-output coercion + repair machinery as the rest of A3S Code. A coercion failure demotes the step to unsuccessful (success: false), so callers never treat unvalidated text as the promised object.

Forced schema applies to parallel / parallelResumable specs. Node and Python pipeline stages ignore it; Go pipeline stages honor it.

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"])

Cost governance & lifecycle

Orchestrated steps run through the same session, so the session's controls apply to them directly. The session budget guard (setBudgetGuard in Node, budget_guard / set_budget_guard in Python, SetBudgetGuard in Go) is installed on every step, and close() fires the session cancellation token that every derived step inherits. The host identity labels (tenant_id, principal, agent_template_id, correlation_id) stay on the parent session for host-side aggregation and billing. See Sessions and Limits for the details of those controls.