Course Content
Multi-Agent Systems and Collaboration
4 sections · 12 lessons
Communication Protocols (Messages, Events, Queues)
A legal-tech team ran a three-agent document pipeline: an OCR agent that turned scanned contracts into text, an extraction agent that pulled out parties, dates and clauses, and an indexing agent that wrote the results into a search store. The agents talked over direct HTTP calls. Extraction called OCR and waited for the answer; indexing called extraction and waited for that.
It worked for months. Then a client uploaded a batch of 400-page scanned faxes. OCR, which normally returned in 4 seconds, started taking 40. The extraction agent's HTTP client had a 30-second timeout, so it raised, so indexing's call to extraction raised too. About 12% of that day's 12,000 documents — 1,440 of them — simply vanished. No queue held them. No log said which ones. The only record was a spike in a 504 counter.
The team's first instinct was to raise the timeout to 120 seconds. That made things worse: now every extraction worker sat blocked for two minutes holding a connection, so the pool exhausted and fast documents started failing too. The real problem was never the timeout. It was that the only way these agents could talk was "call and wait", and that shape of communication cannot survive a slow neighbour.
Why the wiring is the architecture
People spend weeks on agent prompts and ten minutes on how agents exchange information. That is backwards. The communication mechanism decides:
- What happens when a receiver is slow or dead. Does the sender block, fail, or carry on?
- Whether work can be lost. If a process is killed mid-task, does the task come back?
- How you add a fourth agent. Do you edit the existing three, or none of them?
- Whether you can scale one stage independently. Can you run six OCR workers and two indexers?
Three mechanisms cover the field: direct messages, events, and queues. They are not competitors — most real systems use all three, in different places, for different reasons. The skill is knowing which one a given link needs.
Message passing: I am talking to you specifically
A message is addressed. The sender names the recipient. That single property — a known addressee — is what distinguishes it from the other two mechanisms.
1from dataclasses import dataclass, field2from datetime import datetime, timezone3import uuid45@dataclass6class Message:7 sender: str8 recipient: str # the defining field9 performative: str # request | inform | query | propose | error10 content: dict11 conversation_id: str # ties a request to its reply12 message_id: str = field(default_factory=lambda: str(uuid.uuid4()))13 reply_to: str | None = None # message_id this answers14 sent_at: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat())1516class MessageBus:17 def __init__(self):18 self.inboxes: dict[str, list[Message]] = {}1920 def register(self, agent_name: str):21 self.inboxes.setdefault(agent_name, [])2223 def send(self, msg: Message):24 if msg.recipient not in self.inboxes:25 raise UnknownRecipient(msg.recipient)26 self.inboxes[msg.recipient].append(msg)2728 def receive(self, agent_name: str) -> list[Message]:29 msgs, self.inboxes[agent_name] = self.inboxes[agent_name], []30 return msgsMessage types worth distinguishing
The performative field is not decoration. It tells the receiver what kind of reply, if any, is expected — which is exactly the thing that goes wrong when agents talk in free-form text. The vocabulary below comes from agent-communication research and has survived because each entry earns its place.
| Performative | Means | Reply expected | Example content |
|---|---|---|---|
request | Do this action | Yes — result or error | {"action": "ocr", "doc_id": "c-882"} |
inform | Here is a fact you may want | No | {"doc_id": "c-882", "pages": 412} |
query | Tell me something you know | Yes — an inform | {"question": "status of c-882"} |
propose | I offer to do this at this cost | Accept or reject | {"task": "ocr", "eta_s": 38} |
error | Your request could not be served | No | {"code": "TIMEOUT", "retryable": true} |
Without an explicit type, the receiving agent must infer intent from prose, and it will sometimes infer wrong — answering a statement as if it were a question, or acting on a proposal that was never accepted.
What message passing buys and costs
It buys clarity and traceability. Every message has a sender, a recipient and a conversation ID, so reconstructing "what did agent A ask agent B at 14:02" is a single filter. It also buys strict request/response semantics, which are easy to reason about and easy to test.
It costs coupling and availability. The sender must know the recipient exists and be able to name it. In the legal-tech pipeline, extraction had to know OCR's address, and when OCR was slow, extraction had nowhere to put the work.
Direct messaging couples the sender's availability to the receiver's. If the receiver is down, the sender's only options are to block, to fail, or to invent a queue.
Events: something happened, whoever cares
An event is a statement of fact about the past, published without an addressee. The publisher does not know or care who is listening. Subscribers register interest in an event type.
1from collections import defaultdict23class EventBus:4 def __init__(self):5 self.handlers = defaultdict(list)67 def subscribe(self, event_type: str, handler):8 self.handlers[event_type].append(handler)910 def publish(self, event_type: str, payload: dict):11 # Snapshot the handler list: a handler may subscribe during dispatch.12 for handler in list(self.handlers[event_type]):13 try:14 handler(payload)15 except Exception as exc:16 # One bad subscriber must not stop the others.17 self.publish("handler_failed",18 {"event_type": event_type, "error": str(exc)})1920bus = EventBus()21bus.subscribe("document.ocr_completed", extraction_agent.on_ocr_done)22bus.subscribe("document.ocr_completed", metrics.record_ocr)23bus.subscribe("document.ocr_completed", audit_log.append)2425bus.publish("document.ocr_completed", {"doc_id": "c-882", "chars": 184_233})Look at what just happened. The OCR agent published one line. Three separate consumers reacted. Adding a fourth — a compliance sampler, say — means writing one subscribe call and touching no existing agent. That is the entire value proposition of events: you can add consumers without modifying producers.
Naming events correctly
Events are named in the past tense, as facts, not as instructions. document.ocr_completed is an event. extract_document is a command wearing an event's clothing, and it will bite you: the moment two subscribers both handle extract_document, the document gets extracted twice.
| Category | Examples | Typical subscribers |
|---|---|---|
| Lifecycle | agent.started, agent.stopped, agent.heartbeat | Health monitor, supervisor |
| Task progress | task.assigned, task.completed, task.failed | Progress tracker, retry manager, dashboard |
| Domain | document.ocr_completed, clause.flagged | Downstream agents, search indexer |
| Resource | quota.exceeded, rate_limit.hit | Throttler, alerting |
The two traps
Nobody is listening, and nothing tells you. A message to an unknown recipient raises. An event with no subscribers is silently dropped. A team once renamed ocr_completed to ocr.completed in the publisher and forgot the subscriber; the pipeline reported zero errors and processed zero documents for six hours. Guard against it: assert at start-up that every event type your system publishes has at least one subscriber.
Cascades are invisible. Event A triggers handler B which publishes event C which triggers handler D. Nothing in the code shows this chain. Reading the publisher tells you nothing about what will actually happen. This is the price of decoupling, and it is why event-driven systems need tracing more urgently than any other kind.
Events remove the coupling between who produces information and who uses it — and with it, they remove the code path you would have read to understand the system.
Queues: durable, buffered, and load-levelling
A queue sits between producer and consumer as a persistent buffer. The producer writes and moves on. Consumers pull work when they have capacity. Neither needs the other to be alive at the same moment.
This is the mechanism that would have saved the 1,440 documents. With a queue, slow OCR does not fail extraction — it just means the extraction queue grows, and shrinks again when OCR recovers.
1import json, redis23r = redis.Redis(decode_responses=True)45def enqueue(queue: str, task: dict):6 r.lpush(queue, json.dumps(task))78def worker_loop(queue: str, handler, visibility_timeout=300):9 processing = f"{queue}:processing"10 while True:11 # Atomically move one task to a processing list so a crash12 # mid-task does not lose it.13 raw = r.blmove(queue, processing, timeout=5, src="RIGHT", dest="LEFT")14 if raw is None:15 continue16 task = json.loads(raw)17 try:18 handler(task)19 r.lrem(processing, 1, raw) # done: remove20 except Exception:21 r.lrem(processing, 1, raw)22 task["attempts"] = task.get("attempts", 0) + 123 if task["attempts"] < 3:24 r.lpush(queue, json.dumps(task)) # retry25 else:26 r.lpush(f"{queue}:dead", json.dumps(task)) # give up, keep itTwo details in that loop matter more than the rest. blmove (the current replacement for the older brpoplpush) moves the task to a processing list atomically, so a worker killed halfway through leaves a recoverable trace rather than a hole. And the dead-letter list means a task that fails three times is parked for inspection rather than retried forever or silently dropped.
Sizing a queue with real numbers
Queues make capacity planning arithmetic instead of guesswork. Suppose documents arrive at 50 per minute and one OCR worker completes 12 per minute. Workers needed:
Run 5 workers and you have 60 per minute of capacity against 50 of demand — 17% headroom. Run 4 and you have 48 against 50: the backlog grows by 2 per minute, which is 120 per hour, which is 2,880 documents by the next morning. Nothing errors. Nothing alerts. The queue depth just climbs, and by the time somebody notices, every new document waits about an hour (2,880 queued ÷ 48 per minute = 60 minutes) before a worker touches it — and the wait grows every minute.
That is why queue depth and queue age are the two metrics you alert on. Depth tells you the backlog; age tells you how bad it is for the unlucky item at the front. A depth of 3,000 with an age of 40 seconds is a healthy burst. A depth of 120 with an age of 3 hours means a poison message is stuck at the head.
Delivery guarantees, stated honestly
| Guarantee | Means | Cost | Use when |
|---|---|---|---|
| At-most-once | Delivered zero or one times; loss possible | Cheapest, no acknowledgement | Metrics, heartbeats — losing one is fine |
| At-least-once | Never lost, may be delivered twice | Consumer must be idempotent | Almost all agent work — the practical default |
| Exactly-once | Delivered precisely once | Expensive; usually at-least-once plus deduplication | Payments, anything double-charging matters |
The honest statement is this: end-to-end exactly-once does not exist across a network. What real systems do is at-least-once delivery plus an idempotent consumer, which produces exactly-once effects. Idempotent means running the handler twice with the same input has the same result as running it once. Practically:
1def handle_ocr(task):2 doc_id = task["doc_id"]3 # SETNX returns 1 only the first time this key is claimed.4 if not r.set(f"done:ocr:{doc_id}", "1", nx=True, ex=86_400):5 return # already processed; drop silently6 text = run_ocr(doc_id)7 store(doc_id, text)Without that guard, a worker that crashes after doing the work but before acknowledging causes the task to be redelivered and the work to be done twice. With OCR that wastes money. With "send the client an email" it is a support ticket.
Which queue system
| System | Model | Ordering | Retention | Reach for it when |
|---|---|---|---|---|
| Redis lists / RQ | Simple list, pull-based | FIFO per list | In memory (optionally persisted) | You already run Redis and want a queue this afternoon |
| RabbitMQ | Broker with exchanges and routing keys | FIFO per queue | Until acknowledged | You need routing rules, priorities, per-message TTL |
| Kafka | Append-only partitioned log | FIFO per partition | Time- or size-based, replayable | Many consumers, replay matters, very high volume |
| Amazon SQS | Managed queue | FIFO only in FIFO queues | Up to 14 days | You want no broker to operate |
| Celery (on Redis or RabbitMQ) | Task framework over a broker | Broker-dependent | Broker-dependent | You want chains, groups, scheduling and retries built in |
One property that trips people: Kafka guarantees order only within a partition. If document c-882's "ocr_completed" and "extraction_completed" land in different partitions, a consumer can see them out of order. The fix is to partition by a key that keeps related events together — partition_key = doc_id — which is a one-line decision with system-wide consequences.
Choosing, per link, not per system
| Property | Messages | Events | Queues |
|---|---|---|---|
| Recipient known? | Yes, named | No, anonymous | No — any free worker |
| Sender blocks? | Usually yes | No | No |
| Survives receiver being down? | No | No | Yes |
| Consumers per item | One | Many | Exactly one |
| Add a consumer without code change? | No | Yes | Yes (another worker) |
| Natural back-pressure signal | Timeouts | None | Queue depth |
| Debuggability | High | Low | Medium |
Use messages when one specific agent must answer and you need the answer before continuing: a supervisor asking a worker "can you take this task?", or an agent querying a registry. Use events when a fact might interest an unknown number of parties: task completed, threshold crossed, agent joined. Use queues when work must not be lost, when producers and consumers run at different speeds, or when you want to scale a stage by adding processes.
Named failure modes
Request/response over a slow link. The opening story. Symptom: cascading timeouts and connection-pool exhaustion, where raising the timeout makes it worse. Cause: synchronous messaging where the work is long and variable. Fix: put a queue in front of the slow stage and return a task ID immediately.
Events used as commands. Publishing process_document and having exactly one subscriber act on it. Symptom: the day someone adds a second subscriber "just to log it", documents get processed twice. Cause: event semantics are broadcast; commands are point-to-point. Fix: name events in the past tense as facts, and send commands as addressed messages or queued tasks.
Non-idempotent consumers on at-least-once delivery. Symptom: duplicate database rows, duplicate emails, doubled billing — always after a worker restart or a network hiccup. Cause: assuming delivery is exactly-once. Fix: a claim key, as shown above, or a natural unique constraint in the database.
The unbounded queue. No depth limit, no alert. Symptom: memory exhaustion on the broker, or a backlog that takes days to drain. Cause: no back-pressure. Fix: set a maximum depth and decide explicitly what happens on overflow — reject new work with a clear error, shed low-priority work, or scale consumers automatically.
The poison message. One malformed task fails, is retried, fails, is retried, forever — while queue depth looks fine and everything behind it starves. Symptom: high queue age with low depth. Fix: attempt counters and a dead-letter queue, exactly as in the worker loop above.
What this means when you build
Draw your agents as boxes and each communication link as an arrow, then label every arrow with three answers: does the sender need a reply?, may this be lost?, and how many parties care?. Those three answers pick the mechanism, and they pick it per arrow. A realistic pipeline uses all three: a queue from the API into OCR because bursts must not be dropped, an event when OCR finishes because extraction, metrics and audit all care, and a direct message from extraction to a validation agent because it needs an answer before it can continue.
Whatever mechanism you pick, carry a correlation ID through all of them. One identifier, generated when the work enters the system, attached to every message, event and task derived from it. Without it, an incident like the vanished 1,440 documents leaves you with three separate log streams and no way to join them. With it, one grep reconstructs the whole journey.
And put explicit envelopes around everything. The Message dataclass above, the past-tense event names, the task dictionary with an attempts counter — these are not ceremony. They are the difference between a system where you can answer "where did document c-882 go?" in thirty seconds and one where you cannot answer it at all.