Multi-Agent Systems and Collaboration

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.

Three ways to wire one linkNamed recipientSender waitsCallee down, caller diesAnyone who caresFire and forgetNosubscriber, silent lossA worker, laterBuffered, durableBacklog grows unnoticedWho receivesCouplingHow it failsMessageEventQueueDirect HTTP between OCR, extraction and indexing means every agent is only as available as the slowest one.
Choose per link, not per system: the same pipeline can be a call here, an event there, and a queue where load spikes.

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.

Python
from dataclasses import dataclass, fieldfrom datetime import datetime, timezoneimport uuid@dataclassclass Message:    sender: str    recipient: str                 # the defining field    performative: str              # request | inform | query | propose | error    content: dict    conversation_id: str           # ties a request to its reply    message_id: str = field(default_factory=lambda: str(uuid.uuid4()))    reply_to: str | None = None    # message_id this answers    sent_at: str = field(default_factory=lambda: datetime.now(timezone.utc).isoformat())class MessageBus:    def __init__(self):        self.inboxes: dict[str, list[Message]] = {}    def register(self, agent_name: str):        self.inboxes.setdefault(agent_name, [])    def send(self, msg: Message):        if msg.recipient not in self.inboxes:            raise UnknownRecipient(msg.recipient)        self.inboxes[msg.recipient].append(msg)    def receive(self, agent_name: str) -> list[Message]:        msgs, self.inboxes[agent_name] = self.inboxes[agent_name], []        return msgs

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

PerformativeMeansReply expectedExample content
requestDo this actionYes — result or error{"action": "ocr", "doc_id": "c-882"}
informHere is a fact you may wantNo{"doc_id": "c-882", "pages": 412}
queryTell me something you knowYes — an inform{"question": "status of c-882"}
proposeI offer to do this at this costAccept or reject{"task": "ocr", "eta_s": 38}
errorYour request could not be servedNo{"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.

Python
from collections import defaultdictclass EventBus:    def __init__(self):        self.handlers = defaultdict(list)    def subscribe(self, event_type: str, handler):        self.handlers[event_type].append(handler)    def publish(self, event_type: str, payload: dict):        # Snapshot the handler list: a handler may subscribe during dispatch.        for handler in list(self.handlers[event_type]):            try:                handler(payload)            except Exception as exc:                # One bad subscriber must not stop the others.                self.publish("handler_failed",                             {"event_type": event_type, "error": str(exc)})bus = EventBus()bus.subscribe("document.ocr_completed", extraction_agent.on_ocr_done)bus.subscribe("document.ocr_completed", metrics.record_ocr)bus.subscribe("document.ocr_completed", audit_log.append)bus.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.

CategoryExamplesTypical subscribers
Lifecycleagent.started, agent.stopped, agent.heartbeatHealth monitor, supervisor
Task progresstask.assigned, task.completed, task.failedProgress tracker, retry manager, dashboard
Domaindocument.ocr_completed, clause.flaggedDownstream agents, search indexer
Resourcequota.exceeded, rate_limit.hitThrottler, 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.

Python
import json, redisr = redis.Redis(decode_responses=True)def enqueue(queue: str, task: dict):    r.lpush(queue, json.dumps(task))def worker_loop(queue: str, handler, visibility_timeout=300):    processing = f"{queue}:processing"    while True:        # Atomically move one task to a processing list so a crash        # mid-task does not lose it.        raw = r.blmove(queue, processing, timeout=5, src="RIGHT", dest="LEFT")        if raw is None:            continue        task = json.loads(raw)        try:            handler(task)            r.lrem(processing, 1, raw)          # done: remove        except Exception:            r.lrem(processing, 1, raw)            task["attempts"] = task.get("attempts", 0) + 1            if task["attempts"] < 3:                r.lpush(queue, json.dumps(task))          # retry            else:                r.lpush(f"{queue}:dead", json.dumps(task))  # give up, keep it

Two 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:

⌈5012⌉=⌈4.17⌉=5\left\lceil \frac{50}{12} \right\rceil = \lceil 4.17 \rceil = 5

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

GuaranteeMeansCostUse when
At-most-onceDelivered zero or one times; loss possibleCheapest, no acknowledgementMetrics, heartbeats — losing one is fine
At-least-onceNever lost, may be delivered twiceConsumer must be idempotentAlmost all agent work — the practical default
Exactly-onceDelivered precisely onceExpensive; usually at-least-once plus deduplicationPayments, 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:

Python
def handle_ocr(task):    doc_id = task["doc_id"]    # SETNX returns 1 only the first time this key is claimed.    if not r.set(f"done:ocr:{doc_id}", "1", nx=True, ex=86_400):        return                        # already processed; drop silently    text = run_ocr(doc_id)    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

SystemModelOrderingRetentionReach for it when
Redis lists / RQSimple list, pull-basedFIFO per listIn memory (optionally persisted)You already run Redis and want a queue this afternoon
RabbitMQBroker with exchanges and routing keysFIFO per queueUntil acknowledgedYou need routing rules, priorities, per-message TTL
KafkaAppend-only partitioned logFIFO per partitionTime- or size-based, replayableMany consumers, replay matters, very high volume
Amazon SQSManaged queueFIFO only in FIFO queuesUp to 14 daysYou want no broker to operate
Celery (on Redis or RabbitMQ)Task framework over a brokerBroker-dependentBroker-dependentYou 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

PropertyMessagesEventsQueues
Recipient known?Yes, namedNo, anonymousNo — any free worker
Sender blocks?Usually yesNoNo
Survives receiver being down?NoNoYes
Consumers per itemOneManyExactly one
Add a consumer without code change?NoYesYes (another worker)
Natural back-pressure signalTimeoutsNoneQueue depth
DebuggabilityHighLowMedium

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.