Skip to content

Chains

A Chain is a directed graph workflow where nodes execute agents, tools, or logic, and edges define the flow between them. Chains support cycles (retry loops), typed state, checkpointing with resume, human-in-the-loop approval, parallel execution, and conditional branching.

Quick Start

from fastaiagent import Agent, Chain, LLMClient

summarizer = Agent(
    name="summarizer",
    system_prompt="Summarize the input in one sentence.",
    llm=LLMClient(provider="openai", model="gpt-4.1"),
)
translator = Agent(
    name="translator",
    system_prompt="Translate the input to French. Output only the French text.",
    llm=LLMClient(provider="anthropic", model="claude-sonnet-4-6"),
)

chain = Chain("summarize-and-translate")
chain.add_node("summarize", agent=summarizer)
chain.add_node("translate", agent=translator)
chain.connect("summarize", "translate")

result = chain.execute({"message": "Python is a popular programming language"})
print(result.output)
print(result.node_results)     # {"summarize": {...}, "translate": {...}}
print(result.execution_id)     # UUID for checkpointing/resume

Node Types

Type Purpose Example
agent Run an agent (default) chain.add_node("research", agent=my_agent)
tool Execute a tool directly chain.add_node("fetch", tool=my_tool)
condition Branch based on state See Conditional Branching
transformer Render a template chain.add_node("fmt", type=NodeType.transformer, template="Hello {{state.name}}")
parallel Run multiple agents concurrently See Parallel Execution
hitl Pause for human approval See Human-in-the-Loop
start / end Explicit entry/exit points chain.add_node("in", type=NodeType.start)
from fastaiagent.chain import NodeType

chain.add_node("classify", agent=classifier_agent)
chain.add_node("transform", type=NodeType.transformer, template="Category: {{state.category}}")
chain.add_node("approve", type=NodeType.hitl)

Connecting Nodes

Simple Edges

chain.connect("a", "b")  # a → b
chain.connect("b", "c")  # b → c

Conditional Edges

Route flow based on chain state, the initial input, or the output of an earlier node:

chain.connect(
    "classify",
    "billing_agent",
    condition="{{node_results.classify.output}} == billing",
)
chain.connect(
    "classify",
    "tech_agent",
    condition="{{node_results.classify.output}} == technical",
)
chain.connect("classify", "general_agent")  # default fallback

Condition expressions support: ==, !=, >, <, >=, <=, contains, startswith. Values must be resolved via {{path.to.value}} templates against the chain context (input, state, node_results); bare identifiers like category are treated as literal text, not state lookups.

Routing semantics

Full contract

For the authoritative spec (including strict routing, parallel failure modes, validator rules, and the resume contract) see Execution Spec. The section below is a summary.

The executor walks edges in declaration order, applying these rules at each source node:

  1. All-unconditional fan-out — if every outgoing edge from a source has condition=None, every target runs (the original chain behavior).
  2. Mixed conditional + default — if any outgoing edge has condition= set, the first matching condition in declaration order wins. A single unconditional sibling acts as the default fallback when no condition matches.
  3. NodeType.condition nodes — the node returns {"matched": handle}; the outgoing edge whose label equals that handle wins, with the unlabeled or "default"-labeled edge as fallback.
  4. Dead branches — a node that no upstream branch routes to is silently skipped (it does not appear in result.node_results).

chain.validate() enforces two structural rules at design time: a source mixing conditional and unconditional edges may have at most one default, and a condition node's outgoing edges must label every handle it can return (or provide a default).

Typed State

Validate the chain state at every step using JSON Schema:

chain = Chain(
    "typed-pipeline",
    state_schema={
        "type": "object",
        "properties": {
            "message": {"type": "string"},
            "category": {"type": "string"},
            "priority": {"type": "integer"},
        },
        "required": ["message"],
    },
)
chain.add_node("classify", agent=classifier)

# This works — message is present
result = chain.execute({"message": "My order is late", "priority": 1})

# This raises ChainStateValidationError — missing required field
result = chain.execute({"priority": 1})

State is validated: 1. Before execution starts (initial state) 2. After each node updates the state

Accessing State in Nodes

Each node receives the current state in its context. Agent nodes receive the input as their prompt. Transformer nodes can template against the full state:

chain.add_node(
    "format_output",
    type=NodeType.transformer,
    template="Customer {{state.name}} (priority: {{state.priority}}): {{node_results.classify.output}}",
)

Parallel Execution

Run multiple agents concurrently within a single node:

chain = Chain("parallel-pipeline")
chain.add_node("start", type=NodeType.start)
chain.add_node(
    "parallel_research",
    type=NodeType.parallel,
    agents=[researcher_1, researcher_2, researcher_3],  # Run all 3 in parallel
)
chain.add_node("merge", type=NodeType.transformer, template="Results: {{node_results.parallel_research.outputs}}")
chain.connect("start", "parallel_research")
chain.connect("parallel_research", "merge")

Parallel nodes use asyncio.gather() for concurrent execution.

Conditional Branching

Route execution based on state or previous node output:

chain = Chain("routing-pipeline")
chain.add_node("classify", agent=classifier_agent)
chain.add_node(
    "router",
    type=NodeType.condition,
    conditions=[
        {"expression": "category == billing", "handle": "billing"},
        {"expression": "category == technical", "handle": "technical"},
    ],
)
chain.add_node("billing_agent", agent=billing_agent)
chain.add_node("tech_agent", agent=tech_agent)

chain.connect("classify", "router")
chain.connect("router", "billing_agent", condition="category == billing")
chain.connect("router", "tech_agent", condition="category == technical")

