How to Build Structured Workflow Pipelines Using Agno's Workflow Class
Agno's Workflow class orchestrates agents, teams, and Python functions into structured pipelines using composable Step objects and container types like Parallel and Condition.
The Workflow class in the agno-agi/agno repository provides a pipeline-style execution engine for building complex AI applications. Located in libs/agno/agno/workflow/workflow.py, this system enables developers to chain together modular units of work that can persist state, handle retries, and stream real-time events.
Core Architecture Components
Understanding the building blocks of Agno's workflow system is essential for designing robust pipelines.
The Workflow Class
The Workflow class serves as the central orchestrator. According to the source code in libs/agno/agno/workflow/workflow.py, it holds workflow metadata including name, description, session_id, and database configuration. The class manages an ordered list of steps and handles session creation through the read_or_create_session method. When run() is invoked, the workflow iterates over self.steps, passing a StepInput object to each step and collecting StepOutput results.
The Step Class
Individual units of work are encapsulated in the Step class defined in libs/agno/agno/workflow/step.py. A Step can delegate execution to an Agent, a Team, or a custom callable via the executor parameter. The class handles input validation through strict_input_validation, implements retry logic via max_retries, and merges session_state changes back to the persistent session after execution. Steps return structured StepOutput objects containing content, media artifacts, and updated state.
Container Steps for Control Flow
Agno provides specialized container step types in libs/agno/agno/workflow/steps.py and related modules that compose multiple steps into higher-order structures:
- Steps: Sequential execution of child steps
- Parallel: Concurrent execution using asyncio tasks (
libs/agno/agno/workflow/parallel.py) - Loop: Repeated execution of a sub-workflow (
libs/agno/agno/workflow/loop.py) - Condition: Branching based on a predicate function (
libs/agno/agno/workflow/condition.py) - Router: Dynamic dispatch to multiple branches (
libs/agno/agno/workflow/router.py)
Each container implements a run method that iterates over child steps while preserving the StepInput context and propagating WorkflowRunOutput for streaming events.
Implementing Sequential Pipelines
The most common pattern chains agents together where the output of one feeds into the next.
This example from cookbook/05_agent_os/workflow/basic_workflow.py demonstrates a research-to-planning pipeline:
from agno.agent.agent import Agent
from agno.models.openai.chat import OpenAIChat
from agno.workflow.step import Step
from agno.workflow.workflow import Workflow
from agno.db.sqlite import SqliteDb
# Define specialized agents
searcher = Agent(
name="HackerNews Search",
model=OpenAIChat(id="gpt-4o-mini"),
tools=[HackerNewsTools()],
)
planner = Agent(
name="Content Planner",
model=OpenAIChat(id="gpt-4o"),
instructions=[
"Create a 4-week content calendar based on the supplied research.",
"Provide 3 posts per week.",
],
)
# Wrap agents in Steps
research_step = Step(name="Research", agent=searcher)
plan_step = Step(name="Plan", agent=planner)
# Assemble workflow
content_workflow = Workflow(
name="content-creation-pipeline",
description="Generate research → plan a content calendar",
db=SqliteDb(session_table="workflow_session", db_file="tmp/workflow.db"),
steps=[research_step, plan_step],
)
# Execute
output = content_workflow.run(message="Latest AI trends")
print(output.content)
The workflow persists session data to SQLite via the db parameter, enabling stateful multi-run conversations.
Integrating Custom Logic
Agno workflows support heterogeneous execution, allowing standard Python functions to interoperate with Agent steps.
As shown in cookbook/05_agent_os/workflow/workflow_with_custom_function.py, you can pass a callable to the executor parameter:
from agno.workflow.step import Step, StepOutput
from agno.workflow.workflow import Workflow
from agno.db.sqlite import SqliteDb
def extract_keywords(step_input, session_state=None):
text = step_input.input
return StepOutput(content=[w.lower() for w in text.split() if len(w) > 4])
# Function-based step
keyword_step = Step(name="KeywordExtractor", executor=extract_keywords)
# Agent-based refinement
refiner = Agent(
name="Keyword Refiner",
model=OpenAIChat(id="gpt-4o-mini"),
instructions=["Take the list of keywords and rank them by relevance."],
)
refine_step = Step(name="RefineKeywords", agent=refiner)
pipeline = Workflow(
name="keyword-pipeline",
db=SqliteDb(session_table="wf_session", db_file="tmp/kw.db"),
steps=[keyword_step, refine_step],
)
result = pipeline.run(message="Explain the impact of quantum computing on AI.")
The StepInput object passed to custom executors contains the input message, previous step outputs, and session state, enabling complex data transformations between pipeline stages.
Advanced Execution Patterns
For sophisticated pipelines, Agno provides container steps that control execution flow.
Parallel Processing
The Parallel class in libs/agno/agno/workflow/parallel.py executes child steps concurrently using asyncio.gather:
from agno.workflow.parallel import Parallel
from agno.workflow.step import Step
from agno.workflow.workflow import Workflow
agent_a = Agent(name="SummarizerA", model=OpenAIChat(id="gpt-4o-mini"))
agent_b = Agent(name="SummarizerB", model=OpenAIChat(id="gpt-4o-mini"))
step_a = Step(name="SummarizeA", agent=agent_a)
step_b = Step(name="SummarizeB", agent=agent_b)
parallel_branch = Parallel(name="ParallelSummaries", steps=[step_a, step_b])
final_merge = Step(
name="MergeSummaries",
executor=lambda inp: StepOutput(
content=" | ".join([s.content for s in inp.previous_step_outputs])
),
)
workflow = Workflow(
name="parallel-summary",
steps=[parallel_branch, final_merge],
)
out = workflow.run(message="Explain blockchain technology.")
This pattern is ideal for fan-out operations where independent agents process the same input through different strategies or domains.
Conditional Branching
Use Condition from libs/agno/agno/workflow/condition.py to select execution paths based on runtime data:
from agno.workflow.condition import Condition
from agno.workflow.step import Step
from agno.workflow.workflow import Workflow
def is_question(step_input):
return "?" in step_input.input
question_step = Step(name="AnswerQuestion", agent=question_agent)
statement_step = Step(name="ExplainStatement", agent=statement_agent)
conditional = Condition(
name="QuestionOrStatement",
predicate=is_question,
true_steps=[question_step],
false_steps=[statement_step],
)
workflow = Workflow(name="q-or-s", steps=[conditional])
output = workflow.run(message="What is Agno?")
The predicate function receives the current StepInput and returns a boolean determining which branch executes.
Managing State and History
Agno workflows maintain persistent state through the WorkflowSession class in libs/agno/agno/session/workflow.py.
Session Persistence
When run() is called, the workflow invokes read_or_create_session to load existing state from the configured BaseDb or initialize a new session. Steps can read from and write to session_state, which persists across runs and enables multi-turn interactions.
History Injection
Enable context awareness by configuring history parameters:
workflow = Workflow(
name="history-demo",
add_workflow_history_to_steps=True, # Prepend past run messages
num_history_runs=2,
steps=[research_step, plan_step],
)
workflow.run(message="AI in healthcare")
# Second run receives context from the first
workflow.run(message="AI in finance")
This feature modifies the StepInput passed to each step, including recent conversation history from the WorkflowSession.
Streaming and Events
For real-time applications, set stream=True, stream_events=True, or stream_executor_events=True when calling run(). The workflow yields WorkflowRunOutputEvent objects (defined in libs/agno/agno/run/workflow.py) such as StepStartedEvent and StepCompletedEvent, allowing FastAPI WebSocket consumers to display progress as steps execute.
Summary
- Workflow orchestration: The
Workflowclass inlibs/agno/agno/workflow/workflow.pymanages ordered execution of steps with persistent session handling viaWorkflowSession. - Flexible executors:
Stepobjects inlibs/agno/agno/workflow/step.pyaccept Agents, Teams, or custom functions as executors, enabling hybrid pipelines. - Control flow containers: Use Parallel for concurrency, Condition for branching, Loop for iteration, and Router for dynamic dispatch.
- State management: Steps access and mutate
session_statewhich persists across runs through the configured database adapter. - Resilience features: Configure
max_retries,skip_on_failure, andstrict_input_validationon individual steps for robust error handling.
Frequently Asked Questions
Can I mix Python functions with Agent steps in the same workflow?
Yes. The Step class accepts any callable via the executor parameter alongside the agent parameter. Function executors receive step_input and optional session_state arguments and must return a StepOutput object. This allows data preprocessing, API calls, or validation logic to execute between agent steps.
How does session state persistence work across workflow runs?
The Workflow class uses read_or_create_session (defined in libs/agno/agno/workflow/workflow.py) to load or initialize a WorkflowSession from the configured database (e.g., SQLite via SqliteDb). Steps can mutate session_state during execution, and these changes are merged back into the session. Subsequent calls to run() access the same session if the same session_id is provided.
What is the difference between Steps and Parallel containers?
Steps (the base container in libs/agno/agno/workflow/steps.py) executes child steps sequentially, waiting for each to complete before starting the next. Parallel (in libs/agno/agno/workflow/parallel.py) uses asyncio.gather to run child steps concurrently, reducing total execution time when steps are independent. Both propagate the same StepInput context to their children.
How do I enable real-time streaming of workflow events?
Pass stream=True or stream_events=True to the Workflow.run() method. When enabled, the method yields WorkflowRunOutputEvent instances (such as StepStartedEvent or StepCompletedEvent from libs/agno/agno/run/workflow.py) instead of blocking until completion. This integrates with FastAPI WebSocket endpoints to provide live progress updates to clients.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →