Give one language model a large goal and a pile of tools, and it will usually make progress for a while and then wander: it repeats a search, forgets a constraint it read twenty steps ago, spends far more than anyone intended, or declares victory on half the work. The orchestrator is the component that stops this. It owns the goal, breaks it into tasks, decides what runs when, checks what comes back and knows when to stop. The models do the thinking inside each task; the orchestrator does the bookkeeping that models are bad at.
This article builds an orchestrator from first principles as a runtime component. It covers the task graph and its states, how a planner proposes work and why its proposals must be validated, the scheduler loop, budgets that are enforced rather than hoped for, result validation, cancellation and human interrupts, and what to persist. A worked example follows one research job through the system. Patterns for coordinating several specialised agents are compared in multi-agent orchestration, and crash-proof execution is covered in durable agent workflows; this page is about the coordinator in the middle.
What an orchestrator is, and what it is not
An orchestrator is ordinary, deterministic code that holds the state of a job and makes control decisions about it. It is not a smarter prompt and it is not a bigger model. The rule that makes the design work is simple: models propose, the orchestrator disposes. A planner model may propose ten subtasks; the orchestrator decides whether ten is allowed. A worker model may claim it finished; the orchestrator decides whether the output passes. A model may want to call a payment tool; the orchestrator decides whether that needs a human first.
Keeping control decisions in code buys three properties a pure agent loop lacks. Decisions are repeatable, so bugs can be reproduced. Limits are real, because a budget enforced by code cannot be talked out of by a persuasive tool result. And the job is observable, because every state change passes through one place that records it. The cost is that you must model the job explicitly, which is exactly the discipline long-running agent work needs.
The task graph
The central data structure is a directed acyclic graph of tasks. Each node is a unit of work with a kind (LLM call, tool call, human task, or a sub-plan), its inputs, the IDs of nodes it depends on, and a state. Edges mean data dependency: a node becomes ready only when every dependency has succeeded, and its inputs are assembled from their outputs. A graph rather than a list matters because independent work can run in parallel and because a failure only blocks the nodes downstream of it.
Give nodes an explicit state machine and refuse transitions that are not on it. A useful set is PENDING (dependencies not done), READY, RUNNING, WAITING_HUMAN, SUCCEEDED, FAILED and CANCELLED. Every transition is written to an append-only event log with the node ID, the attempt number and a timestamp. The current graph is then a projection of that log, which makes resuming after a crash and answering the question of what happened at step 14 the same operation.
from dataclasses import dataclass, field
from enum import Enum
class S(Enum):
PENDING = 1; READY = 2; RUNNING = 3; WAITING_HUMAN = 4
SUCCEEDED = 5; FAILED = 6; CANCELLED = 7
ALLOWED = {
S.PENDING: {S.READY, S.CANCELLED},
S.READY: {S.RUNNING, S.CANCELLED},
S.RUNNING: {S.SUCCEEDED, S.READY, S.FAILED, S.WAITING_HUMAN, S.CANCELLED},
S.WAITING_HUMAN: {S.READY, S.CANCELLED},
}
@dataclass
class Node:
id: str
kind: str # "llm" | "tool" | "human" | "plan"
spec: dict # prompt template, tool name, args
deps: list = field(default_factory=list)
state: S = S.PENDING
attempts: int = 0
output: object = None
def transition(node, new, log):
if new not in ALLOWED.get(node.state, set()):
raise ValueError(f"{node.id}: {node.state.name} -> {new.name} not allowed")
log.append({"node": node.id, "from": node.state.name, "to": new.name, "attempt": node.attempts})
node.state = new
Planning, bounded
Some graphs are fixed in code: a document pipeline that always extracts, then summarises, then classifies needs no planner. Open-ended goals need dynamic planning, where a planner model returns a list of proposed subtasks with dependencies, and plan nodes can themselves expand later as information arrives. Dynamic planning is where most runaway behaviour starts, so treat a plan as untrusted input and validate it before it touches the graph.
A plan validator checks structure and policy in code. It rejects cycles, references to unknown node IDs, and tools the job is not permitted to use. It enforces a maximum number of nodes per expansion and in total, and a maximum expansion depth, so a plan that keeps spawning sub-plans terminates. It checks that each node spec is complete enough to run. When validation fails, the planner gets the specific errors back and one or two chances to repair the plan; after that the job fails loudly rather than running a malformed graph. How to get good plans in the first place is the subject of agent planner architecture.
The scheduler loop
The heart of the orchestrator is a loop: find ready nodes, admit as many as concurrency and budget allow, dispatch them, wait for any to finish, validate the result, update the graph, and repeat until nothing is runnable. The loop below is a compact asyncio version. It is deliberately boring, which is the point: all the interesting judgement happens inside workers and validators, and the loop only enforces order and limits.
import asyncio
async def run(graph, workers, budget, validate, log, max_parallel=4):
running = {} # asyncio.Task -> Node
while True:
if budget.exhausted():
cancel_all(graph, running, log, reason="budget")
return "BUDGET_EXHAUSTED"
for n in graph.values(): # promote PENDING -> READY
if n.state is S.PENDING and all(graph[d].state is S.SUCCEEDED for d in n.deps):
transition(n, S.READY, log)
ready = [n for n in graph.values() if n.state is S.READY]
for n in ready[: max_parallel - len(running)]:
ticket = budget.reserve(estimate(n)) # None if it would overspend
if ticket is None:
break
transition(n, S.RUNNING, log)
n.attempts += 1
inputs = {d: graph[d].output for d in n.deps}
t = asyncio.create_task(workers[n.kind](n.spec, inputs))
t.ticket = ticket
running[t] = n
if not running:
done = all(n.state in (S.SUCCEEDED, S.CANCELLED) for n in graph.values())
return "DONE" if done else "STUCK" # failed or waiting nodes remain
finished, _ = await asyncio.wait(running, return_when=asyncio.FIRST_COMPLETED)
for t in finished:
n = running.pop(t)
budget.settle(t.ticket, actual_cost(t))
settle_node(n, t, validate, log) # SUCCEEDED, retry as READY, or FAILEDTwo details deserve attention. The loop returns STUCK rather than spinning when work remains but nothing can run, which is what happens when a node fails and its dependants can never become ready, or when a node waits for a human. And the budget is reserved before dispatch and settled after, so a burst of parallel admissions cannot collectively overspend.
Budgets that are enforced
Agent jobs fail expensively, so every job carries limits in several currencies at once: model tokens or money, wall-clock time, total node executions, and per-node attempts. A limit that is only checked after the fact is a report, not a limit. The reserve-then-settle pattern makes it enforceable: before dispatching a node, the orchestrator estimates its worst-case cost (for an LLM node, prompt tokens plus the output token cap), reserves that amount, and refuses dispatch if the reservation would exceed what remains. When the node finishes, the reservation is replaced by the measured cost and the difference returns to the pool.
Make exhaustion a normal, reported outcome. A job that stops at its budget should return what it has, which nodes did not run, and an estimate of what finishing would cost, so a human can extend it. Retries draw on the same budget, which caps how much a flaky tool can burn.
Validating results
A worker reporting success is a claim. The validator turns claims into facts the graph can rely on. Validation runs in layers, cheapest first. Structural checks parse the output against the schema the node promised, such as a JSON object with a list of sources each carrying a URL. Deterministic checks test properties code can verify: cited URLs were actually fetched in this job, numbers sum correctly, a generated query parses. Only then, if the node is important, a verifier model or a rubric judges quality.
Classify every failure before deciding what to do with it. A transient error (timeout, rate limit, connection reset) is retried with backoff and the same input. A correctable error (schema violation, failed check) is retried with the validator's message added to the input, because the model can usually fix what it is told about. A permanent error (forbidden tool, missing permission, invalid task) fails the node immediately. Each class has its own attempt cap. The same classification applied to individual tool calls is covered in tool retry architecture.
Cancellation, timeouts and human interrupts
Cancellation has to propagate both ways. When a user cancels a job, the orchestrator marks every non-terminal node CANCELLED, cancels running tasks, and aborts in-flight external calls where the API allows it. When a node fails permanently, its descendants are cancelled rather than left PENDING forever, while unrelated branches continue. Every running node also gets its own deadline, so one hung tool cannot hold the whole job until the global timeout.
Human interrupts are scheduled work, not exceptions. When policy says an action needs approval, such as sending an email, spending money or writing to production, the node moves to WAITING_HUMAN with the exact proposed action attached, and the loop carries on with everything else. The approval arrives as an event that moves the node back to READY with the approved arguments frozen, so the model cannot quietly change them between approval and execution. A rejection becomes input to a replanning step rather than a crash.
State and durability
Because every transition goes through the event log, persistence is a question of where the log lives. For jobs that finish in seconds, memory is enough. For jobs that run for minutes to days, write the log and node outputs to durable storage on every transition and rebuild the graph from it on restart. Nodes that were RUNNING at the crash are treated as unknown: idempotent work is simply re-run, while side-effecting work is reconciled first by checking whether the external action happened, using an idempotency key derived from the job, node and attempt. Snapshot strategies are covered in agent checkpointing; if you need replay of arbitrary workflow code, use a durable execution engine rather than growing your own.
Worked example: a vendor comparison report
A user asks for a comparison of three observability vendors on pricing, data retention and EU hosting, with a budget of 400,000 tokens and 15 minutes. The planner proposes a graph: three research nodes (one per vendor), each followed by an extraction node that turns findings into a fixed schema, then a compare node, then a write node. The validator accepts it: eight nodes, depth one, only the search and fetch tools.
The scheduler admits all three research nodes at once, reserving 40,000 tokens each. Vendor A and B finish; their extraction nodes start. Vendor C's research times out at its 90-second node deadline, is classified transient, and is retried. Vendor B's extraction fails validation because the retention field is free text instead of a number of days; the retry includes that message and passes. Vendor C's retry succeeds but its extraction cites a URL that was never fetched, a deterministic check failure, so it is retried with the instruction to cite only fetched sources.
With all three extractions accepted, compare runs, then write. Publishing to the shared wiki needs approval, so the write node parks in WAITING_HUMAN with the draft attached. The user approves and the job ends having used 236,000 tokens and eleven node executions, with every retry and its reason in the event log.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Job spends far past its budget | Cost checked after dispatch, or parallel nodes all admitted before any settled | Reserve worst-case cost before dispatch; settle on completion |
| Plan grows without end | Plan nodes may expand recursively with no limit | Cap nodes per expansion, total nodes and depth in the validator |
| Job hangs at 99% | Failed node left dependants PENDING; loop waits forever | Cancel descendants on permanent failure; return STUCK when nothing can run |
| Confident but wrong final answer | Worker success accepted without checks | Layered validation: schema, deterministic checks, then verifier |
| Approved action differs from executed action | Model regenerated arguments after approval | Freeze approved arguments in the event; execute exactly those |
| Duplicate side effects after restart | RUNNING nodes re-run blindly | Idempotency keys per job, node and attempt; reconcile before retrying |
Trade-offs
A static graph is predictable and easy to test but cannot adapt to what the job discovers; dynamic planning adapts but needs validation, caps and evaluation. A central orchestrator is easy to observe but is one more service that must be correct; decentralised agents remove that bottleneck and lose most of the guarantees above. Strict validation catches errors early at the cost of extra calls on every node. Most production systems land on a central orchestrator, a mostly fixed graph with one or two bounded planning points, and validation proportional to how much a node's output matters downstream.
What to do next
- Write down the states a node can be in and the allowed transitions, and make illegal transitions raise.
- Put every transition in an append-only event log and derive the current graph from it.
- Validate every plan in code: no cycles, known tools only, caps on nodes and depth.
- Enforce budgets with reserve-then-settle in tokens, time and executions, and return partial results on exhaustion.
- Add a per-node deadline and classify failures as transient, correctable or permanent with separate attempt caps.
- Route side-effecting actions through WAITING_HUMAN with frozen arguments where policy demands approval.
- Replay a recorded event log in a test and assert the same scheduling decisions are made.