AgentFieldbuild

Orchestration patterns

Multi-agent patterns for parallel, sequential, and fan-out/fan-in execution using cross-agent calls and shared memory.

Three orchestration patterns: sequential, parallel, fan-in

Compose agents into parallel, sequential, and fan-out/fan-in workflows using the primitives you already have.

AgentField does not require a separate orchestration DSL or workflow engine. You orchestrate with normal language primitives, then the control plane turns those calls into a traced workflow.

That means you keep native asyncio, Promise.all, or goroutines, but still get a workflow DAG, propagated context, and execution visibility across the whole fan-out/fan-in flow.

@app.reasoner()
async def underwrite_application(application: dict) -> dict:
    # Fan-out: run three checks in parallel — each is a separate agent
    credit, fraud, compliance = await asyncio.gather(
        app.call("credit-scorer.score", ssn=application["ssn"]),
        app.call("fraud-detector.screen", applicant=application),
        app.call("compliance.check_sanctions", name=application["name"]),
    )
    # Fan-in: aggregate results for final decision
    if fraud["risk"] > 0.7 or not compliance["cleared"]:
        return {"decision": "denied", "reason": "risk threshold exceeded"}
    return {"decision": "approved", "credit_score": credit["score"]}
    # Every call is a node in the execution DAG — full observability for free

What just happened

  • The orchestration logic stayed inside normal application code
  • Parallel child calls became tracked DAG nodes automatically
  • The control plane recorded the fan-out and fan-in structure without a separate workflow definition

Example DAG shape:

{
  "target": "underwriter.underwrite_application",
  "children": [
    { "target": "credit-scorer.score", "status": "completed" },
    { "target": "fraud-detector.screen", "status": "completed" },
    { "target": "compliance.check_sanctions", "status": "completed" }
  ],
  "result": { "decision": "approved" }
}
What you get
  • Parallel execution -- call multiple agents simultaneously with asyncio.gather(), Promise.all(), or goroutines.
  • Sequential pipelines -- chain agents where each step's output feeds the next step's input.
  • Fan-out/fan-in -- distribute work items across agents, then merge results.
  • Workflow DAG -- every cross-agent call creates a tracked edge, visualizable through the control plane API.
  • Automatic context propagation -- workflow ID, session ID, and actor ID flow through every pattern.
Parallel execution

Call multiple agents at the same time and wait for all results. Each call runs as an independent execution but shares the same workflow context.

import asyncio
from agentfield import Agent

app = Agent(node_id="orchestrator")

@app.reasoner()
async def analyze_document(document: str) -> dict:
    # Run three analyses in parallel
    sentiment_task = app.call("nlp-agent.sentiment", text=document)
    entities_task = app.call("nlp-agent.extract_entities", text=document)
    summary_task = app.call("summarizer.summarize", text=document)

    sentiment, entities, summary = await asyncio.gather(
        sentiment_task,
        entities_task,
        summary_task,
    )

    return {
        "sentiment": sentiment,
        "entities": entities,
        "summary": summary,
    }
Sequential pipeline

Chain agents where each step transforms or enriches the data from the previous step. The workflow DAG captures the full chain.

@app.reasoner()
async def claims_pipeline(claim: dict) -> dict:
    # Step 1: Validate the claim
    validation = await app.call("validator.check_claim", claim=claim)
    if not validation.get("valid"):
        return {"status": "rejected", "reason": validation["reason"]}

    # Step 2: Assess risk using validation context
    risk = await app.call(
        "risk-engine.assess",
        claim=claim,
        validation=validation,
    )

    # Step 3: Make decision based on risk
    decision = await app.call(
        "adjudicator.decide",
        claim=claim,
        risk_score=risk["score"],
        risk_factors=risk["factors"],
    )

    # Step 4: Notify the customer
    await app.call(
        "notifier.send_decision",
        claim_id=claim["id"],
        decision=decision["action"],
    )

    return decision
Fan-out / fan-in

Distribute a batch of work items across agents in parallel, then merge the results. This pattern is common for processing lists, running multiple analyses on different data partitions, or parallelizing expensive LLM calls.

import asyncio

@app.reasoner()
async def batch_analyze(documents: list[dict]) -> dict:
    # Fan-out: analyze each document in parallel
    tasks = [
        app.call(
            "analyzer.process",
            document=doc["text"],
            doc_id=doc["id"],
        )
        for doc in documents
    ]
    results = await asyncio.gather(*tasks, return_exceptions=True)

    # Fan-in: merge results, handle failures
    successes = []
    failures = []
    for doc, result in zip(documents, results):
        if isinstance(result, Exception):
            failures.append({"id": doc["id"], "error": str(result)})
        else:
            successes.append(result)

    return {
        "total": len(documents),
        "succeeded": len(successes),
        "failed": len(failures),
        "results": successes,
        "errors": failures,
    }
Workflow DAG visualization

Every app.call() creates a parent-child edge in the execution DAG. The control plane tracks these relationships automatically through the context headers propagated on each call.

Querying the DAG

# Get the full execution DAG for a workflow
curl http://localhost:8080/api/ui/v1/workflows/{workflowId}/dag

The response contains:

  • Every execution node (agent, reasoner/skill name, status, duration)
  • Parent-child edges showing the call graph
  • Timing data for latency analysis

What the DAG Captures

PatternDAG Shape
Sequential pipelineLinear chain: A -> B -> C -> D
Parallel executionFan: A -> [B, C, D]
Fan-out/fan-inDiamond: A -> [B1, B2, B3] -> C
Nested orchestrationTree with multiple levels

The DAG updates in real time as executions complete, so you can monitor long-running workflows in the dashboard.

Advanced patterns

Parallel with Shared State

Use shared memory to coordinate between parallel agents without passing all state through call arguments.

@app.reasoner()
async def coordinated_analysis(document: str) -> dict:
    # Store the document in workflow memory so all agents can access it
    await app.memory.set("source_document", document)

    # Run parallel analyses -- each agent reads from shared memory
    results = await asyncio.gather(
        app.call("legal-agent.review"),
        app.call("financial-agent.review"),
        app.call("compliance-agent.review"),
    )

    # Each agent wrote its findings to shared memory
    legal = await app.memory.get("findings.legal")
    financial = await app.memory.get("findings.financial")
    compliance = await app.memory.get("findings.compliance")

    return {
        "legal": legal,
        "financial": financial,
        "compliance": compliance,
    }

Conditional Branching

Route to different agents based on a classification step.

@app.reasoner()
async def smart_router(request: dict) -> dict:
    # Step 1: Classify the request
    classification = await app.call(
        "classifier.classify",
        text=request["message"],
    )

    # Step 2: Route to the appropriate specialist
    target = {
        "billing": "billing-agent.handle",
        "technical": "tech-support.handle",
        "general": "general-agent.handle",
    }.get(classification["category"], "general-agent.handle")

    result = await app.call(target, **request)
    return result

Retry with Fallback

When a primary agent fails, discover and use an alternative.

@app.reasoner()
async def resilient_pipeline(data: dict) -> dict:
    try:
        return await app.call("primary-analyzer.process", **data)
    except Exception as primary_error:
        app.note(f"Primary failed: {primary_error}", ["fallback"])

        # Discover a healthy alternative
        alternatives = app.discover(
            tags=["analyzer"],
            health_status="active",
        )

        for cap in alternatives.json.capabilities:
            if cap.agent_id == "primary-analyzer":
                continue
            try:
                target = f"{cap.agent_id}.{cap.reasoners[0].id}"
                return await app.call(target, **data)
            except Exception:
                continue

        raise RuntimeError("All analyzers failed")