Orchestration patterns
Multi-agent patterns for parallel, sequential, and fan-out/fan-in execution using cross-agent calls and shared memory.
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 freeWhat 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 decisionFan-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}/dagThe 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
| Pattern | DAG Shape |
|---|---|
| Sequential pipeline | Linear chain: A -> B -> C -> D |
| Parallel execution | Fan: A -> [B, C, D] |
| Fan-out/fan-in | Diamond: A -> [B1, B2, B3] -> C |
| Nested orchestration | Tree 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 resultRetry 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")