# How to Build Structured Workflow Pipelines Using Agno's Workflow Class

> Master Agno's Workflow class to build structured pipelines with agents, teams, and Python functions. Orchestrate complex tasks effectively using composable steps and containers.

- Repository: [Agno/agno](https://github.com/agno-agi/agno)
- Tags: how-to-guide
- Published: 2026-02-23

---

**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](https://github.com/agno-agi/agno) repository provides a pipeline-style execution engine for building complex AI applications. Located in [`libs/agno/agno/workflow/workflow.py`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/workflow/parallel.py))
- **Loop**: Repeated execution of a sub-workflow ([`libs/agno/agno/workflow/loop.py`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/workflow/loop.py))
- **Condition**: Branching based on a predicate function ([`libs/agno/agno/workflow/condition.py`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/workflow/condition.py))
- **Router**: Dynamic dispatch to multiple branches ([`libs/agno/agno/workflow/router.py`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/cookbook/05_agent_os/workflow/basic_workflow.py) demonstrates a research-to-planning pipeline:

```python
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`](https://github.com/agno-agi/agno/blob/main/cookbook/05_agent_os/workflow/workflow_with_custom_function.py), you can pass a callable to the `executor` parameter:

```python
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`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/workflow/parallel.py) executes child steps concurrently using `asyncio.gather`:

```python
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`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/workflow/condition.py) to select execution paths based on runtime data:

```python
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`](https://github.com/agno-agi/agno/blob/main/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:

```python
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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/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`](https://github.com/agno-agi/agno/blob/main/libs/agno/agno/run/workflow.py)) instead of blocking until completion. This integrates with FastAPI WebSocket endpoints to provide live progress updates to clients.