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:

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 Workflow class in libs/agno/agno/workflow/workflow.py manages ordered execution of steps with persistent session handling via WorkflowSession.
  • Flexible executors: Step objects in libs/agno/agno/workflow/step.py accept 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_state which persists across runs through the configured database adapter.
  • Resilience features: Configure max_retries, skip_on_failure, and strict_input_validation on 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →