Communication Between Agents
Learning Outcomes
- Compare message passing, shared state, and event-driven communication and choose the right one per edge
- Design a versioned message envelope with correlation IDs and idempotency keys
- Implement a blackboard with optimistic concurrency control to coordinate agents over shared state
- Build an event-driven agent that reacts to a published event without tight coupling
- Decide what to pass between agents — full context, summaries, or references — and why
Lesson Plan
| Segment | Duration | Topic |
|---|---|---|
| Intro | 3 min | The three paradigms and why this is the hard part |
| Explain | 7 min | Message envelopes, schemas, and serialization |
| Build | 9 min | A typed message-passing protocol with Pydantic |
| Build | 9 min | The blackboard: shared state with optimistic concurrency |
| Build | 8 min | Event-driven agents over a broker |
| Explain | 7 min | What to pass: context, summaries, or references |
| Wrap-up | 2 min | A three-agent protocol and trade-offs recap |
Before You Begin
Pre-work:
- Complete Lesson 2: Orchestration Patterns — communication choices follow from the topology you chose
- Review tool use and agent loops from the Agentic AI course, especially Tool Use and Memory Systems
- Be comfortable reading async Python and distributed-systems vocabulary (idempotency, eventual consistency, back-pressure)
Shopping List:
- Python 3.11+ with
pydanticv2 and theanthropicSDK installed - A broker for the event-driven section — Redis (with Streams) is enough; Kafka or NATS work the same conceptually
- Optional: a running Redis (
docker run -p 6379:6379 redis) to execute the examples
With more than one agent, the architecture stops being about prompts and starts being about how information moves. There are three paradigms, and most production systems blend them — each gets a step below:
| Paradigm | Coupling | Best when |
|---|---|---|
| Message passing | Tighter — sender knows receiver | Request/response, handoffs, supervisor → worker |
| Shared state (blackboard) | Loose on identity, tight on data | Many agents building one evolving artifact |
| Event-driven | Loosest — sender doesn't know receivers | Fan-out, reactive pipelines, adding agents without edits |
The deeper question this lesson keeps returning to is what you put on the wire. An agent's message is tokens that cost money and fill a finite context window — passing the wrong thing is the most common cause of expensive, slow, confused multi-agent systems.
Before any paradigm, agree on a message envelope — metadata wrapping every payload regardless of content. It is what makes a system debuggable and recoverable. Treat it like an HTTP header set.
| Field | Purpose |
|---|---|
message_id |
Unique per message (UUID) — dedup and logging |
correlation_id |
Ties all messages in one task together for tracing |
causation_id |
The message that caused this one — reconstructs the chain |
schema_version |
Lets receivers handle old and new payloads |
idempotency_key |
So a redelivered message isn't processed twice |
JSON is the pragmatic default for serialization: readable in logs, universally supported, and already what LLM tool calls produce. Reach for Protobuf or Avro only for a strict cross-language contract or a measured throughput problem. And validate at every boundary — a malformed message from a hallucinating agent must fail loudly at the edge, not corrupt downstream state.
from pydantic import BaseModel, Field
from datetime import datetime, timezone
from uuid import uuid4
class Envelope(BaseModel):
message_id: str = Field(default_factory=lambda: str(uuid4()))
correlation_id: str
causation_id: str | None = None
sender: str
recipient: str | None = None
schema_version: int = 1
idempotency_key: str | None = None
timestamp: datetime = Field(default_factory=lambda: datetime.now(timezone.utc))
Message passing suits the supervisor pattern: a coordinator hands a typed task to a named worker and awaits a typed result. Define request and response as schemas, then carry them in the envelope.
from typing import Literal
class ResearchRequest(BaseModel):
type: Literal["research_request"] = "research_request"
query: str
deadline_seconds: int = 30
class ResearchResult(BaseModel):
type: Literal["research_result"] = "research_result"
findings: list[str]
confidence: float
class Message(BaseModel):
envelope: Envelope
body: ResearchRequest | ResearchResult
A worker runs its agent loop and returns a Message whose reply envelope copies the request's correlation_id (preserving the task thread) and sets causation_id to the request's message_id — so you can reconstruct which message produced which.
The supervisor must await with a timeout — never assume a worker responds, since an agent can stall on a slow tool call or burn its budget. Wrap the call in asyncio.wait_for(..., timeout=req.deadline_seconds) and retry, fall back, or escalate on TimeoutError. Synchronous request/response is easiest to reason about but couples the supervisor's liveness to the slowest worker; for long-running work, go async — fire the request, return a handle, let the worker post back when ready.
idempotency_key. If the supervisor times out and retries, a worker that already ran can return its cached result instead of paying for the same expensive generation twice.In the blackboard pattern, agents don't address each other. They read from and write to a shared workspace and a controller wakes the agent whose preconditions are met — ideal for agents collaborating on one evolving artifact (a document, a plan, a codebase).
The hard part is concurrency: if two agents read v3, both edit, and both write, the second silently clobbers the first — a lost update. Use optimistic concurrency control — every write declares the version it was based on, and the store rejects stale writes.
from dataclasses import dataclass, field
class StaleWriteError(Exception): ...
@dataclass
class Blackboard:
state: dict = field(default_factory=dict)
version: int = 0
def read(self) -> tuple[dict, int]:
return dict(self.state), self.version
def write(self, patch: dict, base_version: int) -> int:
if base_version != self.version:
raise StaleWriteError(f"based on v{base_version}, store at v{self.version}")
self.state.update(patch)
self.version += 1
return self.version
Agents do read-modify-write with a bounded retry budget — a compare-and-swap loop. Read state and its version, compute a patch (possibly via the model), then write(patch, base_version=ver); on StaleWriteError, re-read and retry, capping attempts so contention can't loop forever. In production the blackboard is Redis, Postgres, or DynamoDB and the version becomes a row version or conditional-write expression — same pattern.
Event-driven communication is the loosest coupling: an agent publishes a fact (source_fetched, draft_ready) with no idea who consumes it. New agents subscribe without editing the publisher — this scales a system's number of behaviors without scaling its interconnections.
Use a durable log, not fire-and-forget pub/sub, so a crashed consumer can resume. Redis Streams gives consumer groups and acknowledgements; Kafka and NATS JetStream offer the same semantics at scale.
import json, redis.asyncio as redis
r = redis.Redis(decode_responses=True)
async def publish(event_type: str, env: Envelope, payload: dict):
await r.xadd("agent.events", {
"type": event_type,
"envelope": env.model_dump_json(),
"payload": json.dumps(payload),
})
async def subscribe(group: str, consumer: str, handler):
while True:
resp = await r.xreadgroup(group, consumer, {"agent.events": ">"},
count=10, block=5000)
for _stream, messages in resp or []:
for msg_id, fields in messages:
await handler(fields) # dedup on idempotency_key
await r.xack("agent.events", group, msg_id) # ack ONLY on success
Acknowledge only after the handler succeeds, so an unacked event is redelivered after a crash; dedup on the envelope's idempotency key to make redelivery safe. Create the group once at startup with xgroup_create(..., mkstream=True).
correlation_id on every event — it's the only thread that reconstructs one task across agents that never directly spoke.This decision determines whether your system is fast and cheap or slow and broke. Every token passed between agents is paid for and competes for the receiver's context window. There are three options; mature systems use all three deliberately.
| Strategy | What moves | Cost | Use when |
|---|---|---|---|
| Full context | Whole transcript / raw documents | Highest, overflow risk | Receiver needs verbatim source (legal review, exact quotes) |
| Summary | A model-generated digest | Low, lossy | Phase handoffs needing conclusions, not raw material |
| Reference | An ID/URI; fetch on demand | Lowest, adds a round-trip | Large artifacts only some receivers open |
The reference approach mirrors REST: pass a handle, dereference only if needed. An Artifact schema carries an artifact_id, a short summary for routing, a uri to the full body, and a token_estimate so the receiver can budget before fetching. A coordinator routes on the cheap summary; only the downstream agent that needs detail pays to dereference the artifact. That is the difference between cost that grows linearly with agents and cost that explodes quadratically because everyone forwards everything.
Now apply all three paradigms to a coordinator / search / synthesis research system — the workflow you'll formalize in Lesson 5 with the Anthropic Agent SDK. The coordinator uses message passing to delegate; the search agent writes artifacts to a blackboard; a findings_ready event triggers synthesis without polling:
coordinator ──ResearchRequest (msg)──▶ search_agent
│ writes artifacts
▼
blackboard ──findings_ready (event)──▶ synthesis
coordinator ◀────────── Report (summary + refs) ──────────────────────────┘
Every edge passes the minimum fidelity its receiver needs: the event carries only a blackboard version and artifact IDs, synthesis fetches full artifacts on demand, and the report returns a summary plus citations rather than the raw transcript.
Questions & Answers
UPDATE; in Redis a WATCH/MULTI block. Pessimistic locking works too, but it serializes your agents and reintroduces the bottleneck orchestration was meant to remove.correlation_id on every event plus distributed tracing that stitches the spans into one timeline. Without that discipline they are genuinely hard to debug — which is why observability (Lesson 8) is a prerequisite before adopting heavy event-driven coupling in production.schema_version) and make consumers tolerant readers: ignore unknown fields, never hard-fail on additions. Treat removals and type changes as breaking changes needing a new major version, run old and new side by side during migration. It's ordinary API contract discipline — multi-agent systems inherit every versioning problem of distributed systems.Key Takeaways
- Three paradigms, chosen per edge — message passing for request/response, shared state for collaboration on one artifact, event-driven for decoupled fan-out. Real systems mix them.
- Envelope everything — message_id, correlation_id, causation_id, schema_version, and idempotency_key turn opaque LLM output into a traceable, recoverable, deduplicatable system.
- Validate at every boundary — parse agent output through a strict schema the moment it arrives; silent garbage propagates and fails three hops downstream.
- Shared state needs concurrency control — optimistic versioning prevents lost updates, and a termination fixpoint stops agents livelocking and billing forever.
- What you pass is the cost lever — prefer summaries and references over full context; verbatim payloads only for agents that need the source. This is linear vs. quadratic cost growth.
- It's a distributed system — idempotency, back-pressure, schema versioning, and observability aren't optional extras; agent communication inherits every hard problem of distributed computing.
Next Steps: Lesson 4: Task Decomposition at Scale