Communication Between Agents

45 min intermediate Lesson 3

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:

Shopping List:

  • Python 3.11+ with pydantic v2 and the anthropic SDK 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

1 The Three Paradigms (and the One Question Behind Them)

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.

NOTE
Key Insight
Topology (Lesson 2) decides who talks to whom. Communication decides how the bytes move and what they contain. A supervisor pattern usually implies message passing; a pipeline often implies shared state or events. Match the paradigm to your topology, then design the payload deliberately.

2 Envelopes, Schemas, and Serialization

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))
WARNING
Watch Out
Never let agents invent message formats from free text. Parse every LLM JSON blob through a strict validator at the receiving boundary. Garbage that passes silently between agents is the multi-agent equivalent of a buffer overflow — it surfaces three hops later as an inexplicable failure.

3 Message Passing — A Typed Request/Response Protocol

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.

TIP
Tip
Make worker handlers idempotent, keyed on the envelope's 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.

4 Shared State — The Blackboard with Optimistic Concurrency

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.

WARNING
Watch Out
A blackboard with no termination condition lets agents thrash forever, each re-triggering the others. Always define a fixpoint — a quiescent state where no agent's preconditions fire — and a hard iteration cap as a safety net. Multi-agent livelock is real, and it bills by the token.

5 Event-Driven — Publish, Subscribe, React

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).

NOTE
Key Insight
Events buy back-pressure for free: a slow consumer just lets events queue instead of being overwhelmed. The trade-off is losing a global view, so observability (Lesson 8) becomes mandatory. Carry the correlation_id on every event — it's the only thread that reconstructs one task across agents that never directly spoke.

6 What to Actually Pass — and a Three-Agent Protocol

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.

WARNING
Watch Out
Summaries are lossy by definition. A fact-checking or compliance agent must receive the source — full context or a fetchable reference — never a summary. You cannot verify a claim against a paraphrase of the evidence. Match fidelity to the receiver's job.
TIP
Tip
Prompt caching changes this math: if many calls share a large stable prefix (a spec, a codebase, a policy doc), passing full context can be cheaper than re-summarizing once you count cache hits. We quantify this in Lesson 9: Cost Management.

Questions & Answers

Q: How do I stop two agents from corrupting shared state when they write concurrently?
Use optimistic concurrency control: every write declares the version it was based on, and the store rejects it if the version moved (Step 4); the agent then re-reads and retries within a bounded budget. In Postgres that's a row version column with a conditional 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.
Q: Won't event-driven make debugging impossible, since no one has the full picture?
That's the real trade-off. You regain visibility with a 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.
Q: An upstream agent changed its output schema and broke a downstream consumer. How do I prevent this?
Version every schema in the envelope (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.
Q: How much context should I pass? Passing everything feels safest.
Passing everything is the most common and most expensive mistake — it inflates cost, slows every call, and degrades quality as the receiver drowns in irrelevant context. Default to summaries plus references (Step 6); pass full context only when the job needs verbatim source. If per-task token spend grows faster than the number of agents, you're over-forwarding.
Q: Do I need a real message broker, or can agents just call each other's functions?
For a prototype on one machine, direct async calls (Step 3) are fine and far simpler. Introduce a broker when you need durability (survive a crash mid-task), independent scaling, back-pressure, or cross-process communication. Don't add Kafka on day one — but design messages with envelopes from the start so swapping the transport later is a plumbing change, not a rewrite.

Key Takeaways

  1. 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.
  2. Envelope everything — message_id, correlation_id, causation_id, schema_version, and idempotency_key turn opaque LLM output into a traceable, recoverable, deduplicatable system.
  3. Validate at every boundary — parse agent output through a strict schema the moment it arrives; silent garbage propagates and fails three hops downstream.
  4. Shared state needs concurrency control — optimistic versioning prevents lost updates, and a termination fixpoint stops agents livelocking and billing forever.
  5. 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.
  6. 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