primer.agents.orchestration
Orchestration: how much autonomy to give the model
Run: python -m primer.agents.orchestration
New to the notation? primer.notation explains every symbol used here from
zero. This lesson builds on the message exchange and tool calling in
primer.agents.llm.
Level 1: The practitioner's guide
In one sentence. Orchestration is the decision of how much of a system's control flow the model gets to decide, from none (code runs fixed steps and asks the model to fill each one) to all of it (a team of models steering each other), and the named patterns in between.
When you need it. You face this decision the moment a task needs more than one model call: a document pipeline, a support desk, a research assistant, anything that must look something up and then act. The tell that you have chosen wrongly in one direction is a fixed pipeline that keeps growing special cases for requests it didn't anticipate; the tell in the other is an agent that spends a minute and a thousand tokens on a request a two-step script answered correctly every time. Anthropic's Building effective agents, the post these pattern names come from, draws the line this way: workflows are model calls "orchestrated through predefined code paths"; agents are systems where the model "dynamically directs its own processes and tool usage". It also gives the rule: find the simplest solution possible, and add complexity only when needed. You don't need orchestration at all when one call with retrieval and a few examples in the prompt answers the question; that post says many applications stop there.
This lesson measures the cost of each step up on one question ("How many PTO days do I get, and what is the meal per-diem?"). A fixed workflow answers it in 2 calls and 93 tokens; a router in 2 calls and 105 tokens; an agent loop in 2 calls and about 990 tokens, ten times the workflow, because it re-sends its tool definitions and results on every call; a supervisor with two specialists in 4 calls and about 600 tokens. Same answer, four prices.
Your options. Seven designs, from the least autonomy to the most:
| Option | What it does | What it guarantees | What it costs | Where it lives |
|---|---|---|---|---|
| Fixed workflow (prompt chaining with gates) | Code runs steps in order; the model fills each one; a code check between steps stops a bad result | The same steps every time; a failed gate stops before the next call is paid for or misled | 2 calls and 93 tokens on the lesson's question; a gate per hand-off to write | Your code |
| Router | One cheap classifier call picks a handler; unknown labels go to a fallback | Exactly one path per request, and never a crash on a label the model invented | One extra call (105 tokens here) and a fallback handler | Your code, with one model decision |
| Parallelization (sectioning, voting) | Independent pieces run at once, or several judges answer the same prompt and the majority wins | Wall-clock time of the slowest branch; three 80%-accurate judges make an 89.6% panel if they err independently | n times the tokens for n branches | Your code (a thread pool) |
| Evaluator-optimizer | One call drafts, another critiques, until it passes or the round limit hits | A checked result, or a flagged best effort at the limit | Up to two calls per round (2 rounds in the lesson's example) | Your code |
| Agent loop | The model picks tools and decides when it is done | Handles requests you did not anticipate | About 990 tokens for the lesson's question, ten times the workflow; less predictability | The model decides; your loop enforces the budget (primer.agents.agent_loop) |
| Orchestrator-workers, multi-agent | A lead model splits the work, specialists each get only their brief, a synthesizer combines | Focused context per specialist and parallel work; subtasks decided at run time | 4 calls and about 600 tokens here; in Anthropic's production research system, about 15 times the tokens of a chat | Several model roles, coordinated by your code |
| Explicit state machine with checkpoints | Code owns every state and transition; the model is consulted at a few named nodes; the run is saved after each step | Auditable, testable transitions and a crash that costs one step, not the run | A state model to design and a store for checkpoints | Your code, or a durable-execution engine |
How to choose. Start from whether the steps vary between requests, or only what happens inside each step.
- Steps known in advance, variation inside them (invoice intake: classify, extract, validate, approve, post): a fixed workflow or state machine, with the model at the judgement nodes only.
- A few kinds of request, each with its own handling: a router in front of the workflows, with a fallback that is safe rather than clever.
- One hard judgement call: voting, but vary the prompt, model or evidence, because copies of the same model tend to make the same mistake.
- Quality that a checker can state as criteria (a citation present, a test passing): evaluator-optimizer, with a round limit.
- The path depends on what is discovered along the way: an agent loop.
- Work that is separable and too big for one context (independent research threads, per-document processing): orchestrator-workers. Not for tasks where every part needs the same context or depends on the others; the same Anthropic report found most coding tasks fall on that side.
- Whatever you pick, code owns the transitions you need to audit, and the model decides only where judgement is needed.
What it costs. Tokens, first: each rung up the spectrum multiplies them, from 93 to 105 to about 990 to about 600 on one toy question, and Anthropic reports agents using about 4 times the tokens of a chat and multi-agent systems about 15 times. Latency: a chain is the sum of its calls, a parallel fan-out is its slowest branch, and every hand-off in a multi-agent design adds a call. Quality: multi-agent can win (that report measured a 90.2% improvement over a single agent on an internal research evaluation) when the work parallelizes, and loses context at every hand-off when it does not. Effort: every pattern in this lesson is under 100 lines of plain Python; a framework adds tracing, persistence and approval hooks, and hidden assumptions. Crashes: without checkpoints, a rerun after a crash repeats every model call (4 instead of 2 in the lesson's invoice run) and risks a different answer the second time.
What breaks.
- A router that trusts the label. The classifier returns free text
(
"Billing-ish, maybe refunds?"). Normalize it, accept only known labels, send the rest to a fallback, and log the misses. - A gate that isn't there. Without a check between chain steps, a one-line outline becomes a paid-for, confident draft of nothing.
- Correlated voters. Five copies of one model with one prompt agree on the same error; the majority formula assumes independence.
- An evaluator that is never satisfied. No round limit, no exit. Put the limit in, and flag what came out at the limit.
- Lost hand-offs. A specialist knows only what the supervisor wrote in its brief. Whatever the supervisor forgot is gone. Test that each specialist sees only its own task, and that a delegation to a specialist that doesn't exist is reported, not crashed on.
- No checkpoint. A crash on the last step restarts from the first, doubling cost and, combined with a non-idempotent side effect, posting an invoice twice.
In the wild. The pattern names come from Anthropic's Building effective agents, which also advises starting with the model API directly before reaching for a framework. LangGraph models a workflow as a graph with explicit state and checkpoints, which is the state-machine row above. The Claude Agent SDK and the OpenAI Agents SDK ship vendor-built agent loops, tools and hand-offs. CrewAI and AutoGen are built around multi-agent teams. Temporal provides durable execution for any code: each completed step's result is stored and a crashed run resumes from there. Anthropic's own research feature is an orchestrator-workers system with a lead model and parallel subagents, and its engineering report is the source of the token multipliers above.
Go deeper. Level 2 builds every pattern in a few dozen lines each and measures it: the four-rung cost table, a chain with a gate, a router with a fallback, voting with the formula for why independent judges help, a supervisor whose specialists see only their brief, and an invoice state machine that survives a crash. If you only needed to choose, you are done.
Level 2: How it works, from scratch
"Agent" covers everything from one model call inside ordinary code to a team of models steering each other. This lesson lays out that range, builds each named pattern in a few dozen lines, and measures what each step up costs. The rule that falls out is simple: use the least autonomy that solves the problem, and say so out loud.
1. The autonomy spectrum
Everyday picture. Four ways to run an office. An assembly line (fixed workflow): the steps never change and people do one task at each station. A receptionist (router) decides which department you go to. A personal assistant (agent loop) decides which errands to run and when the job's done. A team with a manager (multi-agent) splits the work among specialists.
flowchart LR W[Fixed workflow<br/>code decides steps] --> R[Router<br/>LLM picks a path] R --> A[Agent loop<br/>LLM picks tools] A --> M[Multi-agent<br/>agents coordinate]
Reading it: left to right, the model decides more and your code decides less. Each arrow buys flexibility (the system can handle requests you didn't anticipate) and costs predictability, tokens and debuggability. The skill is stopping at the leftmost box that handles your real requests.
Tiny worked example. One question, "How many PTO days do I get, and
what is the meal per-diem when travelling?", answered by each design
(autonomy_costs()):
| design | model calls | why |
|---|---|---|
| fixed workflow | 2 | rewrite as a search, then answer |
| router | 2 | classify, then the chosen handler answers |
| agent loop | 2 | one turn with two parallel tool calls, one answer turn |
| multi-agent | 4 | plan, one call per specialist (2), final answer |
Reading it: the left panel counts model calls and the right counts tokens. Calls barely move until multi-agent doubles them. Tokens tell a sharper story. The agent loop re-sends its tool definitions and tool results on every call, so it uses many times the tokens of the fixed workflow for the same answer. The supervisor pays for planning and hand-offs. None of that is waste if you need the flexibility, and all of it is waste if you don't.
2. The named patterns
These names (from Anthropic's Building effective agents) let you describe a design in one word.
In code: every pattern below is built from ask, one model call with
one user message that returns the reply's text.
Prompt chaining
Everyday picture. A relay race with an inspector at each hand-off: if the baton is dropped, the next runner doesn't start.
Tiny worked example. Step 1 writes an outline, and a code gate checks it has at
least three points. The outline "- scope" fails the gate, so the chain stops after one
call and the draft step never runs.
flowchart LR I[Input] --> S1[LLM: outline] --> G{Gate: 3+ points?} G -->|yes| S2[LLM: draft] --> O[Output] G -->|no| X[Stop with the reason]
Reading it: two model calls in a fixed order, with plain code in the middle. The gate is where you catch a bad intermediate result before paying for, and being misled by, the next step.
In code: run_chain runs a list of ChainSteps in order, feeding each
output into the next prompt and stopping at the first gate that reports a
problem; ChainResult records the outputs and where and why it stopped.
Routing
Everyday picture. The receptionist listens for five seconds and sends you to billing, tech support or the front desk.
Tiny worked example. The classifier replies " Technical\n", which is
normalized to technical and routed. It replies "Billing-ish, maybe refunds?",
which is not a known label, so the request goes to the general fallback
instead of crashing.
flowchart LR Q[Request] --> C[LLM classifier] C -->|billing| B[Billing handler] C -->|technical| T[Tech handler] C -->|anything else| F[Fallback]
Reading it: one cheap model call decides the path, then specialised handling takes over. The "anything else" arrow is essential, because the label is free text from a model and code must never assume it's one of the keys.
In code: route makes the one classifier call, normalizes the label,
swaps anything unknown for the fallback, and hands the request to that
label's handler.
Parallelization: sectioning and voting
Everyday picture. Sectioning: several cooks each make one dish at the same time. Voting: a panel of judges, where the majority decides.
Tiny worked example. Three reviewers judge a SQL snippet: two say
"vulnerable" and one says "safe", so the verdict is ("vulnerable", {"vulnerable": 2, "safe": 1}).
flowchart LR Q[Same prompt] --> R1[Reviewer 1] Q --> R2[Reviewer 2] Q --> R3[Reviewer 3] R1 --> V{Majority} R2 --> V R3 --> V
Reading it: the three calls run at the same time (wall-clock time is the slowest one), and plain code counts the answers. Sectioning looks the same except each branch gets a different sub-task, such as a guardrail check running beside the answer.
Why voting helps, when voters err independently:
Level 3: the formula and its symbols
$$ P(\text{majority right}) = \sum_{k=\lfloor n/2 \rfloor + 1}^{n} \binom{n}{k}\, p^{k} (1-p)^{n-k} $$
Symbols
| Symbol | Meaning |
|---|---|
| $n$ | number of voters (odd, so there are no ties) |
| $p$ | chance a single voter is right |
| $k$ | number of voters who are right |
| $\binom{n}{k}$ | ways to choose which $k$ of the $n$ are right |
In words: add up the chances of every outcome where more than half the voters are right.
On the example: $n = 3$, $p = 0.8$: $0.8^3 + 3 \times 0.8^2 \times 0.2 = 0.512 + 0.384 = 0.896$. Three 80% judges make an 89.6% panel.
Level 3: in Python
In Python:
from math import comb
n, p = 3, 0.8
# the chance exactly k voters are right ...
P = sum(comb(n, k) * p ** k * (1 - p) ** (n - k)
# ... for every k from ⌊n/2⌋ + 1 to n
for k in range(n // 2 + 1, n + 1))
round(P, 3) # → 0.896
Reading it: each curve is a per-voter accuracy, and the x-axis adds voters. Every curve rises: more independent judges, better majority. The catch is the word independent. Five copies of the same model with the same prompt tend to make the same mistake, and then voting buys little. Vary the prompt, the model or the evidence.
In code: run_sections is sectioning: it runs independent pieces of
work on a thread pool and collects every result. vote is voting: it asks
each reviewer the same prompt through run_sections and returns the
majority label with its tally, and majority_accuracy evaluates the formula
above.
Orchestrator-workers
Everyday picture. A project lead reads the brief, splits it into sections, hands each to a writer, then edits the pieces into one report.
Tiny worked example. "What's the meal per-diem, when are receipts due,
and does PTO roll over?" The orchestrator returns three subtasks. Three workers each
find one document (fin-002, fin-004, hr-001), and the synthesizer
writes one answer citing all three.
flowchart TD Q[Request] --> O[Orchestrator LLM:<br/>decide the subtasks] O --> W1[Worker: per-diem] O --> W2[Worker: receipts] O --> W3[Worker: PTO rollover] W1 --> S[Synthesizer LLM:<br/>one cited answer] W2 --> S W3 --> S
Reading it: it looks like sectioning, but the subtasks aren't known in advance. The orchestrator decides them from the request, which is what makes it suited to open-ended requests.
In code: orchestrate asks the orchestrator for subtasks, runs the
workers on them in parallel with run_sections, and asks the synthesizer for
one answer, returning an OrchestratorResult. policy_worker is the worker
in the example: it finds the one current policy document for a subtask.
Evaluator-optimizer
Everyday picture. A writer and an editor. The draft goes back and forth until the editor signs off, or the deadline hits.
Tiny worked example. Draft 1, "PTO rolls over.", gets the feedback "Missing a citation". Draft 2, "Up to 5 unused PTO days roll over [hr-001].", passes. Result: 2 rounds, passed.
flowchart LR T[Task] --> G[Generator drafts] G --> E{Evaluator:<br/>PASS?} E -->|feedback| G E -->|PASS| D[Done] E -->|max rounds| B[Stop: best effort, flagged]
Reading it: the loop only exits on PASS or on the round limit. Without the limit, a critic that's never satisfied loops forever. Clear, checkable criteria make this pattern work.
In code: evaluate_optimize alternates generator drafts and evaluator
verdicts, feeding each critique back, and returns the last draft, the rounds
used and whether it passed.
3. Multi-agent: a supervisor and specialists
Everyday picture. A manager who never does the work, only assigns it and assembles the results. Each specialist gets a short, focused brief.
Tiny worked example. The supervisor sends "PTO days per year?" to hr and
"Meal per-diem?" to finance. The HR specialist never sees the finance question,
and vice versa (tested). A delegation to a non-existent legal specialist
is reported, not crashed on.
sequenceDiagram participant U as User participant S as Supervisor participant H as HR agent participant F as Finance agent U->>S: PTO and meal allowance? S->>H: "PTO days per year?" (only this) S->>F: "Meal per-diem?" (only this) H-->>S: 20 days [hr-001] F-->>S: 75 dollars a day [fin-002] S-->>U: combined answer
Reading it: the supervisor's messages to each specialist are the hand-off, and they're all each specialist knows. That's both the benefit (small, focused context) and the risk (whatever the supervisor forgets to write down is lost). Peer designs, where agents hand control directly to each other, have the same hand-off problem without a central coordinator.
When it's worth it. Use several agents when subtasks are genuinely separable, such as independent research threads or per-document work that would overflow one context. The costs are real: more tokens, context lost at every hand-off, and much harder debugging.
In code: supervise gets a delegation plan from the supervisor, sends
each specialist only its own task (reporting unknown names instead of
crashing), then asks the supervisor for the final answer; SupervisorResult
keeps the delegations and every specialist's answer.
4. An explicit state machine, with the model at specific nodes
Everyday picture. A board game. The squares and the rules for moving between them are fixed. On a few special squares you draw a card (ask the model). Everywhere else, the rules decide.
Tiny worked example. Invoice intake. The model is asked twice: "is this an invoice?" and "extract the number, vendor and amount". Code does everything else: validation, the approval rule (over 5,000 waits for a person), and posting. The runs:
INVOICE INV-2041 ... 1200.00 -> RECEIVED > CLASSIFIED > EXTRACTED > VALIDATED > POSTED > DONE
INVOICE INV-2042 ... 18000.00 -> RECEIVED > CLASSIFIED > EXTRACTED > VALIDATED > AWAITING_APPROVAL
Acme newsletter -> RECEIVED > REJECTED
stateDiagram-v2 [*] --> RECEIVED RECEIVED --> CLASSIFIED: LLM says invoice RECEIVED --> REJECTED: LLM says other CLASSIFIED --> EXTRACTED: LLM extracts fields EXTRACTED --> VALIDATED: code checks pass EXTRACTED --> NEEDS_REVIEW: code checks fail VALIDATED --> AWAITING_APPROVAL: amount > limit VALIDATED --> POSTED: code posts to ledger AWAITING_APPROVAL --> POSTED: human approves POSTED --> DONE
Reading it: every box is a state and every arrow a transition, and only two arrows are labelled "LLM". The rest are rules you can read, test and audit. When something goes wrong, you know exactly which state it was in and why it moved.
In code: InvoiceWorkflow.run walks the states, calling the model only
at RECEIVED and CLASSIFIED; validate_invoice is the code check between
EXTRACTED and VALIDATED; InvoiceWorkflow.approve is the human step that
releases a paused invoice. WorkflowRun holds the current state, its history
and the extracted data.
Checkpoints and durable execution. After every transition the run is saved to disk. Durable execution means a workflow whose progress survives crashes, because each completed step's result is stored and the run resumes from there. Engines like Temporal do this at scale.
Reading it: the same crash happens in both bars. Resuming from the checkpoint finishes with the 2 model calls already made. Restarting from scratch asks the model both questions again, doubling cost, and, worse, risks a different answer the second time.
In code: InvoiceWorkflow saves a checkpoint file after every
transition and loads it at the start of InvoiceWorkflow.run, so a rerun
resumes. resume_comparison crashes the ledger once and counts the model
calls each way.
5. Frameworks vs. plain code
LangGraph models agents as graphs with explicit state and checkpoints. The Claude Agent SDK and OpenAI Agents SDK give you vendor-built agent loops, tools and hand-offs. CrewAI and AutoGen focus on multi-agent teams. Temporal provides durable execution for any code. A framework earns its place with standard patterns, built-in tracing, persistence and human-in-the-loop hooks. Plain code wins for simple flows, full control and fewer dependencies. Everything in this lesson is under 100 lines of plain Python, and knowing that is what lets you judge a framework.
In 20 seconds
- Autonomy spectrum: fixed workflow → router → agent loop → multi-agent. Use the least that works.
- Named patterns: chaining (with gates), routing (with a fallback), parallelization (sectioning, voting), orchestrator-workers, evaluator-optimizer.
- Multi-agent buys focused context and parallelism, and costs tokens, hand-off losses and debuggability.
- Production shape: an explicit state machine, the model consulted only at specific nodes, checkpoints after each step.
Self-test questions
Q: When would you use a fixed workflow instead of an agent? Give an example. A: When the steps are known in advance and the variation is inside each step, not in which steps happen. Invoice intake is the example: classify, extract, validate, approve, post. A workflow is cheaper, faster, testable and auditable. The model is used where judgement is needed (classification, extraction), and code owns the flow. Reach for an agent only when the path genuinely depends on what's discovered along the way.
Q: Single agent vs. multi-agent for a document-processing workflow: what does each side have going for it? A: For one agent (or a workflow): documents flow through the same steps, hand-offs lose context, multi-agent multiplies tokens, and one trace is far easier to debug. For several agents: documents are independent and can be processed in parallel, each document or section may overflow a single context, and specialists (tables, legal clauses, figures) benefit from focused instructions and tools. A common answer is a workflow that fans out per document to a focused worker, which is orchestrator-workers rather than free-form agents talking to each other.
Q: A router's classifier sometimes returns labels that aren't in its list. How should the code handle that?
A: Normalize (case, whitespace), accept only known labels, and send
everything else to a safe fallback. Log the misses and add them to the
classifier's evaluation set. Consider structured output with an enum so
the model can only produce valid labels.
Q: Why checkpoint after every step of a long workflow? A: So a crash costs one step, not the run. Resuming reuses completed work (no repeated model calls, no double-posted invoices when combined with idempotency keys) and avoids getting a different answer on the rerun.
The papers behind this lesson
- Yao et al., ReAct: Synergizing Reasoning and Acting in Language Models (2022). https://arxiv.org/abs/2210.03629. It established the reason-act-observe loop that the "agent loop" rung of the spectrum is built on. annotated companion
Further reading
- Anthropic, Building effective agents (the named patterns): https://www.anthropic.com/engineering/building-effective-agents
- LangGraph docs: https://langchain-ai.github.io/langgraph/
- OpenAI Agents SDK: https://openai.github.io/openai-agents-python/
- CrewAI docs: https://docs.crewai.com/
- AutoGen: https://microsoft.github.io/autogen/
- Temporal, durable execution: https://docs.temporal.io/
1r""" 2# Orchestration: how much autonomy to give the model 3 4Run: `python -m primer.agents.orchestration` 5 6New to the notation? `primer.notation` explains every symbol used here from 7zero. This lesson builds on the message exchange and tool calling in 8`primer.agents.llm`. 9 10## Level 1: The practitioner's guide 11 12**In one sentence.** Orchestration is the decision of how much of a 13system's control flow the model gets to decide, from none (code runs fixed 14steps and asks the model to fill each one) to all of it (a team of models 15steering each other), and the named patterns in between. 16 17**When you need it.** You face this decision the moment a task needs more 18than one model call: a document pipeline, a support desk, a research 19assistant, anything that must look something up and then act. The tell that 20you have chosen wrongly in one direction is a fixed pipeline that keeps 21growing special cases for requests it didn't anticipate; the tell in the 22other is an agent that spends a minute and a thousand tokens on a request a 23two-step script answered correctly every time. Anthropic's *Building 24effective agents*, the post these pattern names come from, draws the line 25this way: workflows are model calls "orchestrated through predefined code 26paths"; agents are systems where the model "dynamically directs its own 27processes and tool usage". It also gives the rule: find the simplest 28solution possible, and add complexity only when needed. You don't need 29orchestration at all when one call with retrieval and a few examples in the 30prompt answers the question; that post says many applications stop there. 31 32This lesson measures the cost of each step up on one question ("How many PTO 33days do I get, and what is the meal per-diem?"). A fixed workflow answers it 34in 2 calls and 93 tokens; a router in 2 calls and 105 tokens; an agent loop 35in 2 calls and about 990 tokens, ten times the workflow, because it re-sends 36its tool definitions and results on every call; a supervisor with two 37specialists in 4 calls and about 600 tokens. Same answer, four prices. 38 39**Your options.** Seven designs, from the least autonomy to the most: 40 41| Option | What it does | What it guarantees | What it costs | Where it lives | 42|---|---|---|---|---| 43| Fixed workflow (prompt chaining with gates) | Code runs steps in order; the model fills each one; a code check between steps stops a bad result | The same steps every time; a failed gate stops before the next call is paid for or misled | 2 calls and 93 tokens on the lesson's question; a gate per hand-off to write | Your code | 44| Router | One cheap classifier call picks a handler; unknown labels go to a fallback | Exactly one path per request, and never a crash on a label the model invented | One extra call (105 tokens here) and a fallback handler | Your code, with one model decision | 45| Parallelization (sectioning, voting) | Independent pieces run at once, or several judges answer the same prompt and the majority wins | Wall-clock time of the slowest branch; three 80%-accurate judges make an 89.6% panel if they err independently | n times the tokens for n branches | Your code (a thread pool) | 46| Evaluator-optimizer | One call drafts, another critiques, until it passes or the round limit hits | A checked result, or a flagged best effort at the limit | Up to two calls per round (2 rounds in the lesson's example) | Your code | 47| Agent loop | The model picks tools and decides when it is done | Handles requests you did not anticipate | About 990 tokens for the lesson's question, ten times the workflow; less predictability | The model decides; your loop enforces the budget (`primer.agents.agent_loop`) | 48| Orchestrator-workers, multi-agent | A lead model splits the work, specialists each get only their brief, a synthesizer combines | Focused context per specialist and parallel work; subtasks decided at run time | 4 calls and about 600 tokens here; in Anthropic's production research system, about 15 times the tokens of a chat | Several model roles, coordinated by your code | 49| Explicit state machine with checkpoints | Code owns every state and transition; the model is consulted at a few named nodes; the run is saved after each step | Auditable, testable transitions and a crash that costs one step, not the run | A state model to design and a store for checkpoints | Your code, or a durable-execution engine | 50 51**How to choose.** Start from whether the *steps* vary between requests, or 52only what happens inside each step. 53 54- Steps known in advance, variation inside them (invoice intake: classify, 55 extract, validate, approve, post): a fixed workflow or state machine, with 56 the model at the judgement nodes only. 57- A few kinds of request, each with its own handling: a router in front of 58 the workflows, with a fallback that is safe rather than clever. 59- One hard judgement call: voting, but vary the prompt, model or evidence, 60 because copies of the same model tend to make the same mistake. 61- Quality that a checker can state as criteria (a citation present, a test 62 passing): evaluator-optimizer, with a round limit. 63- The path depends on what is discovered along the way: an agent loop. 64- Work that is separable and too big for one context (independent research 65 threads, per-document processing): orchestrator-workers. Not for tasks 66 where every part needs the same context or depends on the others; the 67 same Anthropic report found most coding tasks fall on that side. 68- Whatever you pick, code owns the transitions you need to audit, and the 69 model decides only where judgement is needed. 70 71**What it costs.** Tokens, first: each rung up the spectrum multiplies 72them, from 93 to 105 to about 990 to about 600 on one toy question, and 73Anthropic reports agents using about 4 times the tokens of a chat and 74multi-agent systems about 15 times. Latency: a chain is the sum of its 75calls, a parallel fan-out is its slowest branch, and every hand-off in a 76multi-agent design adds a call. Quality: multi-agent can win (that report 77measured a 90.2% improvement over a single agent on an internal research 78evaluation) when the work parallelizes, and loses context at every hand-off 79when it does not. Effort: every pattern in this lesson is under 100 lines 80of plain Python; a framework adds tracing, persistence and approval hooks, 81and hidden assumptions. Crashes: without checkpoints, a rerun after a 82crash repeats every model call (4 instead of 2 in the lesson's invoice 83run) and risks a different answer the second time. 84 85**What breaks.** 86 87- **A router that trusts the label.** The classifier returns free text 88 (`"Billing-ish, maybe refunds?"`). Normalize it, accept only known 89 labels, send the rest to a fallback, and log the misses. 90- **A gate that isn't there.** Without a check between chain steps, a 91 one-line outline becomes a paid-for, confident draft of nothing. 92- **Correlated voters.** Five copies of one model with one prompt agree on 93 the same error; the majority formula assumes independence. 94- **An evaluator that is never satisfied.** No round limit, no exit. Put 95 the limit in, and flag what came out at the limit. 96- **Lost hand-offs.** A specialist knows only what the supervisor wrote in 97 its brief. Whatever the supervisor forgot is gone. Test that each 98 specialist sees only its own task, and that a delegation to a specialist 99 that doesn't exist is reported, not crashed on. 100- **No checkpoint.** A crash on the last step restarts from the first, 101 doubling cost and, combined with a non-idempotent side effect, posting an 102 invoice twice. 103 104**In the wild.** The pattern names come from Anthropic's *Building 105effective agents*, which also advises starting with the model API directly 106before reaching for a framework. LangGraph models a workflow as a graph 107with explicit state and checkpoints, which is the state-machine row above. 108The Claude Agent SDK and the OpenAI Agents SDK ship vendor-built agent 109loops, tools and hand-offs. CrewAI and AutoGen are built around multi-agent 110teams. Temporal provides durable execution for any code: each completed 111step's result is stored and a crashed run resumes from there. Anthropic's 112own research feature is an orchestrator-workers system with a lead model 113and parallel subagents, and its engineering report is the source of the 114token multipliers above. 115 116**Go deeper.** Level 2 builds every pattern in a few dozen lines each and 117measures it: the four-rung cost table, a chain with a gate, a router with a 118fallback, voting with the formula for why independent judges help, a 119supervisor whose specialists see only their brief, and an invoice state 120machine that survives a crash. If you only needed to choose, you are done. 121 122## Level 2: How it works, from scratch 123 124"Agent" covers everything from one model call inside ordinary code to a 125team of models steering each other. This lesson lays out that range, builds 126each named pattern in a few dozen lines, and measures what each step up 127costs. The rule that falls out is simple: **use the least autonomy that 128solves the problem, and say so out loud.** 129 130## 1. The autonomy spectrum 131 132**Everyday picture.** Four ways to run an office. An **assembly line** 133(fixed workflow): the steps never change and people do one task at each 134station. A **receptionist** (router) decides which department you go to. A 135**personal assistant** (agent loop) decides which errands to run and when 136the job's done. A **team with a manager** (multi-agent) splits the work 137among specialists. 138 139```mermaid 140flowchart LR 141 W[Fixed workflow<br/>code decides steps] --> R[Router<br/>LLM picks a path] 142 R --> A[Agent loop<br/>LLM picks tools] 143 A --> M[Multi-agent<br/>agents coordinate] 144``` 145 146**Reading it:** left to right, the model decides more and your code 147decides less. Each arrow buys flexibility (the system can handle requests 148you didn't anticipate) and costs predictability, tokens and debuggability. 149The skill is stopping at the leftmost box that handles your real 150requests. 151 152**Tiny worked example.** One question, "How many PTO days do I get, and 153what is the meal per-diem when travelling?", answered by each design 154(`autonomy_costs()`): 155 156| design | model calls | why | 157|---|---|---| 158| fixed workflow | 2 | rewrite as a search, then answer | 159| router | 2 | classify, then the chosen handler answers | 160| agent loop | 2 | one turn with two parallel tool calls, one answer turn | 161| multi-agent | 4 | plan, one call per specialist (2), final answer | 162 163 164 165**Reading it:** the left panel counts model calls and the right counts tokens. 166Calls barely move until multi-agent doubles them. Tokens tell a sharper story. 167The agent loop re-sends its tool definitions and tool results on every call, 168so it uses many times the tokens of the fixed workflow for the same answer. The 169supervisor pays for planning and hand-offs. None of that is waste *if* you 170need the flexibility, and all of it is waste if you don't. 171 172## 2. The named patterns 173 174These names (from Anthropic's *Building effective agents*) let you describe 175a design in one word. 176 177**In code:** every pattern below is built from `ask`, one model call with 178one user message that returns the reply's text. 179 180### Prompt chaining 181 182**Everyday picture.** A relay race with an inspector at each hand-off: if 183the baton is dropped, the next runner doesn't start. 184 185**Tiny worked example.** Step 1 writes an outline, and a code *gate* checks it has at 186least three points. The outline `"- scope"` fails the gate, so the chain stops after one 187call and the draft step never runs. 188 189```mermaid 190flowchart LR 191 I[Input] --> S1[LLM: outline] --> G{Gate: 3+ points?} 192 G -->|yes| S2[LLM: draft] --> O[Output] 193 G -->|no| X[Stop with the reason] 194``` 195 196**Reading it:** two model calls in a fixed order, with plain code in the 197middle. The gate is where you catch a bad intermediate result before 198paying for, and being misled by, the next step. 199 200**In code:** `run_chain` runs a list of `ChainStep`s in order, feeding each 201output into the next prompt and stopping at the first gate that reports a 202problem; `ChainResult` records the outputs and where and why it stopped. 203 204### Routing 205 206**Everyday picture.** The receptionist listens for five seconds and sends you 207to billing, tech support or the front desk. 208 209**Tiny worked example.** The classifier replies `" Technical\n"`, which is 210normalized to `technical` and routed. It replies `"Billing-ish, maybe refunds?"`, 211which is not a known label, so the request goes to the `general` fallback 212instead of crashing. 213 214```mermaid 215flowchart LR 216 Q[Request] --> C[LLM classifier] 217 C -->|billing| B[Billing handler] 218 C -->|technical| T[Tech handler] 219 C -->|anything else| F[Fallback] 220``` 221 222**Reading it:** one cheap model call decides the path, then specialised 223handling takes over. The "anything else" arrow is essential, because the label is free 224text from a model and code must never assume it's one of the keys. 225 226**In code:** `route` makes the one classifier call, normalizes the label, 227swaps anything unknown for the fallback, and hands the request to that 228label's handler. 229 230### Parallelization: sectioning and voting 231 232**Everyday picture.** *Sectioning:* several cooks each make one dish at the 233same time. *Voting:* a panel of judges, where the majority decides. 234 235**Tiny worked example.** Three reviewers judge a SQL snippet: two say 236"vulnerable" and one says "safe", so the verdict is `("vulnerable", {"vulnerable": 2, "safe": 1})`. 237 238```mermaid 239flowchart LR 240 Q[Same prompt] --> R1[Reviewer 1] 241 Q --> R2[Reviewer 2] 242 Q --> R3[Reviewer 3] 243 R1 --> V{Majority} 244 R2 --> V 245 R3 --> V 246``` 247 248**Reading it:** the three calls run at the same time (wall-clock time is the 249slowest one), and plain code counts the answers. Sectioning looks the same 250except each branch gets a *different* sub-task, such as a guardrail check 251running beside the answer. 252 253Why voting helps, when voters err independently: 254 255$$ 256P(\text{majority right}) = \sum_{k=\lfloor n/2 \rfloor + 1}^{n} \binom{n}{k}\, p^{k} (1-p)^{n-k} 257$$ 258 259**Symbols** 260 261| Symbol | Meaning | 262|---|---| 263| $n$ | number of voters (odd, so there are no ties) | 264| $p$ | chance a single voter is right | 265| $k$ | number of voters who are right | 266| $\binom{n}{k}$ | ways to choose which $k$ of the $n$ are right | 267 268**In words:** add up the chances of every outcome where more than half the 269voters are right. 270 271**On the example:** $n = 3$, $p = 0.8$: $0.8^3 + 3 \times 0.8^2 \times 0.2 = 0.512 + 0.384 = 0.896$. 272Three 80% judges make an 89.6% panel. 273 274**In Python:** 275 276```python 277from math import comb 278n, p = 3, 0.8 279# the chance exactly k voters are right ... 280P = sum(comb(n, k) * p ** k * (1 - p) ** (n - k) 281 # ... for every k from ⌊n/2⌋ + 1 to n 282 for k in range(n // 2 + 1, n + 1)) 283round(P, 3) # → 0.896 284``` 285 286 287 288**Reading it:** each curve is a per-voter accuracy, and the x-axis adds voters. 289Every curve rises: more independent judges, better majority. The catch is 290the word *independent*. Five copies of the same model with the same prompt 291tend to make the *same* mistake, and then voting buys little. Vary the 292prompt, the model or the evidence. 293 294**In code:** `run_sections` is sectioning: it runs independent pieces of 295work on a thread pool and collects every result. `vote` is voting: it asks 296each reviewer the same prompt through `run_sections` and returns the 297majority label with its tally, and `majority_accuracy` evaluates the formula 298above. 299 300### Orchestrator-workers 301 302**Everyday picture.** A project lead reads the brief, splits it into 303sections, hands each to a writer, then edits the pieces into one report. 304 305**Tiny worked example.** "What's the meal per-diem, when are receipts due, 306and does PTO roll over?" The orchestrator returns three subtasks. Three workers each 307find one document (`fin-002`, `fin-004`, `hr-001`), and the synthesizer 308writes one answer citing all three. 309 310```mermaid 311flowchart TD 312 Q[Request] --> O[Orchestrator LLM:<br/>decide the subtasks] 313 O --> W1[Worker: per-diem] 314 O --> W2[Worker: receipts] 315 O --> W3[Worker: PTO rollover] 316 W1 --> S[Synthesizer LLM:<br/>one cited answer] 317 W2 --> S 318 W3 --> S 319``` 320 321**Reading it:** it looks like sectioning, but the subtasks aren't known in 322advance. The orchestrator *decides* them from the request, which is what makes it 323suited to open-ended requests. 324 325**In code:** `orchestrate` asks the orchestrator for subtasks, runs the 326workers on them in parallel with `run_sections`, and asks the synthesizer for 327one answer, returning an `OrchestratorResult`. `policy_worker` is the worker 328in the example: it finds the one current policy document for a subtask. 329 330### Evaluator-optimizer 331 332**Everyday picture.** A writer and an editor. The draft goes back and forth 333until the editor signs off, or the deadline hits. 334 335**Tiny worked example.** Draft 1, "PTO rolls over.", gets the feedback "Missing a citation". 336Draft 2, "Up to 5 unused PTO days roll over [hr-001].", passes. Result: 2 rounds, 337passed. 338 339```mermaid 340flowchart LR 341 T[Task] --> G[Generator drafts] 342 G --> E{Evaluator:<br/>PASS?} 343 E -->|feedback| G 344 E -->|PASS| D[Done] 345 E -->|max rounds| B[Stop: best effort, flagged] 346``` 347 348**Reading it:** the loop only exits on PASS or on the round limit. Without 349the limit, a critic that's never satisfied loops forever. Clear, checkable 350criteria make this pattern work. 351 352**In code:** `evaluate_optimize` alternates generator drafts and evaluator 353verdicts, feeding each critique back, and returns the last draft, the rounds 354used and whether it passed. 355 356## 3. Multi-agent: a supervisor and specialists 357 358**Everyday picture.** A manager who never does the work, only assigns it and 359assembles the results. Each specialist gets a short, focused brief. 360 361**Tiny worked example.** The supervisor sends "PTO days per year?" to `hr` and 362"Meal per-diem?" to `finance`. The HR specialist never sees the finance question, 363and vice versa (tested). A delegation to a non-existent `legal` specialist 364is reported, not crashed on. 365 366```mermaid 367sequenceDiagram 368 participant U as User 369 participant S as Supervisor 370 participant H as HR agent 371 participant F as Finance agent 372 U->>S: PTO and meal allowance? 373 S->>H: "PTO days per year?" (only this) 374 S->>F: "Meal per-diem?" (only this) 375 H-->>S: 20 days [hr-001] 376 F-->>S: 75 dollars a day [fin-002] 377 S-->>U: combined answer 378``` 379 380**Reading it:** the supervisor's messages to each specialist are the 381**hand-off**, and they're *all* each specialist knows. That's both the 382benefit (small, focused context) and the risk (whatever the supervisor 383forgets to write down is lost). Peer designs, where agents hand control 384directly to each other, have the same hand-off problem without a central 385coordinator. 386 387**When it's worth it.** Use several agents when subtasks are genuinely 388separable, such as independent research threads or per-document work that would 389overflow one context. The costs are real: more tokens, context lost at every 390hand-off, and much harder debugging. 391 392**In code:** `supervise` gets a delegation plan from the supervisor, sends 393each specialist only its own task (reporting unknown names instead of 394crashing), then asks the supervisor for the final answer; `SupervisorResult` 395keeps the delegations and every specialist's answer. 396 397## 4. An explicit state machine, with the model at specific nodes 398 399**Everyday picture.** A board game. The squares and the rules for moving 400between them are fixed. On a few special squares you draw a card (ask the 401model). Everywhere else, the rules decide. 402 403**Tiny worked example.** Invoice intake. The model is asked twice: "is 404this an invoice?" and "extract the number, vendor and amount". Code does 405everything else: validation, the approval rule (over 5,000 waits for a person), 406and posting. The runs: 407 408```text 409INVOICE INV-2041 ... 1200.00 -> RECEIVED > CLASSIFIED > EXTRACTED > VALIDATED > POSTED > DONE 410INVOICE INV-2042 ... 18000.00 -> RECEIVED > CLASSIFIED > EXTRACTED > VALIDATED > AWAITING_APPROVAL 411Acme newsletter -> RECEIVED > REJECTED 412``` 413 414```mermaid 415stateDiagram-v2 416 [*] --> RECEIVED 417 RECEIVED --> CLASSIFIED: LLM says invoice 418 RECEIVED --> REJECTED: LLM says other 419 CLASSIFIED --> EXTRACTED: LLM extracts fields 420 EXTRACTED --> VALIDATED: code checks pass 421 EXTRACTED --> NEEDS_REVIEW: code checks fail 422 VALIDATED --> AWAITING_APPROVAL: amount > limit 423 VALIDATED --> POSTED: code posts to ledger 424 AWAITING_APPROVAL --> POSTED: human approves 425 POSTED --> DONE 426``` 427 428**Reading it:** every box is a state and every arrow a transition, and only 429two arrows are labelled "LLM". The rest are rules you can read, test and 430audit. When something goes wrong, you know exactly which state it was in and 431why it moved. 432 433**In code:** `InvoiceWorkflow.run` walks the states, calling the model only 434at RECEIVED and CLASSIFIED; `validate_invoice` is the code check between 435EXTRACTED and VALIDATED; `InvoiceWorkflow.approve` is the human step that 436releases a paused invoice. `WorkflowRun` holds the current state, its history 437and the extracted data. 438 439**Checkpoints and durable execution.** After every transition the run is 440saved to disk. **Durable execution** means a workflow whose progress 441survives crashes, because each completed step's result is stored and the 442run resumes from there. Engines like Temporal do this at scale. 443 444 445 446**Reading it:** the same crash happens in both bars. Resuming from the checkpoint 447finishes with the 2 model calls already made. Restarting from scratch asks 448the model both questions again, doubling cost, and, worse, risks a 449*different* answer the second time. 450 451**In code:** `InvoiceWorkflow` saves a checkpoint file after every 452transition and loads it at the start of `InvoiceWorkflow.run`, so a rerun 453resumes. `resume_comparison` crashes the ledger once and counts the model 454calls each way. 455 456## 5. Frameworks vs. plain code 457 458**LangGraph** models agents as graphs with explicit state and checkpoints. 459The **Claude Agent SDK** and **OpenAI Agents SDK** give you vendor-built agent 460loops, tools and hand-offs. **CrewAI** and **AutoGen** focus on multi-agent 461teams. **Temporal** provides durable execution for any code. A framework 462earns its place with standard patterns, built-in tracing, persistence and 463human-in-the-loop hooks. Plain code wins for simple flows, full control 464and fewer dependencies. Everything in this lesson is under 100 lines of 465plain Python, and knowing that is what lets you judge a framework. 466 467## In 20 seconds 468- Autonomy spectrum: fixed workflow → router → agent loop → multi-agent. Use the least that works. 469- Named patterns: chaining (with gates), routing (with a fallback), parallelization (sectioning, voting), orchestrator-workers, evaluator-optimizer. 470- Multi-agent buys focused context and parallelism, and costs tokens, hand-off losses and debuggability. 471- Production shape: an explicit state machine, the model consulted only at specific nodes, checkpoints after each step. 472 473## Self-test questions 474 475**Q: When would you use a fixed workflow instead of an agent? Give an example.** 476A: When the steps are known in advance and the variation is *inside* each 477step, not in which steps happen. Invoice intake is the example: classify, 478extract, validate, approve, post. A workflow is cheaper, faster, testable 479and auditable. The model is used where judgement is needed (classification, 480extraction), and code owns the flow. Reach for an agent only when the path 481genuinely depends on what's discovered along the way. 482 483**Q: Single agent vs. multi-agent for a document-processing workflow: what does each side have going for it?** 484A: *For one agent (or a workflow):* documents flow through the same steps, 485hand-offs lose context, multi-agent multiplies tokens, and one trace is far 486easier to debug. *For several agents:* documents are independent and can be 487processed in parallel, each document or section may overflow a single 488context, and specialists (tables, legal clauses, figures) benefit from 489focused instructions and tools. A common answer is a workflow that fans out 490per document to a focused worker, which is orchestrator-workers rather than free-form 491agents talking to each other. 492 493**Q: A router's classifier sometimes returns labels that aren't in its list. How should the code handle that?** 494A: Normalize (case, whitespace), accept only known labels, and send 495everything else to a safe fallback. Log the misses and add them to the 496classifier's evaluation set. Consider structured output with an `enum` so 497the model can only produce valid labels. 498 499**Q: Why checkpoint after every step of a long workflow?** 500A: So a crash costs one step, not the run. Resuming reuses completed work 501(no repeated model calls, no double-posted invoices when combined with 502idempotency keys) and avoids getting a different answer on the rerun. 503 504## The papers behind this lesson 505 506- **Yao et al., *ReAct: Synergizing Reasoning and Acting in Language Models* (2022).** 507 https://arxiv.org/abs/2210.03629. It established the reason-act-observe 508 loop that the "agent loop" rung of the spectrum is built on. 509 [annotated companion](../../papers/react.html) 510 511## Further reading 512- Anthropic, *Building effective agents* (the named patterns): https://www.anthropic.com/engineering/building-effective-agents 513- LangGraph docs: https://langchain-ai.github.io/langgraph/ 514- OpenAI Agents SDK: https://openai.github.io/openai-agents-python/ 515- CrewAI docs: https://docs.crewai.com/ 516- AutoGen: https://microsoft.github.io/autogen/ 517- Temporal, durable execution: https://docs.temporal.io/ 518""" 519 520from __future__ import annotations 521 522import json 523import math 524import os 525import re 526from collections import Counter 527from concurrent.futures import ThreadPoolExecutor 528from dataclasses import dataclass, field 529from pathlib import Path 530from typing import Any, Callable 531 532import numpy as np 533 534from primer.agents.llm import LLM, ScriptedLLM, last_user_text 535 536 537def ask(llm: LLM, prompt: str, system: str = "") -> str: 538 """One model call with one user message; returns the text.""" 539 return llm.complete(system=system, messages=[{"role": "user", "content": prompt}]).text 540 541 542# --------------------------------------------------------------------------- 543# Prompt chaining 544# --------------------------------------------------------------------------- 545 546 547@dataclass 548class ChainStep: 549 name: str 550 prompt: str # template; "{input}" is replaced by the previous step's output 551 gate: Callable[[str], str | None] | None = None # returns a problem, or None to continue 552 553 554@dataclass 555class ChainResult: 556 outputs: dict[str, str] 557 stopped_at: str | None = None 558 reason: str | None = None 559 560 561def run_chain(llm: LLM, task: str, steps: list[ChainStep]) -> ChainResult: 562 """Fixed sequence of calls; each output feeds the next, with a code check (gate) in between.""" 563 outputs: dict[str, str] = {} 564 current = task 565 for step in steps: 566 current = ask(llm, step.prompt.format(input=current)) 567 outputs[step.name] = current 568 problem = step.gate(current) if step.gate else None 569 if problem: 570 return ChainResult(outputs, step.name, problem) 571 return ChainResult(outputs) 572 573 574# --------------------------------------------------------------------------- 575# Routing 576# --------------------------------------------------------------------------- 577 578 579def route( 580 classifier: LLM, handlers: dict[str, Callable[[str], str]], request: str, fallback: str = "general" 581) -> tuple[str, str]: 582 """Classify the request with one model call, then hand it to that label's handler. 583 584 The label comes back as free text, so it's normalized and checked against 585 the known labels. Anything else goes to the fallback instead of crashing 586 or guessing. 587 """ 588 raw = ask(classifier, f"Classify into one of {sorted(handlers)}. Reply with the label only.\n\n{request}") 589 label = raw.strip().lower() 590 if label not in handlers: 591 label = fallback 592 return label, handlers[label](request) 593 594 595# --------------------------------------------------------------------------- 596# Parallelization: sectioning and voting 597# --------------------------------------------------------------------------- 598 599 600def run_sections(sections: dict[str, Callable[[], str]]) -> dict[str, str]: 601 """Run independent pieces of work at the same time; wall-clock time is the slowest piece.""" 602 with ThreadPoolExecutor(max_workers=len(sections)) as pool: 603 futures = {name: pool.submit(fn) for name, fn in sections.items()} 604 return {name: f.result() for name, f in futures.items()} 605 606 607def vote(reviewers: list[LLM], prompt: str) -> tuple[str, dict[str, int]]: 608 """Ask several models (or one model several times) the same question; take the majority.""" 609 answers = run_sections({str(i): (lambda r=r: ask(r, prompt).strip().lower()) for i, r in enumerate(reviewers)}) 610 tally = Counter(answers.values()) 611 return tally.most_common(1)[0][0], dict(tally) 612 613 614def majority_accuracy(p: float, n_voters: int) -> float: 615 """Chance a majority of n independent voters, each right with probability p, is right.""" 616 return sum(math.comb(n_voters, k) * p**k * (1 - p) ** (n_voters - k) for k in range(n_voters // 2 + 1, n_voters + 1)) 617 618 619# --------------------------------------------------------------------------- 620# Orchestrator-workers 621# --------------------------------------------------------------------------- 622 623 624@dataclass 625class OrchestratorResult: 626 subtasks: list[str] 627 worker_outputs: dict[str, str] 628 final: str 629 630 631def orchestrate(orchestrator: LLM, worker: Callable[[str], str], synthesizer: LLM, request: str) -> OrchestratorResult: 632 """One model splits the request into subtasks, workers handle them in parallel, one model combines. 633 634 Unlike sectioning, the subtasks aren't known in advance: the orchestrator 635 decides them from the request. 636 """ 637 plan = ask(orchestrator, f'Split into independent subtasks. Reply as JSON {{"subtasks": [...]}}.\n\n{request}') 638 subtasks = list(json.loads(plan)["subtasks"]) 639 outputs = run_sections({s: (lambda s=s: worker(s)) for s in subtasks}) 640 notes = "\n".join(f"- {s}: {outputs[s]}" for s in subtasks) 641 final = ask(synthesizer, f"Question: {request}\n\nFindings:\n{notes}\n\nWrite one answer that cites the [doc ids].") 642 return OrchestratorResult(subtasks, outputs, final) 643 644 645def policy_worker(subtask: str) -> str: 646 """A worker: find the one current policy document that answers the subtask. 647 648 Two things a dense-only lookup gets wrong here, and the fixes real systems use: 649 it happily returns the superseded 2023 travel policy (fix: filter on 650 metadata before ranking), and it confuses PTO with sick leave because both 651 are "time off" (fix: add a keyword-overlap bonus, a tiny hybrid search; 652 see primer.ml.embeddings.retrieval). 653 """ 654 from primer.common.corpus import DOCS 655 from primer.common.embedder import ConceptEmbedder 656 from primer.common.text import tokenize 657 658 current = [d for d in DOCS if "superseded" not in d.title.lower()] 659 emb = ConceptEmbedder() 660 dense = emb.encode([d.title + ". " + d.text for d in current]) @ emb.encode(subtask) 661 words = set(tokenize(subtask)) 662 keyword = np.array([len(words & set(tokenize(d.title + " " + d.text))) for d in current]) 663 best = current[int(np.argmax(dense + 0.2 * keyword))] 664 return f"[{best.id}] {best.text}" 665 666 667def demo_orchestrator() -> ScriptedLLM: 668 return ScriptedLLM(['{"subtasks": ["meal per-diem when travelling", "deadline for expense receipts", "PTO rollover"]}']) 669 670 671def demo_synthesizer() -> ScriptedLLM: 672 """Combines whatever findings it's given, keeping each source id.""" 673 674 def policy(system, messages, tools): 675 findings = [line for line in last_user_text(messages).splitlines() if line.startswith("- ")] 676 parts = [f"{line.split(': ', 1)[0][2:]}: see {re.search(r'\[[a-z]+-\d+\]', line).group(0)}" for line in findings] 677 return "Here's what the policies say. " + "; ".join(parts) + "." 678 679 return ScriptedLLM(policy) 680 681 682# --------------------------------------------------------------------------- 683# Evaluator-optimizer 684# --------------------------------------------------------------------------- 685 686 687def evaluate_optimize(generator: LLM, evaluator: LLM, task: str, max_rounds: int = 3) -> tuple[str, int, bool]: 688 """Writer drafts, editor critiques, writer revises, until the editor says PASS or rounds run out. 689 690 Returns (last draft, rounds used, passed). Stopping at max_rounds matters: 691 a critic that is never satisfied would otherwise loop forever. 692 """ 693 messages: list[dict[str, Any]] = [{"role": "user", "content": task}] 694 draft = "" 695 for rnd in range(1, max_rounds + 1): 696 draft = generator.complete(system="Write the answer.", messages=messages).text 697 verdict = ask(evaluator, f"Task: {task}\nDraft: {draft}\nReply PASS, or say what to fix.", system="You are a strict editor.") 698 if verdict.strip().upper().startswith("PASS"): 699 return draft, rnd, True 700 messages += [{"role": "assistant", "content": draft}, {"role": "user", "content": f"Editor feedback: {verdict}"}] 701 return draft, max_rounds, False 702 703 704# --------------------------------------------------------------------------- 705# Multi-agent: a supervisor delegating to specialists 706# --------------------------------------------------------------------------- 707 708 709@dataclass 710class SupervisorResult: 711 delegations: list[tuple[str, str]] 712 answers: dict[str, str] 713 final: str 714 715 716def supervise(supervisor: LLM, specialists: dict[str, LLM], request: str) -> SupervisorResult: 717 """A supervisor decides who does what; each specialist works in its own, fresh context. 718 719 Each specialist receives only its task text, not the user's full message and not 720 the other specialists' tasks. That focus is the main benefit of several agents. 721 The cost is that anything the supervisor fails to write into the task is lost at 722 the hand-off. 723 """ 724 plan = json.loads(ask(supervisor, f'Delegate to {sorted(specialists)}. Reply as JSON {{"delegate": [{{"to": ..., "task": ...}}]}}.\n\n{request}')) 725 delegations = [(d["to"], d["task"]) for d in plan["delegate"]] 726 answers: dict[str, str] = {} 727 for name, task in delegations: 728 if name not in specialists: 729 answers[name] = f"No specialist named {name!r}. Available: {', '.join(sorted(specialists))}." 730 continue 731 answers[name] = ask(specialists[name], task) 732 findings = "\n".join(f"- {n}: {a}" for n, a in answers.items()) 733 final = ask(supervisor, f"Question: {request}\n\nSpecialist answers:\n{findings}\n\nWrite the final answer.") 734 return SupervisorResult(delegations, answers, final) 735 736 737# --------------------------------------------------------------------------- 738# An explicit state machine with LLM decisions only at specific nodes 739# --------------------------------------------------------------------------- 740 741TERMINAL = ("DONE", "REJECTED", "NEEDS_REVIEW", "AWAITING_APPROVAL") 742 743 744@dataclass 745class WorkflowRun: 746 state: str 747 history: list[str] 748 data: dict[str, Any] = field(default_factory=dict) 749 750 751class InvoiceWorkflow: 752 """Invoice intake as a state machine. Code owns every transition; the model is 753 consulted at exactly two nodes (classify, extract). 754 755 After every transition the run is saved to a checkpoint file, so a crash 756 resumes from the last completed state instead of starting over, and the 757 model is never asked the same question twice. This is durable execution in 758 miniature. 759 """ 760 761 def __init__(self, llm: LLM, checkpoint: str | Path, post: Callable[[dict[str, Any]], Any], approval_limit: float = 5000.0): 762 self.llm = llm 763 self.checkpoint = Path(checkpoint) 764 self.post = post 765 self.approval_limit = approval_limit 766 767 def _save(self, run: WorkflowRun) -> None: 768 # Write to a temp file then rename: a crash mid-write can't leave a half-written checkpoint. 769 tmp = self.checkpoint.with_suffix(".tmp") 770 tmp.write_text(json.dumps({"state": run.state, "history": run.history, "data": run.data})) 771 os.replace(tmp, self.checkpoint) 772 773 def _load(self) -> WorkflowRun | None: 774 if not self.checkpoint.exists(): 775 return None 776 saved = json.loads(self.checkpoint.read_text()) 777 return WorkflowRun(saved["state"], saved["history"], saved["data"]) 778 779 def _move(self, run: WorkflowRun, state: str) -> None: 780 run.state = state 781 run.history.append(state) 782 self._save(run) 783 784 def run(self, document: str) -> WorkflowRun: 785 run = self._load() or WorkflowRun("RECEIVED", ["RECEIVED"], {"document": document}) 786 while run.state not in TERMINAL: 787 if run.state == "RECEIVED": 788 # LLM node 1: a judgement call code can't easily make. 789 kind = ask(self.llm, f"Classify this document as 'invoice' or 'other':\n{run.data['document']}").strip().lower() 790 self._move(run, "CLASSIFIED" if kind == "invoice" else "REJECTED") 791 elif run.state == "CLASSIFIED": 792 # LLM node 2: pull structured fields out of messy text. 793 run.data["invoice"] = json.loads(ask(self.llm, f"Extract invoice_no, vendor, amount as JSON:\n{run.data['document']}")) 794 self._move(run, "EXTRACTED") 795 elif run.state == "EXTRACTED": 796 problem = validate_invoice(run.data["invoice"]) 797 if problem: 798 run.data["problem"] = problem 799 self._move(run, "NEEDS_REVIEW") 800 else: 801 self._move(run, "VALIDATED") 802 elif run.state == "VALIDATED": 803 if run.data["invoice"]["amount"] > self.approval_limit: 804 self._move(run, "AWAITING_APPROVAL") 805 else: 806 # If this raises, the checkpoint still says VALIDATED and a rerun posts again. 807 # A real ledger call would carry invoice_no as an idempotency key. 808 self.post(run.data["invoice"]) 809 self._move(run, "POSTED") 810 elif run.state == "POSTED": 811 self._move(run, "DONE") 812 return run 813 814 815 def approve(self, approver: str) -> WorkflowRun: 816 """A person approves a paused invoice; the run continues from the checkpoint.""" 817 run = self._load() 818 if run is None or run.state != "AWAITING_APPROVAL": 819 raise ValueError("nothing is waiting for approval") 820 run.data["approved_by"] = approver 821 self.post(run.data["invoice"]) 822 self._move(run, "POSTED") 823 self._move(run, "DONE") 824 return run 825 826 827def validate_invoice(invoice: dict[str, Any]) -> str | None: 828 if not re.fullmatch(r"INV-\d+", str(invoice.get("invoice_no", ""))): 829 return f"invoice_no must look like INV-1234, got {invoice.get('invoice_no')!r}" 830 if not isinstance(invoice.get("amount"), (int, float)) or invoice["amount"] <= 0: 831 return f"amount must be positive, got {invoice.get('amount')!r}" 832 return None 833 834 835def demo_document_llm() -> ScriptedLLM: 836 """A stand-in model for the two LLM nodes: classification and extraction.""" 837 838 def policy(system, messages, tools): 839 prompt = last_user_text(messages) 840 document = prompt.split("\n", 1)[1] 841 if prompt.startswith("Classify"): 842 return "invoice" if document.startswith("INVOICE") else "other" 843 m = re.search(r"INVOICE (\S+) from (.+?)\. Amount due: (-?[\d.]+)", document) 844 return json.dumps({"invoice_no": m.group(1), "vendor": m.group(2), "amount": float(m.group(3))}) 845 846 return ScriptedLLM(policy) 847 848 849 850# --------------------------------------------------------------------------- 851# One question, four levels of autonomy 852# --------------------------------------------------------------------------- 853 854AUTONOMY_QUESTION = "How many PTO days do I get, and what is the meal per-diem when travelling?" 855 856 857def autonomy_costs() -> list[dict[str, Any]]: 858 """Answer the same two-part question at each rung of the autonomy ladder; count calls and tokens.""" 859 from primer.agents.agent_loop import DEMO_TOOL_DEFS, DEMO_TOOLS, run_agent, search_kb 860 from primer.agents.llm import ToolCall, tool_results 861 862 rows = [] 863 864 # 1. Fixed workflow: code decides the steps (rewrite the question as a search, then answer). 865 llm = ScriptedLLM(["PTO days per year; meal per-diem", "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]) 866 run_chain(llm, AUTONOMY_QUESTION, [ChainStep("query", "Rewrite as a search query: {input}"), ChainStep("answer", "Answer using: {input}")]) 867 rows.append({"design": "fixed workflow", "llm": llm}) 868 869 # 2. Router: the model picks a path; the chosen handler makes one more call. 870 llm = ScriptedLLM(["policy", "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]) 871 route(llm, {"policy": lambda q: ask(llm, q), "general": lambda q: ask(llm, q)}, AUTONOMY_QUESTION) 872 rows.append({"design": "router", "llm": llm}) 873 874 # 3. Agent loop: the model chooses tools (two searches in parallel), then answers. 875 def agent_policy(system, messages, tools): 876 if not tool_results(messages): 877 return [ToolCall("", "search_kb", {"query": "PTO days"}), ToolCall("", "search_kb", {"query": "meal per-diem"})] 878 return "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]." 879 880 llm = ScriptedLLM(agent_policy) 881 run_agent(llm, AUTONOMY_QUESTION, DEMO_TOOLS, DEMO_TOOL_DEFS) 882 rows.append({"design": "agent loop", "llm": llm}) 883 884 # 4. Multi-agent: a supervisor plus two specialists, each with its own context. 885 sup = ScriptedLLM([ 886 '{"delegate": [{"to": "hr", "task": "PTO days per year?"}, {"to": "finance", "task": "Meal per-diem?"}]}', 887 "20 PTO days a year [hr-001]; 75 dollars a day [fin-002].", 888 ]) 889 hr, fin = ScriptedLLM([search_kb("PTO days")]), ScriptedLLM([search_kb("meal per-diem")]) 890 supervise(sup, {"hr": hr, "finance": fin}, AUTONOMY_QUESTION) 891 rows.append({"design": "multi-agent", "llm": sup, "extra": [hr, fin]}) 892 893 out = [] 894 for r in rows: 895 models = [r["llm"], *r.get("extra", [])] 896 out.append({ 897 "design": r["design"], 898 "calls": sum(m.calls for m in models), 899 "input_tokens": sum(m.total_usage.input_tokens for m in models), 900 "output_tokens": sum(m.total_usage.output_tokens for m in models), 901 }) 902 return out 903 904 905def resume_comparison(tmp_dir: str | Path) -> dict[str, int]: 906 """Model calls spent finishing one invoice when posting crashes once: resume vs. restart.""" 907 results = {} 908 for mode in ("resume from checkpoint", "restart from scratch"): 909 llm = demo_document_llm() 910 attempts: list[dict] = [] 911 912 def flaky_post(invoice: dict[str, Any]) -> None: 913 attempts.append(invoice) 914 if len(attempts) == 1: 915 raise ConnectionError("ledger unavailable") 916 917 ckpt = Path(tmp_dir) / f"{mode.split()[0]}.json" 918 wf = InvoiceWorkflow(llm, ckpt, post=flaky_post) 919 try: 920 wf.run("INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.") 921 except ConnectionError: 922 if mode == "restart from scratch": 923 ckpt.unlink() # no durable state: the second attempt starts over 924 wf.run("INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.") 925 results[mode] = llm.calls 926 return results 927 928 929# --------------------------------------------------------------------------- 930# Figures and walkthrough 931# --------------------------------------------------------------------------- 932 933 934def figures() -> dict: 935 """Plots computed from this lesson's own code (matplotlib imported here).""" 936 import tempfile 937 938 import matplotlib 939 940 matplotlib.use("Agg") 941 import matplotlib.pyplot as plt 942 943 figs = {} 944 945 rows = autonomy_costs() 946 names = [r["design"] for r in rows] 947 fig, (ax1, ax2) = plt.subplots(1, 2, figsize=(9, 3.8)) 948 ax1.bar(names, [r["calls"] for r in rows], color="#5b8fd6") 949 ax1.set(ylabel="model calls", title="Calls for the same question") 950 ax2.bar(names, [r["input_tokens"] + r["output_tokens"] for r in rows], color="#d98c3a") 951 ax2.set(ylabel="tokens (input + output)", title="Tokens for the same question") 952 for ax in (ax1, ax2): 953 ax.tick_params(axis="x", rotation=20) 954 fig.tight_layout() 955 figs["autonomy_costs"] = fig 956 957 voters = list(range(1, 16, 2)) 958 fig, ax = plt.subplots(figsize=(7, 4)) 959 for p in (0.6, 0.7, 0.8, 0.9): 960 ax.plot(voters, [majority_accuracy(p, n) for n in voters], "o-", ms=4, label=f"each voter right {p:.0%}") 961 ax.set(xlabel="number of independent voters (odd)", ylabel="P(majority is right)", ylim=(0.5, 1.01), 962 title="Voting: many noisy judges beat one (if their errors are independent)") 963 ax.legend() 964 figs["voting"] = fig 965 966 with tempfile.TemporaryDirectory() as d: 967 costs = resume_comparison(d) 968 fig, ax = plt.subplots(figsize=(6, 3.8)) 969 bars = ax.bar(list(costs), list(costs.values()), color=["#6bb36b", "#c0392b"]) 970 ax.bar_label(bars) 971 ax.set(ylabel="model calls to finish one invoice", title="Posting crashed once: what did it cost?") 972 figs["resume_cost"] = fig 973 return figs 974 975 976def demo() -> None: 977 import tempfile 978 979 from primer._show import banner, say, table, takeaway 980 981 banner("1. One question, four levels of autonomy") 982 table(["design", "model calls", "input tokens", "output tokens"], 983 [(r["design"], r["calls"], r["input_tokens"], r["output_tokens"]) for r in autonomy_costs()]) 984 takeaway("Use the least autonomy that solves the problem. Every rung up costs calls, tokens and predictability.") 985 986 banner("2. Prompt chaining with a gate") 987 gate = lambda text: None if text.count("- ") >= 3 else "outline needs at least 3 points" # noqa: E731 988 for outline in ("- scope\n- risks\n- timeline", "- scope"): 989 llm = ScriptedLLM([outline, "A draft that follows the outline."]) 990 r = run_chain(llm, "Q3 plan", [ChainStep("outline", "Outline: {input}", gate), ChainStep("draft", "Draft from: {input}")]) 991 print(f" outline {outline!r:32} -> stopped_at={r.stopped_at}, calls={llm.calls}, reason={r.reason}") 992 print() 993 994 banner("3. Routing (with a fallback for labels the model invents)") 995 handlers = {"billing": lambda q: "billing team", "technical": lambda q: "tech support", "general": lambda q: "front desk"} 996 for raw in ("billing", " Technical\n", "Billing-ish, maybe refunds?"): 997 print(f" classifier said {raw!r:32} -> {route(ScriptedLLM([raw]), handlers, 'q')[0]}") 998 print() 999 1000 banner("4. Parallelization: voting") 1001 reviewers = [ScriptedLLM(["vulnerable"]), ScriptedLLM(["safe"]), ScriptedLLM(["vulnerable"])] 1002 print(f" three reviewers -> {vote(reviewers, 'Is this SQL safe?')}") 1003 table(["voters", "p=0.7", "p=0.8"], [(n, majority_accuracy(0.7, n), majority_accuracy(0.8, n)) for n in (1, 3, 5, 9)], floatfmt=".3f") 1004 1005 banner("5. Orchestrator-workers") 1006 r = orchestrate(demo_orchestrator(), policy_worker, demo_synthesizer(), 1007 "What's the meal per-diem, when are receipts due, and does PTO roll over?") 1008 for s in r.subtasks: 1009 print(f" worker[{s}] -> {r.worker_outputs[s][:70]}...") 1010 say(f"Final: {r.final}") 1011 1012 banner("6. Evaluator-optimizer") 1013 writer = ScriptedLLM(["PTO rolls over.", "Up to 5 unused PTO days roll over [hr-001]."]) 1014 editor = ScriptedLLM(["Missing a citation: cite the policy id in brackets.", "PASS"]) 1015 say(f"(draft, rounds, passed) = {evaluate_optimize(writer, editor, 'Does PTO roll over?')}") 1016 1017 banner("7. Multi-agent: supervisor and specialists") 1018 res = supervise( 1019 ScriptedLLM(['{"delegate": [{"to": "hr", "task": "PTO days per year?"}, {"to": "finance", "task": "Meal per-diem?"}]}', 1020 "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]), 1021 {"hr": ScriptedLLM(["20 days [hr-001]"]), "finance": ScriptedLLM(["75 dollars a day [fin-002]"])}, 1022 "PTO and per-diem?", 1023 ) 1024 for (who, task), answer in zip(res.delegations, res.answers.values()): 1025 print(f" {who:8} got only: {task!r:24} -> {answer}") 1026 say(f"Final: {res.final}") 1027 1028 banner("8. A state machine with LLM decisions at two nodes, and checkpoints") 1029 with tempfile.TemporaryDirectory() as d: 1030 for i, doc in enumerate(["INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.", 1031 "INVOICE INV-2042 from Globex. Amount due: 18000.00 USD.", 1032 "Acme Supplies newsletter: our autumn catalogue is here!"]): 1033 run = InvoiceWorkflow(demo_document_llm(), Path(d) / f"c{i}.json", post=lambda inv: None).run(doc) 1034 print(f" {doc[:40]:42} -> {' > '.join(run.history)}") 1035 print() 1036 say(f"Posting crashes once. Model calls to finish: {resume_comparison(d)}") 1037 takeaway("Code owns the transitions; the model decides only where judgement is needed. Checkpoints make crashes cheap.") 1038 1039 1040if __name__ == "__main__": 1041 demo()
538def ask(llm: LLM, prompt: str, system: str = "") -> str: 539 """One model call with one user message; returns the text.""" 540 return llm.complete(system=system, messages=[{"role": "user", "content": prompt}]).text
One model call with one user message; returns the text.
548@dataclass 549class ChainStep: 550 name: str 551 prompt: str # template; "{input}" is replaced by the previous step's output 552 gate: Callable[[str], str | None] | None = None # returns a problem, or None to continue
555@dataclass 556class ChainResult: 557 outputs: dict[str, str] 558 stopped_at: str | None = None 559 reason: str | None = None
562def run_chain(llm: LLM, task: str, steps: list[ChainStep]) -> ChainResult: 563 """Fixed sequence of calls; each output feeds the next, with a code check (gate) in between.""" 564 outputs: dict[str, str] = {} 565 current = task 566 for step in steps: 567 current = ask(llm, step.prompt.format(input=current)) 568 outputs[step.name] = current 569 problem = step.gate(current) if step.gate else None 570 if problem: 571 return ChainResult(outputs, step.name, problem) 572 return ChainResult(outputs)
Fixed sequence of calls; each output feeds the next, with a code check (gate) in between.
580def route( 581 classifier: LLM, handlers: dict[str, Callable[[str], str]], request: str, fallback: str = "general" 582) -> tuple[str, str]: 583 """Classify the request with one model call, then hand it to that label's handler. 584 585 The label comes back as free text, so it's normalized and checked against 586 the known labels. Anything else goes to the fallback instead of crashing 587 or guessing. 588 """ 589 raw = ask(classifier, f"Classify into one of {sorted(handlers)}. Reply with the label only.\n\n{request}") 590 label = raw.strip().lower() 591 if label not in handlers: 592 label = fallback 593 return label, handlers[label](request)
Classify the request with one model call, then hand it to that label's handler.
The label comes back as free text, so it's normalized and checked against the known labels. Anything else goes to the fallback instead of crashing or guessing.
601def run_sections(sections: dict[str, Callable[[], str]]) -> dict[str, str]: 602 """Run independent pieces of work at the same time; wall-clock time is the slowest piece.""" 603 with ThreadPoolExecutor(max_workers=len(sections)) as pool: 604 futures = {name: pool.submit(fn) for name, fn in sections.items()} 605 return {name: f.result() for name, f in futures.items()}
Run independent pieces of work at the same time; wall-clock time is the slowest piece.
608def vote(reviewers: list[LLM], prompt: str) -> tuple[str, dict[str, int]]: 609 """Ask several models (or one model several times) the same question; take the majority.""" 610 answers = run_sections({str(i): (lambda r=r: ask(r, prompt).strip().lower()) for i, r in enumerate(reviewers)}) 611 tally = Counter(answers.values()) 612 return tally.most_common(1)[0][0], dict(tally)
Ask several models (or one model several times) the same question; take the majority.
615def majority_accuracy(p: float, n_voters: int) -> float: 616 """Chance a majority of n independent voters, each right with probability p, is right.""" 617 return sum(math.comb(n_voters, k) * p**k * (1 - p) ** (n_voters - k) for k in range(n_voters // 2 + 1, n_voters + 1))
Chance a majority of n independent voters, each right with probability p, is right.
625@dataclass 626class OrchestratorResult: 627 subtasks: list[str] 628 worker_outputs: dict[str, str] 629 final: str
632def orchestrate(orchestrator: LLM, worker: Callable[[str], str], synthesizer: LLM, request: str) -> OrchestratorResult: 633 """One model splits the request into subtasks, workers handle them in parallel, one model combines. 634 635 Unlike sectioning, the subtasks aren't known in advance: the orchestrator 636 decides them from the request. 637 """ 638 plan = ask(orchestrator, f'Split into independent subtasks. Reply as JSON {{"subtasks": [...]}}.\n\n{request}') 639 subtasks = list(json.loads(plan)["subtasks"]) 640 outputs = run_sections({s: (lambda s=s: worker(s)) for s in subtasks}) 641 notes = "\n".join(f"- {s}: {outputs[s]}" for s in subtasks) 642 final = ask(synthesizer, f"Question: {request}\n\nFindings:\n{notes}\n\nWrite one answer that cites the [doc ids].") 643 return OrchestratorResult(subtasks, outputs, final)
One model splits the request into subtasks, workers handle them in parallel, one model combines.
Unlike sectioning, the subtasks aren't known in advance: the orchestrator decides them from the request.
646def policy_worker(subtask: str) -> str: 647 """A worker: find the one current policy document that answers the subtask. 648 649 Two things a dense-only lookup gets wrong here, and the fixes real systems use: 650 it happily returns the superseded 2023 travel policy (fix: filter on 651 metadata before ranking), and it confuses PTO with sick leave because both 652 are "time off" (fix: add a keyword-overlap bonus, a tiny hybrid search; 653 see primer.ml.embeddings.retrieval). 654 """ 655 from primer.common.corpus import DOCS 656 from primer.common.embedder import ConceptEmbedder 657 from primer.common.text import tokenize 658 659 current = [d for d in DOCS if "superseded" not in d.title.lower()] 660 emb = ConceptEmbedder() 661 dense = emb.encode([d.title + ". " + d.text for d in current]) @ emb.encode(subtask) 662 words = set(tokenize(subtask)) 663 keyword = np.array([len(words & set(tokenize(d.title + " " + d.text))) for d in current]) 664 best = current[int(np.argmax(dense + 0.2 * keyword))] 665 return f"[{best.id}] {best.text}"
A worker: find the one current policy document that answers the subtask.
Two things a dense-only lookup gets wrong here, and the fixes real systems use: it happily returns the superseded 2023 travel policy (fix: filter on metadata before ranking), and it confuses PTO with sick leave because both are "time off" (fix: add a keyword-overlap bonus, a tiny hybrid search; see primer.ml.embeddings.retrieval).
672def demo_synthesizer() -> ScriptedLLM: 673 """Combines whatever findings it's given, keeping each source id.""" 674 675 def policy(system, messages, tools): 676 findings = [line for line in last_user_text(messages).splitlines() if line.startswith("- ")] 677 parts = [f"{line.split(': ', 1)[0][2:]}: see {re.search(r'\[[a-z]+-\d+\]', line).group(0)}" for line in findings] 678 return "Here's what the policies say. " + "; ".join(parts) + "." 679 680 return ScriptedLLM(policy)
Combines whatever findings it's given, keeping each source id.
688def evaluate_optimize(generator: LLM, evaluator: LLM, task: str, max_rounds: int = 3) -> tuple[str, int, bool]: 689 """Writer drafts, editor critiques, writer revises, until the editor says PASS or rounds run out. 690 691 Returns (last draft, rounds used, passed). Stopping at max_rounds matters: 692 a critic that is never satisfied would otherwise loop forever. 693 """ 694 messages: list[dict[str, Any]] = [{"role": "user", "content": task}] 695 draft = "" 696 for rnd in range(1, max_rounds + 1): 697 draft = generator.complete(system="Write the answer.", messages=messages).text 698 verdict = ask(evaluator, f"Task: {task}\nDraft: {draft}\nReply PASS, or say what to fix.", system="You are a strict editor.") 699 if verdict.strip().upper().startswith("PASS"): 700 return draft, rnd, True 701 messages += [{"role": "assistant", "content": draft}, {"role": "user", "content": f"Editor feedback: {verdict}"}] 702 return draft, max_rounds, False
Writer drafts, editor critiques, writer revises, until the editor says PASS or rounds run out.
Returns (last draft, rounds used, passed). Stopping at max_rounds matters: a critic that is never satisfied would otherwise loop forever.
710@dataclass 711class SupervisorResult: 712 delegations: list[tuple[str, str]] 713 answers: dict[str, str] 714 final: str
717def supervise(supervisor: LLM, specialists: dict[str, LLM], request: str) -> SupervisorResult: 718 """A supervisor decides who does what; each specialist works in its own, fresh context. 719 720 Each specialist receives only its task text, not the user's full message and not 721 the other specialists' tasks. That focus is the main benefit of several agents. 722 The cost is that anything the supervisor fails to write into the task is lost at 723 the hand-off. 724 """ 725 plan = json.loads(ask(supervisor, f'Delegate to {sorted(specialists)}. Reply as JSON {{"delegate": [{{"to": ..., "task": ...}}]}}.\n\n{request}')) 726 delegations = [(d["to"], d["task"]) for d in plan["delegate"]] 727 answers: dict[str, str] = {} 728 for name, task in delegations: 729 if name not in specialists: 730 answers[name] = f"No specialist named {name!r}. Available: {', '.join(sorted(specialists))}." 731 continue 732 answers[name] = ask(specialists[name], task) 733 findings = "\n".join(f"- {n}: {a}" for n, a in answers.items()) 734 final = ask(supervisor, f"Question: {request}\n\nSpecialist answers:\n{findings}\n\nWrite the final answer.") 735 return SupervisorResult(delegations, answers, final)
A supervisor decides who does what; each specialist works in its own, fresh context.
Each specialist receives only its task text, not the user's full message and not the other specialists' tasks. That focus is the main benefit of several agents. The cost is that anything the supervisor fails to write into the task is lost at the hand-off.
745@dataclass 746class WorkflowRun: 747 state: str 748 history: list[str] 749 data: dict[str, Any] = field(default_factory=dict)
752class InvoiceWorkflow: 753 """Invoice intake as a state machine. Code owns every transition; the model is 754 consulted at exactly two nodes (classify, extract). 755 756 After every transition the run is saved to a checkpoint file, so a crash 757 resumes from the last completed state instead of starting over, and the 758 model is never asked the same question twice. This is durable execution in 759 miniature. 760 """ 761 762 def __init__(self, llm: LLM, checkpoint: str | Path, post: Callable[[dict[str, Any]], Any], approval_limit: float = 5000.0): 763 self.llm = llm 764 self.checkpoint = Path(checkpoint) 765 self.post = post 766 self.approval_limit = approval_limit 767 768 def _save(self, run: WorkflowRun) -> None: 769 # Write to a temp file then rename: a crash mid-write can't leave a half-written checkpoint. 770 tmp = self.checkpoint.with_suffix(".tmp") 771 tmp.write_text(json.dumps({"state": run.state, "history": run.history, "data": run.data})) 772 os.replace(tmp, self.checkpoint) 773 774 def _load(self) -> WorkflowRun | None: 775 if not self.checkpoint.exists(): 776 return None 777 saved = json.loads(self.checkpoint.read_text()) 778 return WorkflowRun(saved["state"], saved["history"], saved["data"]) 779 780 def _move(self, run: WorkflowRun, state: str) -> None: 781 run.state = state 782 run.history.append(state) 783 self._save(run) 784 785 def run(self, document: str) -> WorkflowRun: 786 run = self._load() or WorkflowRun("RECEIVED", ["RECEIVED"], {"document": document}) 787 while run.state not in TERMINAL: 788 if run.state == "RECEIVED": 789 # LLM node 1: a judgement call code can't easily make. 790 kind = ask(self.llm, f"Classify this document as 'invoice' or 'other':\n{run.data['document']}").strip().lower() 791 self._move(run, "CLASSIFIED" if kind == "invoice" else "REJECTED") 792 elif run.state == "CLASSIFIED": 793 # LLM node 2: pull structured fields out of messy text. 794 run.data["invoice"] = json.loads(ask(self.llm, f"Extract invoice_no, vendor, amount as JSON:\n{run.data['document']}")) 795 self._move(run, "EXTRACTED") 796 elif run.state == "EXTRACTED": 797 problem = validate_invoice(run.data["invoice"]) 798 if problem: 799 run.data["problem"] = problem 800 self._move(run, "NEEDS_REVIEW") 801 else: 802 self._move(run, "VALIDATED") 803 elif run.state == "VALIDATED": 804 if run.data["invoice"]["amount"] > self.approval_limit: 805 self._move(run, "AWAITING_APPROVAL") 806 else: 807 # If this raises, the checkpoint still says VALIDATED and a rerun posts again. 808 # A real ledger call would carry invoice_no as an idempotency key. 809 self.post(run.data["invoice"]) 810 self._move(run, "POSTED") 811 elif run.state == "POSTED": 812 self._move(run, "DONE") 813 return run 814 815 816 def approve(self, approver: str) -> WorkflowRun: 817 """A person approves a paused invoice; the run continues from the checkpoint.""" 818 run = self._load() 819 if run is None or run.state != "AWAITING_APPROVAL": 820 raise ValueError("nothing is waiting for approval") 821 run.data["approved_by"] = approver 822 self.post(run.data["invoice"]) 823 self._move(run, "POSTED") 824 self._move(run, "DONE") 825 return run
Invoice intake as a state machine. Code owns every transition; the model is consulted at exactly two nodes (classify, extract).
After every transition the run is saved to a checkpoint file, so a crash resumes from the last completed state instead of starting over, and the model is never asked the same question twice. This is durable execution in miniature.
785 def run(self, document: str) -> WorkflowRun: 786 run = self._load() or WorkflowRun("RECEIVED", ["RECEIVED"], {"document": document}) 787 while run.state not in TERMINAL: 788 if run.state == "RECEIVED": 789 # LLM node 1: a judgement call code can't easily make. 790 kind = ask(self.llm, f"Classify this document as 'invoice' or 'other':\n{run.data['document']}").strip().lower() 791 self._move(run, "CLASSIFIED" if kind == "invoice" else "REJECTED") 792 elif run.state == "CLASSIFIED": 793 # LLM node 2: pull structured fields out of messy text. 794 run.data["invoice"] = json.loads(ask(self.llm, f"Extract invoice_no, vendor, amount as JSON:\n{run.data['document']}")) 795 self._move(run, "EXTRACTED") 796 elif run.state == "EXTRACTED": 797 problem = validate_invoice(run.data["invoice"]) 798 if problem: 799 run.data["problem"] = problem 800 self._move(run, "NEEDS_REVIEW") 801 else: 802 self._move(run, "VALIDATED") 803 elif run.state == "VALIDATED": 804 if run.data["invoice"]["amount"] > self.approval_limit: 805 self._move(run, "AWAITING_APPROVAL") 806 else: 807 # If this raises, the checkpoint still says VALIDATED and a rerun posts again. 808 # A real ledger call would carry invoice_no as an idempotency key. 809 self.post(run.data["invoice"]) 810 self._move(run, "POSTED") 811 elif run.state == "POSTED": 812 self._move(run, "DONE") 813 return run
816 def approve(self, approver: str) -> WorkflowRun: 817 """A person approves a paused invoice; the run continues from the checkpoint.""" 818 run = self._load() 819 if run is None or run.state != "AWAITING_APPROVAL": 820 raise ValueError("nothing is waiting for approval") 821 run.data["approved_by"] = approver 822 self.post(run.data["invoice"]) 823 self._move(run, "POSTED") 824 self._move(run, "DONE") 825 return run
A person approves a paused invoice; the run continues from the checkpoint.
828def validate_invoice(invoice: dict[str, Any]) -> str | None: 829 if not re.fullmatch(r"INV-\d+", str(invoice.get("invoice_no", ""))): 830 return f"invoice_no must look like INV-1234, got {invoice.get('invoice_no')!r}" 831 if not isinstance(invoice.get("amount"), (int, float)) or invoice["amount"] <= 0: 832 return f"amount must be positive, got {invoice.get('amount')!r}" 833 return None
836def demo_document_llm() -> ScriptedLLM: 837 """A stand-in model for the two LLM nodes: classification and extraction.""" 838 839 def policy(system, messages, tools): 840 prompt = last_user_text(messages) 841 document = prompt.split("\n", 1)[1] 842 if prompt.startswith("Classify"): 843 return "invoice" if document.startswith("INVOICE") else "other" 844 m = re.search(r"INVOICE (\S+) from (.+?)\. Amount due: (-?[\d.]+)", document) 845 return json.dumps({"invoice_no": m.group(1), "vendor": m.group(2), "amount": float(m.group(3))}) 846 847 return ScriptedLLM(policy)
A stand-in model for the two LLM nodes: classification and extraction.
858def autonomy_costs() -> list[dict[str, Any]]: 859 """Answer the same two-part question at each rung of the autonomy ladder; count calls and tokens.""" 860 from primer.agents.agent_loop import DEMO_TOOL_DEFS, DEMO_TOOLS, run_agent, search_kb 861 from primer.agents.llm import ToolCall, tool_results 862 863 rows = [] 864 865 # 1. Fixed workflow: code decides the steps (rewrite the question as a search, then answer). 866 llm = ScriptedLLM(["PTO days per year; meal per-diem", "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]) 867 run_chain(llm, AUTONOMY_QUESTION, [ChainStep("query", "Rewrite as a search query: {input}"), ChainStep("answer", "Answer using: {input}")]) 868 rows.append({"design": "fixed workflow", "llm": llm}) 869 870 # 2. Router: the model picks a path; the chosen handler makes one more call. 871 llm = ScriptedLLM(["policy", "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]) 872 route(llm, {"policy": lambda q: ask(llm, q), "general": lambda q: ask(llm, q)}, AUTONOMY_QUESTION) 873 rows.append({"design": "router", "llm": llm}) 874 875 # 3. Agent loop: the model chooses tools (two searches in parallel), then answers. 876 def agent_policy(system, messages, tools): 877 if not tool_results(messages): 878 return [ToolCall("", "search_kb", {"query": "PTO days"}), ToolCall("", "search_kb", {"query": "meal per-diem"})] 879 return "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]." 880 881 llm = ScriptedLLM(agent_policy) 882 run_agent(llm, AUTONOMY_QUESTION, DEMO_TOOLS, DEMO_TOOL_DEFS) 883 rows.append({"design": "agent loop", "llm": llm}) 884 885 # 4. Multi-agent: a supervisor plus two specialists, each with its own context. 886 sup = ScriptedLLM([ 887 '{"delegate": [{"to": "hr", "task": "PTO days per year?"}, {"to": "finance", "task": "Meal per-diem?"}]}', 888 "20 PTO days a year [hr-001]; 75 dollars a day [fin-002].", 889 ]) 890 hr, fin = ScriptedLLM([search_kb("PTO days")]), ScriptedLLM([search_kb("meal per-diem")]) 891 supervise(sup, {"hr": hr, "finance": fin}, AUTONOMY_QUESTION) 892 rows.append({"design": "multi-agent", "llm": sup, "extra": [hr, fin]}) 893 894 out = [] 895 for r in rows: 896 models = [r["llm"], *r.get("extra", [])] 897 out.append({ 898 "design": r["design"], 899 "calls": sum(m.calls for m in models), 900 "input_tokens": sum(m.total_usage.input_tokens for m in models), 901 "output_tokens": sum(m.total_usage.output_tokens for m in models), 902 }) 903 return out
Answer the same two-part question at each rung of the autonomy ladder; count calls and tokens.
906def resume_comparison(tmp_dir: str | Path) -> dict[str, int]: 907 """Model calls spent finishing one invoice when posting crashes once: resume vs. restart.""" 908 results = {} 909 for mode in ("resume from checkpoint", "restart from scratch"): 910 llm = demo_document_llm() 911 attempts: list[dict] = [] 912 913 def flaky_post(invoice: dict[str, Any]) -> None: 914 attempts.append(invoice) 915 if len(attempts) == 1: 916 raise ConnectionError("ledger unavailable") 917 918 ckpt = Path(tmp_dir) / f"{mode.split()[0]}.json" 919 wf = InvoiceWorkflow(llm, ckpt, post=flaky_post) 920 try: 921 wf.run("INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.") 922 except ConnectionError: 923 if mode == "restart from scratch": 924 ckpt.unlink() # no durable state: the second attempt starts over 925 wf.run("INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.") 926 results[mode] = llm.calls 927 return results
Model calls spent finishing one invoice when posting crashes once: resume vs. restart.
935def figures() -> dict: 936 """Plots computed from this lesson's own code (matplotlib imported here).""" 937 import tempfile 938 939 import matplotlib 940 941 matplotlib.use("Agg") 942 import matplotlib.pyplot as plt 943 944 figs = {} 945 946 rows = autonomy_costs() 947 names = [r["design"] for r in rows] 948 fig, (ax1, ax2) = plt.subplots(1, 2, figsize=(9, 3.8)) 949 ax1.bar(names, [r["calls"] for r in rows], color="#5b8fd6") 950 ax1.set(ylabel="model calls", title="Calls for the same question") 951 ax2.bar(names, [r["input_tokens"] + r["output_tokens"] for r in rows], color="#d98c3a") 952 ax2.set(ylabel="tokens (input + output)", title="Tokens for the same question") 953 for ax in (ax1, ax2): 954 ax.tick_params(axis="x", rotation=20) 955 fig.tight_layout() 956 figs["autonomy_costs"] = fig 957 958 voters = list(range(1, 16, 2)) 959 fig, ax = plt.subplots(figsize=(7, 4)) 960 for p in (0.6, 0.7, 0.8, 0.9): 961 ax.plot(voters, [majority_accuracy(p, n) for n in voters], "o-", ms=4, label=f"each voter right {p:.0%}") 962 ax.set(xlabel="number of independent voters (odd)", ylabel="P(majority is right)", ylim=(0.5, 1.01), 963 title="Voting: many noisy judges beat one (if their errors are independent)") 964 ax.legend() 965 figs["voting"] = fig 966 967 with tempfile.TemporaryDirectory() as d: 968 costs = resume_comparison(d) 969 fig, ax = plt.subplots(figsize=(6, 3.8)) 970 bars = ax.bar(list(costs), list(costs.values()), color=["#6bb36b", "#c0392b"]) 971 ax.bar_label(bars) 972 ax.set(ylabel="model calls to finish one invoice", title="Posting crashed once: what did it cost?") 973 figs["resume_cost"] = fig 974 return figs
Plots computed from this lesson's own code (matplotlib imported here).
977def demo() -> None: 978 import tempfile 979 980 from primer._show import banner, say, table, takeaway 981 982 banner("1. One question, four levels of autonomy") 983 table(["design", "model calls", "input tokens", "output tokens"], 984 [(r["design"], r["calls"], r["input_tokens"], r["output_tokens"]) for r in autonomy_costs()]) 985 takeaway("Use the least autonomy that solves the problem. Every rung up costs calls, tokens and predictability.") 986 987 banner("2. Prompt chaining with a gate") 988 gate = lambda text: None if text.count("- ") >= 3 else "outline needs at least 3 points" # noqa: E731 989 for outline in ("- scope\n- risks\n- timeline", "- scope"): 990 llm = ScriptedLLM([outline, "A draft that follows the outline."]) 991 r = run_chain(llm, "Q3 plan", [ChainStep("outline", "Outline: {input}", gate), ChainStep("draft", "Draft from: {input}")]) 992 print(f" outline {outline!r:32} -> stopped_at={r.stopped_at}, calls={llm.calls}, reason={r.reason}") 993 print() 994 995 banner("3. Routing (with a fallback for labels the model invents)") 996 handlers = {"billing": lambda q: "billing team", "technical": lambda q: "tech support", "general": lambda q: "front desk"} 997 for raw in ("billing", " Technical\n", "Billing-ish, maybe refunds?"): 998 print(f" classifier said {raw!r:32} -> {route(ScriptedLLM([raw]), handlers, 'q')[0]}") 999 print() 1000 1001 banner("4. Parallelization: voting") 1002 reviewers = [ScriptedLLM(["vulnerable"]), ScriptedLLM(["safe"]), ScriptedLLM(["vulnerable"])] 1003 print(f" three reviewers -> {vote(reviewers, 'Is this SQL safe?')}") 1004 table(["voters", "p=0.7", "p=0.8"], [(n, majority_accuracy(0.7, n), majority_accuracy(0.8, n)) for n in (1, 3, 5, 9)], floatfmt=".3f") 1005 1006 banner("5. Orchestrator-workers") 1007 r = orchestrate(demo_orchestrator(), policy_worker, demo_synthesizer(), 1008 "What's the meal per-diem, when are receipts due, and does PTO roll over?") 1009 for s in r.subtasks: 1010 print(f" worker[{s}] -> {r.worker_outputs[s][:70]}...") 1011 say(f"Final: {r.final}") 1012 1013 banner("6. Evaluator-optimizer") 1014 writer = ScriptedLLM(["PTO rolls over.", "Up to 5 unused PTO days roll over [hr-001]."]) 1015 editor = ScriptedLLM(["Missing a citation: cite the policy id in brackets.", "PASS"]) 1016 say(f"(draft, rounds, passed) = {evaluate_optimize(writer, editor, 'Does PTO roll over?')}") 1017 1018 banner("7. Multi-agent: supervisor and specialists") 1019 res = supervise( 1020 ScriptedLLM(['{"delegate": [{"to": "hr", "task": "PTO days per year?"}, {"to": "finance", "task": "Meal per-diem?"}]}', 1021 "20 PTO days a year [hr-001]; 75 dollars a day [fin-002]."]), 1022 {"hr": ScriptedLLM(["20 days [hr-001]"]), "finance": ScriptedLLM(["75 dollars a day [fin-002]"])}, 1023 "PTO and per-diem?", 1024 ) 1025 for (who, task), answer in zip(res.delegations, res.answers.values()): 1026 print(f" {who:8} got only: {task!r:24} -> {answer}") 1027 say(f"Final: {res.final}") 1028 1029 banner("8. A state machine with LLM decisions at two nodes, and checkpoints") 1030 with tempfile.TemporaryDirectory() as d: 1031 for i, doc in enumerate(["INVOICE INV-2041 from Acme Supplies. Amount due: 1200.00 USD.", 1032 "INVOICE INV-2042 from Globex. Amount due: 18000.00 USD.", 1033 "Acme Supplies newsletter: our autumn catalogue is here!"]): 1034 run = InvoiceWorkflow(demo_document_llm(), Path(d) / f"c{i}.json", post=lambda inv: None).run(doc) 1035 print(f" {doc[:40]:42} -> {' > '.join(run.history)}") 1036 print() 1037 say(f"Posting crashes once. Model calls to finish: {resume_comparison(d)}") 1038 takeaway("Code owns the transitions; the model decides only where judgement is needed. Checkpoints make crashes cheap.")