Chain Validation

Validate chain structure before execution:

errors = chain.validate()
if errors:
    for e in errors:
        print(f"Error: {e}")
else:
    print("Chain is valid")

Checks for: - Missing edge targets (referencing nonexistent nodes) - Orphaned nodes (no incoming or outgoing edges) - Cyclic edges without max_iterations

Tool Node State Behavior

When a tool node executes, its return value is wrapped in {"output": <return_value>, "error": <error_or_None>} and merged into chain state. This means each successive tool node overwrites state.output with its own wrapped result.

If you need to thread a value across multiple tool nodes (e.g., a seed_value that step A produces and step C reads), put it on the top-level state via initial_state to chain.execute() or modified_state to chain.resume() — not as a return value from a tool node. Top-level state keys persist because nothing overwrites them.

# This is fragile — step_c can't reliably read step_a's output
# because step_b's output overwrites state.output

# This is reliable — seed_value persists at the top level
result = chain.execute({"seed_value": "original", "message": "go"})
# In the tool node, read via input_mapping:
#   input_mapping={"seed": "{{state.seed_value}}"}

Agent nodes do not have this wrapping quirk — their output is stored under _{node_id}_output in state, preserving it across nodes.

ChainResult

Every chain execution returns a ChainResult:

Field Type Description
output Any Final node's output
final_state dict Chain state after all nodes complete
execution_id str UUID for checkpointing and resume
node_results dict Map of node_id -> output for each executed node
result = chain.execute({"message": "Hello"})

# Access individual node outputs
for node_id, output in result.node_results.items():
    print(f"{node_id}: {output}")

Serialization

Chains serialize to a ReactFlow-compatible JSON format (used by the platform visual editor):

# Serialize
data = chain.to_dict()
# {
#   "name": "my-pipeline",
#   "nodes": [{"id": "a", "type": "agent", "label": "A", "position": {"x": 0, "y": 0}, ...}],
#   "edges": [{"source": "a", "target": "b", "is_cyclic": false, ...}],
#   "state_schema": {...}
# }

# Restore
chain = Chain.from_dict(data)

This canonical format can be used to serialize/deserialize chains for storage or transfer.

Sync vs Async

# Sync
result = chain.execute({"input": "data"})

# Async
result = await chain.aexecute({"input": "data"})

# Async resume
result = await chain.resume(execution_id="<id>")

Error Handling

from fastaiagent._internal.errors import (
    ChainError,                   # Base chain error
    ChainCycleError,              # Cycle exceeded max_iterations
    ChainCheckpointError,         # Checkpoint save/load failed
    ChainStateValidationError,    # State failed schema validation
)

try:
    result = chain.execute({"input": "data"})
except ChainCycleError as e:
    print(f"Cycle limit hit: {e}")
except ChainStateValidationError as e:
    print(f"Invalid state: {e}")
except ChainCheckpointError as e:
    print(f"Checkpoint error: {e}")

Complete Example

A support pipeline with classification, conditional routing, retry loop, and approval:

from fastaiagent import Agent, Chain, LLMClient
from fastaiagent.chain import NodeType

llm = LLMClient(provider="openai", model="gpt-4.1")

chain = Chain(
    "support-pipeline",
    state_schema={
        "type": "object",
        "properties": {
            "message": {"type": "string"},
            "category": {"type": "string"},
        },
        "required": ["message"],
    },
)

# Nodes
chain.add_node("classify", agent=Agent(name="classifier",
    system_prompt="Classify the support request. Set category.", llm=llm))
chain.add_node("research", agent=Agent(name="researcher",
    system_prompt="Research the issue thoroughly.", llm=llm))
chain.add_node("draft", agent=Agent(name="drafter",
    system_prompt="Draft a helpful response.", llm=llm))
chain.add_node("review", type=NodeType.hitl)
chain.add_node("send", agent=Agent(name="sender",
    system_prompt="Finalize and send the response.", llm=llm))

# Flow
chain.connect("classify", "research")
chain.connect("research", "draft")
chain.connect("draft", "review")
chain.connect("review", "send")

result = chain.execute(
    {"message": "My order hasn't arrived"},
    hitl_handler=lambda n, c, s: True,  # Auto-approve for demo
)
print(result.output)

Dependency injection — RunContext

Pass a RunContext to chain.execute(...) / chain.aexecute(...) to make shared dependencies available to every tool and agent node, identical to the contract you'd use on a plain Agent:

from dataclasses import dataclass
import fastaiagent as fa
from fastaiagent.chain import Chain, NodeType

@dataclass
class Deps:
    db: ...           # connection pool, API clients, tenant id, etc.

@fa.tool()
def write_record(payload: dict, ctx: fa.RunContext[Deps]) -> dict:
    return ctx.state.db.insert(payload)

chain = Chain("ingest", checkpoint_enabled=True)
chain.add_node("persist", type=NodeType.tool, tool=write_record,
               input_mapping={"payload": "{{state.input}}"})

ctx = fa.RunContext(state=Deps(db=my_pool))
result = await chain.aexecute({"input": {...}}, context=ctx)

Context flows through:

  • every tool node — tool.aexecute(args, context=ctx)
  • every agent node — agent.arun(input, context=ctx)
  • every parallel child agent
  • recursive cycle re-entries via cyclic edges
  • resume runs — pass the same context to chain.aresume(execution_id, ..., context=ctx) so a tool re-firing after interrupt() still sees ctx.state

Context-propagation is opt-in: tools whose ctx parameter has a default of None continue to run unchanged when no context is supplied.


Next Steps