Pipeline¶
Fixed sequential steps, each one's output feeding the next. Also called a prompt chain.
A pipeline breaks a task into stages that always run in the same order, each its own agent with a focused prompt and a typed output: here an outline, then a draft written from it, then a polished version. Code, not the model, decides the order, so the control flow is ordinary Python that is cheap to test and to bound. Between stages, code can check the previous output and stop early, so later stages never spend tokens on nothing.
Use it when
- A task decomposes into stages that always run in the same order.
- Each stage is more reliable as its own focused prompt than inside one long prompt.
- You want to check intermediate results and stop early.
Look elsewhere when
- The steps vary by input:
routerorplanner_executor. - The stages don't depend on each other and could run together:
fan_out.
What it shows
- Each step is its own agent with a typed output (
Outline,Draft,Polished) - A gate between steps: plain code checks the previous output and fails fast
(
EmptyOutlineError), so later steps never spend tokens on nothing - One
RunUsageshared by every step, soUSAGE_LIMITSbounds the whole chain - Trace labels
<name>.outline,<name>.draft,<name>.polish
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.
run_pipeline returns a RunResult: .output is the piece, .steps the three steps in order,
and .usage the total across all of them.
Adapt it by changing the steps and what each one returns. The generic smoke test runs the whole
chain offline; examples/pipeline/test_example.py tests the gate and the shared budget.
Source¶
All of it is in examples/pipeline/.
"""Pipeline (prompt chain): fixed sequential steps, each one's output feeding the next.
Use this pattern when:
- The task decomposes into stages that always run in the same order
- Each stage is easier and more reliable as its own focused prompt
- You want to check intermediate results and stop early, rather than trust one long prompt
outline → draft → polish (code decides the order; the model never does)
Unlike `supervisor`, nothing here is chosen by the LLM: the control flow is ordinary Python,
so it is cheap to reason about, test and bound.
"""
from __future__ import annotations
from dataclasses import dataclass
from pydantic import BaseModel
from pydantic_ai import Agent
from pydantic_ai.capabilities import RaiseContentFilterError
from pydantic_ai.usage import UsageLimits
from agent.config import settings
from agent.logging import agent_label, configure_logging, get_logger
from agent.runs import Flow, RunResult
logger = get_logger(__name__)
LABEL = agent_label(__name__) # names this agent's run spans in Logfire traces
# One budget for the whole chain: every step runs in one Flow (see run_pipeline).
USAGE_LIMITS = UsageLimits(
request_limit=15, total_tokens_limit=150_000, cost_limit=settings.cost_limit
)
@dataclass
class PipelineDeps:
"""Runtime dependencies shared by every step."""
pass
def _step(role: str, output_type: type[BaseModel], instructions: str) -> Agent:
return Agent(
settings.model,
name=f"{LABEL}.{role}",
output_type=output_type,
deps_type=PipelineDeps,
capabilities=[RaiseContentFilterError()],
instructions=instructions,
)
# --- Step 1: outline ---
class Outline(BaseModel):
points: list[str]
outline_agent: Agent[PipelineDeps, Outline] = _step(
"outline",
Outline,
"Write a short outline for a piece on the given topic: three to five key points, each a "
"single sentence.",
)
# --- Step 2: draft ---
class Draft(BaseModel):
text: str
draft_agent: Agent[PipelineDeps, Draft] = _step(
"draft",
Draft,
"Write a first draft from the outline you are given. Cover every point, in order.",
)
# --- Step 3: polish ---
class PipelineOutput(BaseModel):
result: str
outline: list[str]
class Polished(BaseModel):
result: str
polish_agent: Agent[PipelineDeps, Polished] = _step(
"polish",
Polished,
"Edit the draft for clarity and concision. Keep its meaning and structure.",
)
class EmptyOutlineError(Exception):
"""The outline step produced nothing to draft from."""
async def run_pipeline(
user_input: str, deps: PipelineDeps | None = None
) -> RunResult[PipelineOutput]:
"""Write a short piece on `user_input` in three chained steps.
Returns:
A RunResult: `.output` is the PipelineOutput; `.steps` holds the outline, draft and
polish steps in order.
Raises:
EmptyOutlineError: When step 1 returns no points. Failing here is the payoff of a
chain: the later steps never spend tokens on nothing.
"""
if deps is None:
deps = PipelineDeps()
flow = Flow(USAGE_LIMITS) # one shared budget, so USAGE_LIMITS bounds the whole chain
outline = (await flow.run(outline_agent, f"Topic: {user_input}", deps=deps)).output
# A gate between steps: plain code checking the previous step's typed output.
if not outline.points:
raise EmptyOutlineError(f"No outline was produced for: {user_input!r}")
logger.info("Outline ready", extra={"points": len(outline.points)})
numbered = "\n".join(f"{i}. {point}" for i, point in enumerate(outline.points, 1))
draft = (await flow.run(draft_agent, f"Outline:\n{numbered}", deps=deps)).output
polished = (await flow.run(polish_agent, f"Draft:\n{draft.text}", deps=deps)).output
return flow.finish(PipelineOutput(result=polished.result, outline=outline.points))
if __name__ == "__main__":
import asyncio
configure_logging()
print(asyncio.run(run_pipeline("Why unit tests are worth writing")).output)
"""The chain runs in order, passes each step's output on, and stops at a failed gate."""
import pytest
from pydantic_ai.exceptions import UsageLimitExceeded
from pydantic_ai.messages import ModelResponse, ToolCallPart, UserPromptPart
from pydantic_ai.models.function import AgentInfo, FunctionModel
from pydantic_ai.usage import UsageLimits
from examples.pipeline import agent as pipeline
from examples.pipeline.agent import (
EmptyOutlineError,
draft_agent,
outline_agent,
polish_agent,
run_pipeline,
)
def prompt_of(messages) -> str:
return "\n".join(str(p.content) for p in messages[-1].parts if isinstance(p, UserPromptPart))
def returns(fields: dict, seen: list[str] | None = None):
def model_fn(messages, info: AgentInfo) -> ModelResponse:
if seen is not None:
seen.append(prompt_of(messages))
return ModelResponse(parts=[ToolCallPart(info.output_tools[0].name, fields)])
return FunctionModel(model_fn)
async def test_each_step_receives_the_previous_steps_output():
draft_in: list[str] = []
polish_in: list[str] = []
with (
outline_agent.override(model=returns({"points": ["first point", "second point"]})),
draft_agent.override(model=returns({"text": "a rough draft"}, draft_in)),
polish_agent.override(model=returns({"result": "a polished piece"}, polish_in)),
):
result = await run_pipeline("testing")
output = result.output
assert output.result == "a polished piece"
assert [step.agent for step in result.steps] == [
"pipeline.outline",
"pipeline.draft",
"pipeline.polish",
]
assert result.usage.requests == 3
assert output.outline == ["first point", "second point"]
assert "1. first point" in draft_in[0] and "2. second point" in draft_in[0]
assert "a rough draft" in polish_in[0]
async def test_an_empty_outline_stops_the_chain_before_later_steps_run():
def must_not_run(messages, info):
raise AssertionError("a step after the failed gate ran")
with (
outline_agent.override(model=returns({"points": []})),
draft_agent.override(model=FunctionModel(must_not_run)),
polish_agent.override(model=FunctionModel(must_not_run)),
pytest.raises(EmptyOutlineError),
):
await run_pipeline("testing")
async def test_one_budget_covers_every_step(monkeypatch):
"""Three steps need three requests; request_limit=2 stops the chain at the third."""
monkeypatch.setattr(pipeline, "USAGE_LIMITS", UsageLimits(request_limit=2))
with (
outline_agent.override(model=returns({"points": ["a"]})),
draft_agent.override(model=returns({"text": "d"})),
polish_agent.override(model=returns({"result": "p"})),
pytest.raises(UsageLimitExceeded),
):
await run_pipeline("testing")
"""Live check: the chain runs end to end and each step feeds the next. Run with `pytest -m eval`."""
import pytest
from evals.trace import traced_run
from examples.live_support import assert_every_agent_ran, run_as_script
from examples.pipeline import agent as module
pytestmark = pytest.mark.eval
async def test_the_three_steps_run_in_order_and_produce_a_piece():
traced = await traced_run(module.run_pipeline, "Why unit tests are worth writing")
output = traced.result.output
assert len(output.outline) >= 3 # the outline step was asked for three to five points
assert len(output.result.split()) >= 30 # a real piece of writing, not a stub
assert "test" in output.result.lower()
assert [step.agent for step in traced.result.steps] == [
"pipeline.outline",
"pipeline.draft",
"pipeline.polish",
]
assert traced.result.usage.requests >= 3 # one per step, plus any output retries
assert_every_agent_ran(module, traced.agents_ran)
async def test_the_demo_script_runs():
assert "outline=" in await run_as_script("examples.pipeline.agent")
Recorded run · gemini-3.1-flash-lite · 3 steps · $0.0011
Recorded 2026-10-07 with google:gemini-3.1-flash-lite · 3 steps · 1,227 tokens · $0.0011 · 4.5 s.
Model output varies between runs. Regenerate with uv run python scripts/record_example.py pipeline.
Input
Why unit tests are worth writing
Steps
1. pipeline.outline
188 tokens · $0.0002
Prompt
Topic: Why unit tests are worth writing
Output
{
"points": [
"Unit tests provide immediate feedback that helps identify and fix bugs early in the development cycle.",
"They act as living documentation, making it easier for developers to understand how specific parts of the codebase should function.",
"Automated tests allow developers to refactor code with confidence, knowing that regressions will be quickly detected.",
"Maintaining a comprehensive test suite reduces long-term technical debt and lowers the overall cost of software maintenance."
]
}
2. pipeline.draft
461 tokens · $0.0005
Prompt
Outline: 1. Unit tests provide immediate feedback that helps identify and fix bugs early in the development cycle. 2. They act as living documentation, making it easier for developers to understand how specific parts of the codebase should function. 3. Automated tests allow developers to refactor code with confidence, knowing that regressions will be quickly detected. 4. Maintaining a comprehensive test suite reduces long-term technical debt and lowers the overall cost of software maintenance.
Output
{
"text": "Unit testing is a cornerstone of modern software development, providing critical benefits that extend well beyond simple bug detection.\n\nFirst, unit tests provide immediate feedback, allowing developers to identify and resolve defects early in the development cycle. By catching issues at the moment they are introduced, teams can address problems while the context is still fresh, avoiding the significantly higher costs associated with fixing bugs later in production.\n\nSecond, these tests act as living documentation for the project. For developers joining a new team or revisiting an old codebase, unit tests serve as a clear, executable reference for how specific functions and modules are expected to behave. Rather than relying solely on outdated manuals, developers can look at the tests to understand the intended functionality and logic of the code.\n\nThird, automated testing provides developers the confidence to refactor code. When a comprehensive suite of tests exists, developers can improve the design, readability, and efficiency of the codebase without the fear of breaking existing features. If a change introduces a regression, the automated suite will catch it immediately, ensuring that the software remains stable throughout the development process.\n\nFinally, maintaining a comprehensive test suite is an investment that reduces long-term technical debt. While writing tests requires an upfront time commitment, it significantly lowers the overall cost of software maintenance over the life of the product. By preventing the accumulation of hidden bugs and ensuring the codebase remains modular and understandable, unit tests make the system more sustainable and easier to evolve over time."
}
3. pipeline.polish
578 tokens · $0.0004
Prompt
Draft: Unit testing is a cornerstone of modern software development, providing critical benefits that extend well beyond simple bug detection.
First, unit tests provide immediate feedback, allowing developers to identify and resolve defects early in the development cycle. By catching issues at the moment they are introduced, teams can address problems while the context is still fresh, avoiding the significantly higher costs associated with fixing bugs later in production.
Second, these tests a …
Output
{
"result": "Unit testing is a cornerstone of modern software development, offering benefits that extend well beyond basic bug detection.\n\nFirst, unit tests provide immediate feedback, allowing developers to resolve defects early. By catching issues as they are introduced, teams address problems while the context is fresh, avoiding the high costs of fixing bugs in production.\n\nSecond, tests serve as living documentation. For developers onboarding or revisiting a codebase, unit tests provide an executable reference for how functions and modules should behave, offering clearer insights than outdated manuals.\n\nThird, automated tests empower developers to refactor with confidence. A comprehensive suite allows for improvements in design, readability, and efficiency without the fear of breaking existing features. If a regression occurs, the tests catch it immediately, ensuring system stability.\n\nFinally, maintaining a test suite is an investment that reduces long-term technical debt. While writing tests requires an upfront time commitment, it lowers maintenance costs over the product's life. By preventing hidden bugs and encouraging modularity, unit tests make the system more sustainable and easier to evolve."
}
Result
run_pipeline(...).output
{
"result": "Unit testing is a cornerstone of modern software development, offering benefits that extend well beyond basic bug detection.\n\nFirst, unit tests provide immediate feedback, allowing developers to resolve defects early. By catching issues as they are introduced, teams address problems while the context is fresh, avoiding the high costs of fixing bugs in production.\n\nSecond, tests serve as living documentation. For developers onboarding or revisiting a codebase, unit tests provide an executable reference for how functions and modules should behave, offering clearer insights than outdated manuals.\n\nThird, automated tests empower developers to refactor with confidence. A comprehensive suite allows for improvements in design, readability, and efficiency without the fear of breaking existing features. If a regression occurs, the tests catch it immediately, ensuring system stability.\n\nFinally, maintaining a test suite is an investment that reduces long-term technical debt. While writing tests requires an upfront time commitment, it lowers maintenance costs over the product's life. By preventing hidden bugs and encouraging modularity, unit tests make the system more sustainable and easier to evolve.",
"outline": [
"Unit tests provide immediate feedback that helps identify and fix bugs early in the development cycle.",
"They act as living documentation, making it easier for developers to understand how specific parts of the codebase should function.",
"Automated tests allow developers to refactor code with confidence, knowing that regressions will be quickly detected.",
"Maintaining a comprehensive test suite reduces long-term technical debt and lowers the overall cost of software maintenance."
]
}