Skip to content

Planner-executor

One agent writes the whole plan; code checks it and carries it out; a last agent writes the answer.

A planner agent writes the whole plan up front as data: a list of steps, each saying what to do and which earlier steps it needs. It never answers and never calls a tool. Code then checks the plan (no repeated ids, no missing dependencies, no cycles, one final step) and sends a bad plan back to the planner to fix. Once the plan is valid, code runs it a round at a time: every step whose dependencies are done runs now, all of those in parallel, and each executor sees only its own step and the results it depends on. A last agent writes the answer from whatever completed. The model chooses the shape once; code holds it to the rules and does the scheduling.

Use it when

  • A question needs several different pieces of work, and what they are depends on the question.
  • Some of the pieces are independent (so they can run at the same time) and some need others' results.
  • You want the plan to exist as data you can check, log, show or approve before anything runs.

Look elsewhere when

  • The steps are always the same: pipeline is simpler and cheaper.
  • The steps can't be known until you are partway through: supervisor.
  • The pieces never depend on each other: fan_out.
question → planner → Plan(steps) → check in code → executors, a round at a time → synthesizer
                                    (ids, needs,      (independent steps in parallel)
                                     no cycles,
                                     one final step)

What it shows

  • The model writes the shape once, up front. The planner never answers and never calls a tool: it returns a Plan of steps, each with an id, an instruction, and depends_on. Compare supervisor (the model decides what to delegate one turn at a time) and pipeline / fan_out (the shape is fixed in code)
  • Code holds the plan to rules. check_plan rejects an empty or over-long plan, repeated ids, dependencies on steps that don't exist, a step that depends on itself, a cycle, and a plan with more than one "final" step. The reason goes back to the planner, which rewrites the plan (retries={"output": 3})
  • Rounds, in parallel. Every step whose dependencies are done runs now, all of those at once (asyncio.gather). output.waves records which steps ran together
  • Each executor sees only what it needs: its own instruction and the results of the steps it depends on, plus the question for context. Not its siblings' results, and not steps upstream that it did not name
  • A failure takes down only what needed it. A step that raises is recorded in failed; every step that needed it, directly or through others, is skipped and never reaches an executor; the rest finish. The synthesizer is told which steps failed or were skipped. If nothing completes, AllStepsFailedError; if the shared budget runs out, UsageLimitExceeded ends the run (it is not one step's failure)
  • One budget for the planner, every executor and the synthesizer (USAGE_LIMITS on one Flow)

See it run: sample_run.md is a recorded run against a real model: what each agent was asked, which tools it called, and what it returned.

What the real model taught us

The model does not reliably say which steps a step needs. Asked to plan "look up A, look up B, then compare", Gemini wrote the comparison with an empty depends_on in 9 of 10 sampled plans, even with the rule spelled out and an example in the prompt. Such a step runs in the same round as the lookups, without their results, and the answer still came out right only because the executors can look cities up for themselves and the synthesizer does the arithmetic: the wrong plan was hidden by a right answer.

So the rule lives in code, not in the prompt. A plan must end in one final step that every other step feeds into. Lookups plus a compare step that forgot its dependencies leave several loose ends, and the check names them and sends the plan back, and the model fixes it (the planner makes two requests instead of one; 11 of 12 sampled plans needed exactly one correction). The wording of the rejection matters as much as the rule. Saying only what was wrong, the model repeated the same plan three times and used up its retries; saying what to change ("make 's4' the final step by setting its depends_on to include ['s1', 's2', 's3']", which the code can work out) fixed it on the first retry. That is the stronger guarantee, and it has a limit worth knowing: the check sees the graph, not what a step's words need. A calculation step that declares no dependencies but is used by the final step still passes (2 of the 10 plans we sampled after adding the rule were like that). Here that is harmless; with executors that could not look things up themselves, such a step would fail visibly instead of being rescued.

Running it

uv run python -m examples.planner_executor.agent

To use it in your project:

uv run python scripts/add_agent.py planner_executor --name research

The cities are invented, so a model cannot know them and has to use get_city. To adapt it, replace get_city and the executor's and planner's instructions (the planner is told what its executors can do, which is how it plans only steps they can carry out), and keep check_plan, the scheduling loop in run_planner_executor, and the failure handling.

run_planner_executor returns a RunResult: .output carries the answer (result, subject, value), the plan, each step's results, the waves, and any failed / skipped steps; .steps holds the planner, then each executor that finished, then the synthesizer. The tests are in examples/planner_executor/test_example.py.

Source

All of it is in examples/planner_executor/.

"""Planner-executor: one agent writes the whole plan, then code carries it out.

Use this pattern when:
- A question needs several different pieces of work, and what they are depends on the question
- Some of those pieces are independent (so they can run at the same time) and some need others' results
- You want the plan to exist as data you can check, log, show or approve *before* anything runs

How it works:

    question → planner → Plan(steps) → check it in code → executors, a wave at a time → synthesizer
                                         (ids, dependencies, no cycles)   (independent steps in parallel)

The planner never answers and never calls a tool: it only writes a `Plan`, a list of steps that each
say what to do and which earlier steps they need. Code validates the plan (a bad one is sent back to
the planner to fix), then runs it: every step whose dependencies are done runs now, all of those in
parallel, and each executor sees only its own step and the results it depends on. A step that fails
takes down only the steps that needed it. A last agent writes the answer from whatever completed.

Compare `supervisor`, where the model decides what to delegate one turn at a time and nothing is
known in advance, and `pipeline` / `fan_out`, where the shape is fixed in code. Here the *model*
chooses the shape, once, up front, and *code* holds it to the rules.

The cities are invented, so a model cannot know them and must use the tool. Replace `get_city` and
the executor's instructions with your own tools; the planner, the checking and the scheduling stay.
"""

from __future__ import annotations

import asyncio
from dataclasses import dataclass, field

from pydantic import BaseModel, Field
from pydantic_ai import Agent, ModelRetry, RunContext, ToolFailed
from pydantic_ai.capabilities import RaiseContentFilterError
from pydantic_ai.exceptions import UsageLimitExceeded
from pydantic_ai.usage import UsageLimits

from agent.config import settings
from agent.logging import agent_label, configure_logging, get_logger
from agent.prompts.templates import load_prompt
from agent.runs import Flow, RunResult

logger = get_logger(__name__)
LABEL = agent_label(__name__)  # names this agent's run spans in Logfire traces

MAX_STEPS = 6  # the most steps a plan may have (the planner is told, and the check enforces it)

# One budget for the whole run: the planner, every executor and the synthesizer share it. Sized for
# a plan of MAX_STEPS steps that each use the tool, with headroom for retries.
USAGE_LIMITS = UsageLimits(
    request_limit=40, total_tokens_limit=200_000, cost_limit=settings.cost_limit
)


# --- The data ---
@dataclass(frozen=True)
class City:
    name: str
    population: int
    area_km2: float
    founded: int


CITIES = {
    city.name: city
    for city in (
        City("Brindlemoor", 84_000, 120.0, 1721),
        City("Quillhaven", 156_000, 130.0, 1648),
        City("Tarnby", 39_000, 130.0, 1893),
        City("Ashwick", 210_000, 350.0, 1502),
    )
}


def density(city: City) -> float:
    """People per km², which is what the example's questions are about."""
    return round(city.population / city.area_km2, 1)


@dataclass
class PlanDeps:
    """Runtime dependencies shared by every agent in the run."""

    cities: dict[str, City] = field(default_factory=lambda: dict(CITIES))
    # Every city name the executors asked the tool about, in order: the evidence of what ran.
    calls: list[str] = field(default_factory=list)


# --- The plan ---
class PlanStep(BaseModel):
    id: str = Field(description="A short unique name for the step, such as 's1'.")
    instruction: str = Field(description="What this one step must do, able to be done on its own.")
    depends_on: list[str] = Field(
        default_factory=list, description="Ids of the steps whose results this step needs."
    )


class Plan(BaseModel):
    steps: list[PlanStep]


class PlanError(ValueError):
    """The plan breaks a rule. The message says which, so the planner can fix it."""


def check_plan(plan: Plan) -> None:
    """Raise PlanError unless the plan is one that can be carried out.

    A plan has between 1 and MAX_STEPS steps, unique ids, dependencies that name other steps in the
    plan, no cycles (a step cannot, even indirectly, wait for itself), and **one final step**: the
    only step nothing else depends on, so every other step feeds into it. That last rule is how a
    step that uses another's result but forgot to say so is caught. A planner that writes
    "look up A", "look up B" and "compare them" with the third step's `depends_on` left empty
    leaves three loose ends, and the check names them.
    """
    if not 1 <= len(plan.steps) <= MAX_STEPS:
        raise PlanError(f"a plan needs between 1 and {MAX_STEPS} steps, not {len(plan.steps)}")
    ids = [step.id for step in plan.steps]
    duplicates = sorted({i for i in ids if ids.count(i) > 1})
    if duplicates:
        raise PlanError(f"step ids must be unique; repeated: {duplicates}")
    for step in plan.steps:
        unknown = [d for d in step.depends_on if d not in ids]
        if unknown:
            raise PlanError(
                f"step {step.id!r} depends on {unknown}, which are not steps in the plan"
            )
        if step.id in step.depends_on:
            raise PlanError(f"step {step.id!r} depends on itself")
    # Peel off steps whose dependencies are all done; if some remain, they wait on each other.
    done: set[str] = set()
    remaining = list(plan.steps)
    while remaining:
        runnable = [s for s in remaining if set(s.depends_on) <= done]
        if not runnable:
            raise PlanError(
                f"the steps {sorted(s.id for s in remaining)} wait on each other (a cycle)"
            )
        done |= {s.id for s in runnable}
        remaining = [s for s in remaining if s.id not in done]
    needed = {dep for step in plan.steps for dep in step.depends_on}
    loose_ends = [step.id for step in plan.steps if step.id not in needed]
    if len(loose_ends) > 1:
        # Say what to change, not only what is wrong: the last loose end is the likeliest final step.
        *others, final = loose_ends
        raise PlanError(
            f"the steps {loose_ends} are not used by any other step, so the plan has no single "
            "final step that answers the question. A step that uses another step's result must "
            f"list it in depends_on. For example, make {final!r} the final step by setting its "
            f"depends_on to include {others}"
        )


# --- Agents ---
planner_agent: Agent[PlanDeps, Plan] = Agent(
    settings.model,
    name=f"{LABEL}.planner",
    output_type=Plan,
    deps_type=PlanDeps,
    capabilities=[RaiseContentFilterError()],
    instructions=load_prompt("planner_executor_planner"),
    retries={"output": 3},  # how many times a rejected plan may be sent back to be rewritten
)


@planner_agent.output_validator
def the_plan_must_be_carryable(plan: Plan) -> Plan:
    try:
        check_plan(plan)
    except PlanError as exc:
        raise ModelRetry(f"The plan is not valid: {exc}. Write the whole plan again.") from exc
    return plan


class StepResult(BaseModel):
    result: str


executor_agent: Agent[PlanDeps, StepResult] = Agent(
    settings.model,
    name=f"{LABEL}.executor",
    output_type=StepResult,
    deps_type=PlanDeps,
    capabilities=[RaiseContentFilterError()],
    instructions=load_prompt("planner_executor_executor"),
)


@executor_agent.tool
def get_city(ctx: RunContext[PlanDeps], name: str) -> dict[str, str | int | float]:
    """Look up a city's population, its area in square kilometres, and the year it was founded."""
    ctx.deps.calls.append(name)
    match = next(
        (c for c in ctx.deps.cities.values() if c.name.lower() == name.strip().lower()), None
    )
    if match is None:
        known = ", ".join(sorted(ctx.deps.cities))
        raise ToolFailed(f"There is no city called {name!r}. The cities are: {known}")
    return {
        "name": match.name,
        "population": match.population,
        "area_km2": match.area_km2,
        "founded": match.founded,
    }


class Answer(BaseModel):
    # `result` is the conventional output field in these examples; the generated
    # eval starter reads it when present (see evals/helpers.py).
    result: str
    subject: str | None = None
    value: float | None = None


synthesizer_agent: Agent[PlanDeps, Answer] = Agent(
    settings.model,
    name=f"{LABEL}.synthesizer",
    output_type=Answer,
    deps_type=PlanDeps,
    capabilities=[RaiseContentFilterError()],
    instructions=load_prompt("planner_executor_synthesizer"),
)


# --- The run ---
class PlannerExecutorOutput(BaseModel):
    result: str
    subject: str | None
    value: float | None
    plan: Plan
    results: dict[str, str]  # step id -> what its executor reported, for the steps that completed
    waves: list[list[str]]  # the step ids that ran together, in the order the waves ran
    failed: list[str]  # steps whose executor raised
    skipped: list[str]  # steps never run because a step they needed failed or was skipped


class AllStepsFailedError(Exception):
    """No step of the plan completed, so there is nothing to answer from."""


def step_prompt(question: str, step: PlanStep, results: dict[str, str]) -> str:
    """What one executor is told: the question, its own step, and only the results it depends on."""
    needed = "\n".join(f"- {dep}: {results[dep]}" for dep in step.depends_on) or "(none)"
    return (
        f"Overall question (for context): {question}\n\n"
        f"Your step ({step.id}): {step.instruction}\n\n"
        f"Results of the steps yours depends on:\n{needed}"
    )


def outcome_of(step: PlanStep, results: dict[str, str], failed: list[str]) -> str:
    """How a step ended, for the synthesizer: its result, or why there is none."""
    if step.id in results:
        return results[step.id]
    return "FAILED" if step.id in failed else "SKIPPED (a step it needed did not complete)"


async def run_planner_executor(
    user_input: str, deps: PlanDeps | None = None
) -> RunResult[PlannerExecutorOutput]:
    """Plan an answer to `user_input`, carry the plan out, and write the answer.

    Returns:
        A RunResult: `.output` is the PlannerExecutorOutput; `.steps` holds the planner, then every
        executor that finished (in completion order), then the synthesizer.

    Raises:
        AllStepsFailedError: When no step of the plan completed.
        UsageLimitExceeded: When the run's shared budget runs out (never treated as a failed step).
    """
    if deps is None:
        deps = PlanDeps()
    logger.info("Running planner-executor", extra={"user_input": user_input})
    flow = Flow(USAGE_LIMITS)  # shared by every agent below, so it bounds the whole run

    plan = (await flow.run(planner_agent, user_input, deps=deps)).output

    pending = {step.id: step for step in plan.steps}
    results: dict[str, str] = {}
    waves: list[list[str]] = []
    failed: list[str] = []
    skipped: list[str] = []
    while True:
        # A step that needs a failed or skipped step can never run. Skipping one can doom another.
        doomed = True
        while doomed:
            doomed = [s for s in pending.values() if set(s.depends_on) & set(failed + skipped)]
            for step in doomed:
                skipped.append(step.id)
                del pending[step.id]
        if not pending:
            break
        ready = [s for s in pending.values() if set(s.depends_on) <= set(results)]
        assert ready, "a checked plan always has a step that can run"  # check_plan guarantees it
        waves.append([s.id for s in ready])
        # return_exceptions=True keeps one failing step from cancelling its siblings' paid-for work.
        outcomes = await asyncio.gather(
            *(
                flow.run(executor_agent, step_prompt(user_input, s, results), deps=deps)
                for s in ready
            ),
            return_exceptions=True,
        )
        for step, outcome in zip(ready, outcomes, strict=True):
            del pending[step.id]
            if isinstance(outcome, BaseException):
                # Running out of budget ends the run; it is not one step's failure. Cancellation
                # and Ctrl-C pass through too.
                if isinstance(outcome, UsageLimitExceeded) or not isinstance(outcome, Exception):
                    raise outcome
                logger.warning("Step failed", extra={"step": step.id, "error": str(outcome)})
                failed.append(step.id)
            else:
                results[step.id] = outcome.output.result

    if not results:
        raise AllStepsFailedError(f"No step of the plan completed for: {user_input!r}")

    report = "\n\n".join(
        f"Step {step.id} ({step.instruction}): {outcome_of(step, results, failed)}"
        for step in plan.steps
    )
    answer = (
        await flow.run(synthesizer_agent, f"Question: {user_input}\n\n{report}", deps=deps)
    ).output
    return flow.finish(
        PlannerExecutorOutput(
            result=answer.result,
            subject=answer.subject,
            value=answer.value,
            plan=plan,
            results=results,
            waves=waves,
            failed=failed,
            skipped=skipped,
        )
    )


if __name__ == "__main__":
    configure_logging()
    demo_deps = PlanDeps()
    run = asyncio.run(
        run_planner_executor(
            "Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, "
            "and what is its density in people per km²?",
            demo_deps,
        )
    )
    print(run.output.result)
    print(f"plan: {[(s.id, s.depends_on) for s in run.output.plan.steps]}")
    print(f"waves: {run.output.waves}")
    print(f"{len(run.steps)} agent runs; cities looked up: {sorted(demo_deps.calls)}")
You carry out one step of a plan for answering a question about cities. You are given the overall
question for context, your step, and the results of the earlier steps your step depends on.

- Do only your step. Do not answer the overall question unless your step is to do so.
- Use `get_city` for any fact about a city. Never use a figure that is not in a tool result or in
  the results you were given.
- If a fact cannot be found (for example a city that does not exist), say so plainly in your
  result instead of guessing.
- Show the figures your result rests on, and calculate exactly.
You plan how to answer a question about cities. You do not answer it. You write a plan that other
agents will carry out.

The agents that carry out your steps can do exactly one thing besides reasoning: look up a city's
facts with `get_city` (its population, its area in km², and the year it was founded). They cannot
look up anything else and they have no calculator, so a step that needs a calculation says what to
calculate, using the results of earlier steps.

Write the plan as steps:
- `id`: a short unique name, such as "s1".
- `instruction`: what that one step must do, written so it can be done on its own. A step sees only
  its own instruction and the results of the steps it depends on, not the original question.
- `depends_on`: the ids of the steps whose results this step needs.

How `depends_on` works matters. Steps run in rounds: every step whose `depends_on` steps are done
runs now, all of those at the same time. So:
- A step that uses a figure or a finding from another step MUST list that step in `depends_on`.
  Otherwise it runs at the same time as that step and does not have its result.
- A step with an empty `depends_on` runs in the first round and must be able to do its work with
  nothing but the tools. Only look-ups can do that.
- Steps that do not need each other must not depend on each other, so that they run together.

Look up each city in its own step. Put every calculation, comparison and conclusion in a later
step whose `depends_on` lists all the steps it uses. For example, to find which of two things is
larger: s1 looks up the first, s2 looks up the second (both with an empty `depends_on`), and s3
compares them with `depends_on` of ["s1", "s2"].

The plan must end in one final step, the one that gives the answer to the question, and every
other step must be used by it, directly or through other steps. A plan with several steps that nothing
else uses is wrong.

Use as few steps as the question needs, and never more than six.
You write the final answer to a question about cities, from the results of the steps that were
carried out to answer it. You are given the question and each step with its result.

- Use only the step results. Never add a figure that is not in them.
- Some steps may be marked failed or skipped. Answer as much of the question as the completed
  steps allow, and say plainly what could not be done and why.
- Fill `subject` with the name of the city (or other thing) the answer is about, and `value` with the
  number the question asks for, rounded to one decimal place. Leave either empty if the question
  does not ask for one or the steps do not give it.
title = "Planner-executor"
pattern = "planner_executor"
summary = "A planner writes the whole plan as data, code checks it and runs it (independent steps in parallel, each executor seeing only what it needs), and a last agent writes the answer."
smoke_input = "Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?"
# A real model's executors must look the cities up rather than guess (checked by the release check).
expected_tools = ["get_city"]

# The offline smoke test scripts the planner's plan, so the executors and the synthesizer really run.
[smoke.planner_agent]
output = { steps = [{ id = "s1", instruction = "Look up Quillhaven.", depends_on = [] }] }

[entrypoint]
deps = "PlanDeps"
run = "run_planner_executor"
"""The plan check, the scheduling of steps, and what a failed step takes down: all offline.

The three agents run on scripted models, so the planner writes exactly the plan a test wants and the
executors record what they were told and how many were running at once. The tool is the real one.
"""

import asyncio
import re

import pytest
from pydantic_ai import ToolFailed
from pydantic_ai.exceptions import UsageLimitExceeded
from pydantic_ai.messages import (
    ModelResponse,
    ToolCallPart,
    ToolReturnPart,
    UserPromptPart,
)
from pydantic_ai.models.function import AgentInfo, FunctionModel

from examples.planner_executor.agent import (
    CITIES,
    MAX_STEPS,
    AllStepsFailedError,
    Plan,
    PlanDeps,
    PlanError,
    PlanStep,
    check_plan,
    density,
    executor_agent,
    get_city,
    planner_agent,
    run_planner_executor,
    step_prompt,
    synthesizer_agent,
)

QUESTION = "Which city is densest?"


def plan_of(*steps: tuple[str, list[str]]) -> Plan:
    return Plan(steps=[PlanStep(id=i, instruction=f"do {i}", depends_on=deps) for i, deps in steps])


def prompt_of(messages) -> str:
    return "\n".join(str(p.content) for p in messages[-1].parts if isinstance(p, UserPromptPart))


# --- Scripted agents ---


def planner(plans: list[Plan], seen: list[str] | None = None):
    """Returns the given plans in order, one per request, and records what it was sent."""
    remaining = list(plans)

    def model_fn(messages, info: AgentInfo) -> ModelResponse:
        if seen is not None:
            seen.append(
                " | ".join(
                    str(p.content) for m in messages for p in m.parts if hasattr(p, "content")
                )
            )
        plan = remaining.pop(0)
        return ModelResponse(parts=[ToolCallPart(info.output_tools[0].name, plan.model_dump())])

    return FunctionModel(model_fn)


class Executors:
    """A scripted executor that records every prompt, how many ran at once, and who was asked."""

    def __init__(self, fail_on=(), lookups: dict[str, str] | None = None, pause: float = 0.02):
        self.fail_on = set(fail_on)
        self.lookups = lookups or {}  # step id -> a city to look up with the real tool first
        self.pause = pause
        self.prompts: dict[str, str] = {}
        self.running = 0
        self.peak = 0

    @staticmethod
    def step_id(prompt: str) -> str:
        return re.search(r"Your step \((\w+)\)", prompt).group(1)

    @property
    def model(self) -> FunctionModel:
        async def model_fn(messages, info: AgentInfo) -> ModelResponse:
            step = self.step_id(str(messages[0].parts[-1].content))
            returned = any(isinstance(p, ToolReturnPart) for m in messages for p in m.parts)
            if step in self.lookups and not returned:
                return ModelResponse(parts=[ToolCallPart("get_city", {"name": self.lookups[step]})])
            self.prompts[step] = str(messages[0].parts[-1].content)
            self.running += 1
            self.peak = max(self.peak, self.running)
            try:
                await asyncio.sleep(self.pause)  # long enough for siblings to overlap
                if step in self.fail_on:
                    raise RuntimeError(f"{step} blew up")
            finally:
                self.running -= 1
            return ModelResponse(
                parts=[ToolCallPart(info.output_tools[0].name, {"result": f"result of {step}"})]
            )

        return FunctionModel(model_fn)


def synthesizer(seen: list[str]):
    def model_fn(messages, info: AgentInfo) -> ModelResponse:
        seen.append(prompt_of(messages))
        fields = {"result": "the answer", "subject": "Quillhaven", "value": 1200.0}
        return ModelResponse(parts=[ToolCallPart(info.output_tools[0].name, fields)])

    return FunctionModel(model_fn)


async def run(plans: list[Plan], executors: Executors, deps: PlanDeps | None = None):
    seen: list[str] = []
    with (
        planner_agent.override(model=planner(plans)),
        executor_agent.override(model=executors.model),
        synthesizer_agent.override(model=synthesizer(seen)),
    ):
        return await run_planner_executor(QUESTION, deps), seen


# --- The data and the tool ---


def test_the_densities_are_what_the_questions_expect():
    assert {name: density(city) for name, city in CITIES.items()} == {
        "Brindlemoor": 700.0,
        "Quillhaven": 1200.0,
        "Tarnby": 300.0,
        "Ashwick": 600.0,
    }
    assert max(CITIES.values(), key=density).name == "Quillhaven"  # no ties for first place


def test_the_tool_finds_a_city_by_any_capitalization_and_records_the_call():
    deps = PlanDeps()
    ctx = type("Ctx", (), {"deps": deps})()
    found = get_city(ctx, "  qUILLhaven ")
    assert found == {
        "name": "Quillhaven",
        "population": 156_000,
        "area_km2": 130.0,
        "founded": 1648,
    }
    assert deps.calls == ["  qUILLhaven "]


def test_the_tool_says_which_cities_exist_when_asked_for_one_that_does_not():
    ctx = type("Ctx", (), {"deps": PlanDeps()})()
    with pytest.raises(ToolFailed, match="no city called 'Atlantis'.*Ashwick, Brindlemoor"):
        get_city(ctx, "Atlantis")


def test_each_run_gets_its_own_cities_and_ledger():
    first, second = PlanDeps(), PlanDeps()
    first.calls.append("Tarnby")
    first.cities.pop("Tarnby")
    assert second.calls == [] and "Tarnby" in second.cities and "Tarnby" in CITIES


# --- The plan check ---


def test_a_good_plan_passes():
    check_plan(plan_of(("a", []), ("b", []), ("c", ["a", "b"])))
    check_plan(plan_of(("only", [])))
    check_plan(plan_of(("a", []), ("b", ["a"]), ("c", ["a"]), ("d", ["b", "c"])))  # a diamond
    check_plan(plan_of(*[(f"s{i}", [f"s{i - 1}"] if i else []) for i in range(MAX_STEPS)]))


@pytest.mark.parametrize(
    ("plan", "message"),
    [
        (Plan(steps=[]), "between 1 and"),
        (plan_of(*[(f"s{i}", []) for i in range(MAX_STEPS + 1)]), "between 1 and"),
        (plan_of(("a", []), ("a", [])), r"unique; repeated: \['a'\]"),
        (plan_of(("a", ["nope"])), r"depends on \['nope'\]"),
        (plan_of(("a", ["a"])), "depends on itself"),
        (plan_of(("a", ["b"]), ("b", ["a"])), r"\['a', 'b'\] wait on each other"),
        (plan_of(("ok", []), ("a", ["c"]), ("b", ["a"]), ("c", ["b"])), r"\['a', 'b', 'c'\] wait"),
        (plan_of(("a", []), ("b", [])), r"\['a', 'b'\] are not used by any other step"),
        (plan_of(("a", []), ("b", []), ("c", ["a"])), r"\['b', 'c'\] are not used"),
        # What the rule is for: two lookups and a "compare them" step that forgot to depend on them.
        (plan_of(("s1", []), ("s2", []), ("s3", [])), r"\['s1', 's2', 's3'\] are not used"),
    ],
)
def test_a_plan_that_cannot_be_carried_out_is_rejected_with_the_reason(plan, message):
    with pytest.raises(PlanError, match=message):
        check_plan(plan)


def test_a_rejected_plan_is_told_what_to_change_not_only_what_is_wrong():
    plan = plan_of(("s1", []), ("s2", []), ("s3", []), ("s4", []))
    with pytest.raises(PlanError) as caught:
        check_plan(plan)
    # The last loose end is offered as the final step, with the others as what it should depend on.
    assert (
        "make 's4' the final step by setting its depends_on to include ['s1', 's2', 's3']"
        in str(caught.value)
    )


def test_the_planner_is_told_what_its_executors_can_do():
    from agent.prompts.templates import load_prompt

    text = load_prompt("planner_executor_planner")
    assert "get_city" in text  # the executor's one tool
    assert "six" in text and MAX_STEPS == 6  # the limit it is told is the one that is enforced


async def test_a_rejected_plan_is_sent_back_with_the_reason_and_the_fixed_one_runs():
    seen: list[str] = []
    cyclic = plan_of(("a", ["b"]), ("b", ["a"]))
    good = plan_of(("a", []))
    with (
        planner_agent.override(model=planner([cyclic, good], seen)),
        executor_agent.override(model=Executors().model),
        synthesizer_agent.override(model=synthesizer([])),
    ):
        result = await run_planner_executor(QUESTION)

    assert [s.id for s in result.output.plan.steps] == ["a"]
    assert "wait on each other" in seen[1]  # the second request carried the planner's mistake
    assert [s.agent for s in result.steps].count("planner_executor.planner") == 1  # one planner run


# --- Running the plan ---


async def test_independent_steps_run_together_and_the_step_that_needs_them_runs_after():
    executors = Executors()
    plan = plan_of(("s1", []), ("s2", []), ("s3", []), ("s4", ["s1", "s2", "s3"]))
    result, seen = await run([plan], executors)
    output = result.output

    assert output.waves == [["s1", "s2", "s3"], ["s4"]]
    assert executors.peak == 3  # the three lookups really overlapped
    assert output.failed == [] and output.skipped == []
    assert set(output.results) == {"s1", "s2", "s3", "s4"}
    assert [s.agent for s in result.steps][0] == "planner_executor.planner"
    assert [s.agent for s in result.steps][-1] == "planner_executor.synthesizer"
    assert len(result.steps) == 6  # planner, four executors, synthesizer


async def test_a_step_sees_its_own_instruction_and_only_the_results_it_depends_on():
    executors = Executors()
    plan = plan_of(("s1", []), ("s2", []), ("s3", ["s1"]), ("s4", ["s2", "s3"]))
    await run([plan], executors)

    assert "do s1" in executors.prompts["s1"] and "(none)" in executors.prompts["s1"]
    assert "result of s2" not in executors.prompts["s1"]  # a sibling's result is not shared
    assert "result of s1" in executors.prompts["s3"]  # its dependency's is
    assert "result of s2" not in executors.prompts["s3"]  # a step it does not depend on is not
    assert QUESTION in executors.prompts["s3"]  # the question is given for context
    assert "result of s2" in executors.prompts["s4"] and "result of s3" in executors.prompts["s4"]
    assert "result of s1" not in executors.prompts["s4"]  # only what it named, not what is upstream


async def test_a_chain_of_steps_runs_one_wave_at_a_time():
    executors = Executors()
    result, _ = await run([plan_of(("s1", []), ("s2", ["s1"]), ("s3", ["s2"]))], executors)
    assert result.output.waves == [["s1"], ["s2"], ["s3"]]
    assert executors.peak == 1


async def test_executors_use_the_real_tool_and_the_ledger_records_it():
    deps = PlanDeps()
    executors = Executors(lookups={"s1": "Quillhaven", "s2": "Tarnby"})
    result, _ = await run([plan_of(("s1", []), ("s2", []), ("s3", ["s1", "s2"]))], executors, deps)

    assert sorted(deps.calls) == ["Quillhaven", "Tarnby"]
    returns = [
        p.content
        for m in result.all_messages()
        for p in m.parts
        if isinstance(p, ToolReturnPart) and p.tool_name == "get_city"
    ]
    assert {r["name"] for r in returns} == {"Quillhaven", "Tarnby"}


async def test_the_synthesizer_is_given_the_question_and_every_step_result():
    result, seen = await run([plan_of(("s1", []), ("s2", ["s1"]))], Executors())
    assert QUESTION in seen[0]
    assert "result of s1" in seen[0] and "result of s2" in seen[0]
    assert (result.output.result, result.output.subject, result.output.value) == (
        "the answer",
        "Quillhaven",
        1200.0,
    )


async def test_one_budget_covers_every_agent_in_the_run():
    result, _ = await run([plan_of(("s1", []), ("s2", []), ("s3", ["s1", "s2"]))], Executors())
    assert result.usage.requests == 5  # planner + three executors + synthesizer


def test_a_steps_prompt_says_so_when_it_depends_on_nothing():
    step = PlanStep(id="s1", instruction="look it up")
    text = step_prompt("the question", step, {})
    assert "the question" in text and "look it up" in text and "(none)" in text


# --- When a step fails ---


async def test_a_failed_step_skips_what_needed_it_and_nothing_else():
    executors = Executors(fail_on={"s2"})
    plan = plan_of(
        ("s1", []),
        ("s2", []),
        ("s3", []),
        ("s4", ["s2", "s3"]),
        ("s5", ["s4"]),
        ("s6", ["s5", "s1"]),
    )
    result, seen = await run([plan], executors)
    output = result.output

    assert output.failed == ["s2"]
    assert output.skipped == ["s4", "s5", "s6"]  # s5 only needed s4, which was skipped: it cascades
    assert set(output.results) == {"s1", "s3"}  # the siblings of the failed step still ran
    assert output.waves == [["s1", "s2", "s3"]]
    assert set(executors.prompts) == {"s1", "s2", "s3"}  # s4, s5 and s6 never reached an executor
    # The synthesizer is told what did not happen, step by step, rather than left to guess.
    assert "Step s2 (do s2): FAILED" in seen[0]
    assert all(f"Step {s} (do {s}): SKIPPED" in seen[0] for s in ("s4", "s5", "s6"))
    assert "result of s2" not in seen[0]


async def test_a_failure_with_a_long_chain_behind_it_skips_the_whole_chain_and_ends():
    """No wave is left to run when the chain is skipped, so the skipping itself must reach its end."""
    executors = Executors(fail_on={"s1"})
    plan = plan_of(("s0", []), ("s1", []), ("s2", ["s1"]), ("s3", ["s2"]), ("s4", ["s3", "s0"]))
    result, _ = await run([plan], executors)

    assert result.output.failed == ["s1"]
    assert result.output.skipped == ["s2", "s3", "s4"]
    assert result.output.waves == [["s0", "s1"]]


async def test_when_no_step_completes_there_is_nothing_to_answer_from():
    def must_not_run(messages, info):
        raise AssertionError("the synthesizer ran with nothing to go on")

    with (
        planner_agent.override(model=planner([plan_of(("s1", []), ("s2", ["s1"]))])),
        executor_agent.override(model=Executors(fail_on={"s1"}).model),
        synthesizer_agent.override(model=FunctionModel(must_not_run)),
        pytest.raises(AllStepsFailedError),
    ):
        await run_planner_executor(QUESTION)


async def test_running_out_of_budget_ends_the_run_instead_of_being_a_failed_step():
    def over_budget(messages, info):
        raise UsageLimitExceeded("the request limit was reached")

    with (
        planner_agent.override(
            model=planner([plan_of(("s1", []), ("s2", []), ("s3", ["s1", "s2"]))])
        ),
        executor_agent.override(model=FunctionModel(over_budget)),
        synthesizer_agent.override(model=synthesizer([])),
        pytest.raises(UsageLimitExceeded),
    ):
        await run_planner_executor(QUESTION)


async def test_cancellation_is_not_mistaken_for_a_failed_step():
    async def cancelled(messages, info: AgentInfo) -> ModelResponse:
        raise asyncio.CancelledError

    with (
        planner_agent.override(model=planner([plan_of(("s1", []))])),
        executor_agent.override(model=FunctionModel(cancelled)),
        synthesizer_agent.override(model=synthesizer([])),
        pytest.raises(asyncio.CancelledError),
    ):
        await run_planner_executor(QUESTION)
"""Live check: a real model plans, real executors look cities up, and the answer is right. `-m eval`.

The expected answers are worked out here from the data, not read from the model's output, and the
plan is checked for its shape (independent lookups together, a step that needs them after), not for
exact wording, so another model that plans differently still passes.
"""

import pytest

from evals.trace import traced_run
from examples.live_support import assert_every_agent_ran, run_as_script
from examples.planner_executor import agent as module
from examples.planner_executor.agent import CITIES, density

pytestmark = pytest.mark.eval

DENSEST_QUESTION = (
    "Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, "
    "and what is its density in people per km²?"
)


async def ask(question: str):
    """Run with a fresh ledger: returns the traced run and the deps it used."""
    deps = module.PlanDeps()

    async def helper(text: str):
        return await module.run_planner_executor(text, deps)

    return await traced_run(helper, question), deps


def looked_up(deps: module.PlanDeps) -> set[str]:
    return {name.strip().lower() for name in deps.calls}


@pytest.fixture(scope="module")
async def densest():
    return await ask(DENSEST_QUESTION)


@pytest.fixture(scope="module")
async def percent_larger():
    return await ask(
        "How much larger, in percent, is the population of Ashwick than Brindlemoor's?"
    )


@pytest.fixture(scope="module")
async def missing_city():
    return await ask("What is the combined population of Quillhaven and Atlantis?")


async def test_the_answer_matches_what_the_data_says(densest):
    traced, _ = densest
    output = traced.result.output
    best = max(CITIES.values(), key=density)  # Quillhaven, 1200.0: no tie for first place

    assert output.subject == best.name
    assert output.value == pytest.approx(density(best), abs=0.1)
    assert output.failed == [] and output.skipped == []


async def test_every_city_in_the_question_was_really_looked_up(densest):
    _, deps = densest
    assert {"brindlemoor", "quillhaven", "tarnby"} <= looked_up(deps)


async def test_the_plan_groups_independent_lookups_and_has_a_step_that_needs_them(densest):
    traced, _ = densest
    output = traced.result.output
    plan = output.plan

    assert len(plan.steps) >= 3
    assert len(output.waves[0]) >= 2  # lookups that need nothing ran together, not one by one
    assert any(step.depends_on for step in plan.steps)  # something was left to do with the lookups
    assert {s for wave in output.waves for s in wave} == {step.id for step in plan.steps}


async def test_each_agent_ran_in_the_order_the_pattern_says(densest):
    traced, _ = densest
    agents = [step.agent for step in traced.result.steps]
    assert agents[0] == "planner_executor.planner"
    assert agents[-1] == "planner_executor.synthesizer"
    assert agents[1:-1] and set(agents[1:-1]) == {"planner_executor.executor"}
    assert "get_city" in traced.tools_called


async def test_a_calculation_that_needs_two_lookups_is_exact(percent_larger):
    traced, deps = percent_larger
    ashwick, brindlemoor = CITIES["Ashwick"].population, CITIES["Brindlemoor"].population
    expected = (ashwick - brindlemoor) / brindlemoor * 100  # 150.0

    assert traced.result.output.value == pytest.approx(expected, abs=0.1)
    assert {"ashwick", "brindlemoor"} <= looked_up(deps)


async def test_a_city_that_does_not_exist_is_reported_not_invented(missing_city):
    traced, deps = missing_city
    output = traced.result.output

    assert "atlantis" in looked_up(deps)  # an executor really asked the tool, and was told no
    assert "Atlantis" in output.result  # the answer says what it could not do
    # No invented total: either no figure, or only the one city that does exist.
    assert output.value is None or output.value == pytest.approx(CITIES["Quillhaven"].population)


async def test_every_agent_ran(densest, percent_larger, missing_city):
    ran = densest[0].agents_ran | percent_larger[0].agents_ran | missing_city[0].agents_ran
    assert_every_agent_ran(module, ran)


async def test_the_demo_script_runs():
    out = await run_as_script("examples.planner_executor.agent")
    assert "waves:" in out and "cities looked up:" in out

Recorded run · gemini-3.1-flash-lite · 6 steps · $0.0021

Recorded 2026-10-07 with google:gemini-3.1-flash-lite · 6 steps · 5,278 tokens · $0.0021 · 6.2 s. Model output varies between runs. Regenerate with uv run python scripts/record_example.py planner_executor.

Input

Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Steps

1. planner_executor.planner 1,918 tokens · $0.0009

Prompt

Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

What happened - asked to retry: The plan is not valid: the steps ['s1', 's2', 's3', 's4'] are not used by any other step, so the plan has no single final step that answers the question. A step that uses another step's result must li …

Output

{
  "steps": [
    {
      "id": "s1",
      "instruction": "Look up the population and area of Brindlemoor.",
      "depends_on": []
    },
    {
      "id": "s2",
      "instruction": "Look up the population and area of Quillhaven.",
      "depends_on": []
    },
    {
      "id": "s3",
      "instruction": "Look up the population and area of Tarnby.",
      "depends_on": []
    },
    {
      "id": "s4",
      "instruction": "Calculate the population density (population / area) for each city using the results from s1, s2, and s3. Compare these values to determine which city is the most densely populated and state its name and its density.",
      "depends_on": [
        "s1",
        "s2",
        "s3"
      ]
    }
  ]
}

2. planner_executor.executor 735 tokens · $0.0003

Prompt

Overall question (for context): Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Your step (s2): Look up the population and area of Quillhaven.

Results of the steps yours depends on: (none)

What happened - called get_city({"name": "Quillhaven"}) - get_city returned: {"name": "Quillhaven", "population": 156000, "area_km2": 130.0, "founded": 1648}

Output

{
  "result": "Quillhaven has a population of 156,000 and an area of 130 km²."
}

3. planner_executor.executor 735 tokens · $0.0003

Prompt

Overall question (for context): Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Your step (s1): Look up the population and area of Brindlemoor.

Results of the steps yours depends on: (none)

What happened - called get_city({"name": "Brindlemoor"}) - get_city returned: {"name": "Brindlemoor", "population": 84000, "area_km2": 120.0, "founded": 1721}

Output

{
  "result": "Brindlemoor has a population of 84,000 and an area of 120 km²."
}

4. planner_executor.executor 733 tokens · $0.0003

Prompt

Overall question (for context): Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Your step (s3): Look up the population and area of Tarnby.

Results of the steps yours depends on: (none)

What happened - called get_city({"name": "Tarnby"}) - get_city returned: {"name": "Tarnby", "population": 39000, "area_km2": 130.0, "founded": 1893}

Output

{
  "result": "Tarnby has a population of 39,000 and an area of 130 km²."
}

5. planner_executor.executor 555 tokens · $0.0003

Prompt

Overall question (for context): Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Your step (s4): Calculate the population density (population / area) for each city using the results from s1, s2, and s3. Compare these values to determine which city is the most densely populated and state its name and its density.

Results of the steps yours depends on: - s1: Brindlemoor has a population of 84,000 and an area of 120 km². - s2: Q …

Output

{
  "result": "- Brindlemoor: 84,000 / 120 km² = 700 people/km²\n- Quillhaven: 156,000 / 130 km² = 1,200 people/km²\n- Tarnby: 39,000 / 130 km² = 300 people/km²\n\nComparing these, Quillhaven is the most densely populated city with a density of 1,200 people/km²."
}

6. planner_executor.synthesizer 602 tokens · $0.0002

Prompt

Question: Which of Brindlemoor, Quillhaven and Tarnby is the most densely populated, and what is its density in people per km²?

Step s1 (Look up the population and area of Brindlemoor.): Brindlemoor has a population of 84,000 and an area of 120 km².

Step s2 (Look up the population and area of Quillhaven.): Quillhaven has a population of 156,000 and an area of 130 km².

Step s3 (Look up the population and area of Tarnby.): Tarnby has a population of 39,000 and an area of 130 km².

Step s4 (Calc …

Output

{
  "result": "Quillhaven is the most densely populated city with a density of 1,200 people/km².",
  "subject": "Quillhaven",
  "value": 1200.0
}

Result

run_planner_executor(...).output

{
  "result": "Quillhaven is the most densely populated city with a density of 1,200 people/km².",
  "subject": "Quillhaven",
  "value": 1200.0,
  "plan": {
    "steps": [
      {
        "id": "s1",
        "instruction": "Look up the population and area of Brindlemoor.",
        "depends_on": []
      },
      {
        "id": "s2",
        "instruction": "Look up the population and area of Quillhaven.",
        "depends_on": []
      },
      {
        "id": "s3",
        "instruction": "Look up the population and area of Tarnby.",
        "depends_on": []
      },
      {
        "id": "s4",
        "instruction": "Calculate the population density (population / area) for each city using the results from s1, s2, and s3. Compare these values to determine which city is the most densely populated and state its name and its density.",
        "depends_on": [
          "s1",
          "s2",
          "s3"
        ]
      }
    ]
  },
  "results": {
    "s1": "Brindlemoor has a population of 84,000 and an area of 120 km².",
    "s2": "Quillhaven has a population of 156,000 and an area of 130 km².",
    "s3": "Tarnby has a population of 39,000 and an area of 130 km².",
    "s4": "- Brindlemoor: 84,000 / 120 km² = 700 people/km²\n- Quillhaven: 156,000 / 130 km² = 1,200 people/km²\n- Tarnby: 39,000 / 130 km² = 300 people/km²\n\nComparing these, Quillhaven is the most densely populated city with a density of 1,200 people/km²."
  },
  "waves": [
    [
      "s1",
      "s2",
      "s3"
    ],
    [
      "s4"
    ]
  ],
  "failed": [],
  "skipped": []
}