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—falsefor a failed or panicked step (never a dropped sibling).structured(optional) — schema-validated object, present only when the spec carried anoutput_schema.source_anchors— sources (toolplus a normalized URL or workspace-relative path) that successful child research tool calls observed.
The key casing differs by SDK:
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.
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 | nullwherectx = { previous: StepOutcome | null, item } - Python:
stage(ctx) -> spec | Nonewherectx = {"previous": <dict|None>, "item": <item>} - Go:
code.PipelineStage,func(ctx, code.PipelineContext) (*code.AgentStepSpec, error), wherePipelineContext{Previous *StepOutcome, Item any}; returnnil, nilto 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, default30000) fails closed — it is treated asnull, stopping only that chain.outputSchemaon a stage spec is ignored. - Python: a stage callable that raises is caught and treated as
None(stops only that chain).output_schemaon a stage spec is ignored. - Go: each stage step runs through
WorkflowStep, soOutputSchemais honored. A stage (or step request) that returns an error stops that chain, andPipelinereturns the first error together with the outcomes collected so far.
Use parallel for schema-validated steps from
Node or Python.
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.
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 }, wherebudgetis 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.
Go always returns a *code.ParallelResult with a Budget snapshot;
Budget.LimitTokens is nil when no limit was passed.
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.
From the host SDKs you don't need a dedicated
loopverb — write the loop in your own language (while/for) aroundparallel, deciding the next round from the outcomes.execute_loopexists 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-failinglog. 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 emitsWorkflowEvent::PhaseStart/PhaseEndon a broadcast you read withsubscribe().log()emitsWorkflowEvent::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 emitsWorkflowEvent::BudgetExhausted.
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.
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.