Course Content
Multi-Agent Systems and Collaboration
4 sections · 12 lessons
Coordination Strategies (Centralized vs Distributed)
A market-intelligence team ran six research agents against a list of 300 competitor URLs. There was no coordinator. Each agent picked a URL it thought looked interesting and fetched it. After the first run they checked the logs: 300 URLs on the list, 487 fetches performed, 172 URLs never touched at all. The 487 fetches covered only 128 distinct URLs, so nearly three-quarters of the work was duplicated, and more than half the list was missed.
The obvious fix was a coordinator: one agent hands out URLs, nobody picks their own. That worked, and the duplicate rate fell to zero. Then they scaled from 6 agents to 30, and throughput fell. The coordinator was making one assignment every 800 milliseconds — it had to read state, decide, write state, reply — and thirty agents finishing a page every five seconds were asking for six assignments per second. The coordinator could serve 1.25. Agents spent most of their lives waiting for permission to work.
Both failures are coordination failures, in opposite directions. Too little and agents collide. Too much and agents queue. Getting this right is not about picking the "better" architecture; it is about knowing which failure you are currently closer to.
What coordination actually has to solve
Coordination is the set of mechanisms that stop independent agents from interfering with each other and make their separate work add up to a coherent result. Concretely, it has to solve five problems, and every architecture below is an answer to some subset of them.
| Problem | Question it answers | Symptom when unsolved |
|---|---|---|
| Task allocation | Who does what? | Duplicated work and untouched work, as above |
| Resource contention | Who may use the shared thing now? | Rate-limit bans, corrupted files, lost writes |
| Ordering and dependency | What must finish before what? | Summarising a document before it was fetched |
| Consistency | Whose view of the state is right? | Two agents acting on contradictory beliefs |
| Termination detection | Are we done? | Runs that end only on a step limit or a timeout |
Termination detection is the one people forget. In a single-agent loop, "done" is obvious — the agent stops calling tools. With eight agents working concurrently, an agent being idle does not mean the system is finished; it may be about to receive work from a peer that is still thinking. Detecting global quiescence needs an explicit protocol, and skipping it is why so many multi-agent runs terminate on the step cap rather than on completion.
An idle agent is not evidence that the system has finished. Without an explicit termination rule, a multi-agent system does not know when to stop.
Centralised coordination
One agent holds the plan, the task list, and the authority to assign. Workers do not talk to each other; they talk to the coordinator.
┌──────────────────────────────┐ │ Coordinator │ │ task queue | assignments │ │ results | global budget │ └───┬──────────┬──────────┬────┘ │ │ │ ┌───▼───┐ ┌───▼───┐ ┌───▼───┐ │Worker1│ │Worker2│ │Worker3│ └───────┘ └───────┘ └───────┘1import threading2from collections import deque34class Coordinator:5 def __init__(self, tasks, token_budget=200_000):6 self.pending = deque(tasks)7 self.assigned: dict[str, str] = {} # task_id -> worker8 self.results: dict[str, dict] = {}9 self.tokens_spent = 010 self.token_budget = token_budget11 self.lock = threading.Lock()1213 def request_task(self, worker: str) -> dict | None:14 with self.lock:15 if self.tokens_spent >= self.token_budget:16 return None # global stop: only possible here17 if not self.pending:18 return None19 task = self.pending.popleft()20 self.assigned[task["id"]] = worker21 return task2223 def submit(self, worker: str, task_id: str, result: dict, tokens: int):24 with self.lock:25 self.results[task_id] = result26 self.tokens_spent += tokens27 self.assigned.pop(task_id, None)2829 def finished(self) -> bool:30 with self.lock:31 return not self.pending and not self.assignedLook at finished(). Termination detection is a two-line check, because one place knows both what is queued and what is in flight. That is the real argument for centralisation, and it is stronger than "it is simpler".
The strengths
- Global optimisation. The coordinator sees every task and every worker, so it can balance load, respect priorities, and enforce a total budget. A distributed system cannot enforce "spend no more than 200,000 tokens across everyone" without an extra protocol.
- No duplicated work by construction. A task is popped from the queue exactly once.
- Cheap communication. With
nworkers you havenlinks and roughly two messages per task — a request and a submission. - One place to debug. The coordinator's state is the system state.
The weaknesses, with the arithmetic
The bottleneck is not vague; you can compute it. Let the coordinator take D seconds per assignment decision and let each worker complete a task in T seconds. Each worker generates 1/T requests per second, so n workers generate n/T. The coordinator serves 1/D. Saturation is at:
With the market-intelligence numbers, T = 5 s and D = 0.8 s, so n_max = 6.25. Six workers is the ceiling. At 30 workers the arrival rate is 6 per second against a service rate of 1.25 — utilisation of 480% — so the request queue grows without bound and every worker's effective throughput collapses. Adding workers past n_max does not just fail to help; it makes latency worse for everyone already there.
Two fixes follow directly from the formula. Reduce D: batch assignments so one decision hands out ten URLs instead of one, cutting effective D to 0.08 s and raising the ceiling to 62 workers. Or increase T by giving workers bigger chunks. Both work; batching is usually the easier one.
The other weakness is availability. If the coordinator dies, in-flight assignments are lost and nothing new is handed out. Every worker is healthy and every worker is idle.
Distributed coordination
No agent holds authority. Agents decide locally, based on what they can observe and what peers tell them.
1class PeerAgent:2 def __init__(self, name, peers, shared_claims):3 self.name = name4 self.peers = peers5 self.claims = shared_claims # e.g. a Redis client67 def try_claim(self, task_id: str, ttl_s: int = 600) -> bool:8 # Atomic claim: exactly one agent's SET can succeed.9 return bool(self.claims.set(f"claim:{task_id}", self.name,10 nx=True, ex=ttl_s))1112 def work_loop(self, candidate_tasks):13 for task in candidate_tasks:14 if not self.try_claim(task["id"]):15 continue # someone else has it16 try:17 result = self.execute(task)18 self.broadcast("task.completed",19 {"id": task["id"], "by": self.name})20 except Exception:21 self.claims.delete(f"claim:{task['id']}") # release itThe claim key with nx=True is doing all the work here. It is a distributed lock: the first agent to set the key wins, everyone else's set returns false. The TTL matters just as much — if the winning agent dies, the claim expires after ten minutes and another agent can pick the task up. Without the TTL, a crashed agent parks that task forever.
The strengths
- No single point of failure. Lose an agent and the rest continue; its claims expire and get retried.
- Horizontal scaling. Nothing serialises through one process, so there is no
n_max. - Locality. An agent that already has a page cached can decide to process it without asking anyone.
The weaknesses
- No global optimum. Each agent makes a locally sensible choice. Together those can be poor — six agents all claiming cheap fast pages and leaving the expensive ones for last, so the tail is served by one agent alone.
- Message growth. A fully connected group of
nagents hasn(n-1)/2links. At 6 agents that is 15; at 20 it is 190. - Consistency is now your problem. Agents can hold different views of the world at the same instant. Reasoning about that is genuinely hard.
- Termination detection needs a protocol. There is no
finished(). You need something like: every agent broadcasts "I am idle with no pending sends"; the system is done only when all agents are idle simultaneously and no message is in flight.
| Dimension | Centralised | Distributed |
|---|---|---|
| Messages per task | ~2 | 1 claim attempt + 1 broadcast, but attempted by many |
| Scaling limit | T / D workers | Limited by the shared claim store |
| Failure of one node | Coordinator: fatal. Worker: retried | Any node: claims expire, work continues |
| Global budget enforcement | Trivial | Needs a shared counter and is approximate |
| Termination detection | Two-line check | Explicit protocol required |
| Debugging | Read the coordinator | Join traces across agents |
| Best at | Under ~10 agents, clear decomposition | Many agents, high fault tolerance, independent work |
Centralised coordination trades availability for control. Distributed coordination trades control for availability. There is no third option that gives you both for free.
Three coordination patterns worth knowing
The blackboard model
Agents never address each other. They read from and write to one shared structured store — the blackboard — and each agent watches for the conditions under which it can contribute.
1class Blackboard:2 def __init__(self):3 self.data = {}4 self.lock = threading.Lock()5 self.version = 067 def write(self, key: str, value, author: str):8 with self.lock:9 self.data[key] = {"value": value, "author": author,10 "version": self.version}11 self.version += 11213 def read(self, key: str):14 with self.lock:15 entry = self.data.get(key)16 return entry["value"] if entry else None1718class KnowledgeSource:19 """An agent that fires only when its preconditions are on the board."""20 requires: list[str] = []21 produces: str = ""2223 def can_contribute(self, bb: Blackboard) -> bool:24 return (all(bb.read(k) is not None for k in self.requires)25 and bb.read(self.produces) is None)2627def control_loop(bb, sources, max_rounds=50):28 for _ in range(max_rounds):29 ready = [s for s in sources if s.can_contribute(bb)]30 if not ready:31 return bb # quiescence: nothing can fire32 for s in ready:33 bb.write(s.produces, s.run(bb), author=type(s).__name__)The blackboard's virtue is that agents need no knowledge of each other at all — a new agent declares what it needs and what it produces, and the control loop schedules it. Its vice is the lost update: if two agents write the same key concurrently, the second overwrites the first with no error. The version counter above lets you detect it; the cleaner fix is to give each agent exclusive ownership of the keys it writes, or to make values append-only lists rather than replaceable scalars.
The Contract Net Protocol
A market. The agent with work broadcasts a call for proposals, capable agents bid, and the best bid wins the contract. It is the standard answer to "who should do this?" when the answer depends on runtime conditions rather than static roles.
1def contract_net(manager, task, contractors, deadline_s=2.0):2 # 1. Announce3 cfp = {"task": task, "deadline_s": deadline_s}4 # 2. Collect bids; agents that cannot do the task return None5 bids = []6 for c in contractors:7 bid = c.evaluate(cfp) # {"agent":…, "cost":…, "eta_s":…} or None8 if bid is not None:9 bids.append(bid)10 if not bids:11 return None # nobody capable: escalate12 # 3. Award to the best bid13 winner = min(bids, key=lambda b: b["cost"])14 winner["agent"].award(task)15 # 4. Reject the rest explicitly, so they free reserved capacity16 for b in bids:17 if b is not winner:18 b["agent"].reject(task)19 return winnerA worked round. Three contractors evaluate a "summarise 40 pages" task. Contractor A is idle and has the right model: cost 7.5. Contractor B has three tasks queued: cost 12.0. Contractor C lacks the tool: no bid. The manager awards to A at 7.5 and sends B an explicit rejection so B does not hold capacity in reserve.
Count the messages: 1 broadcast + 2 bids + 1 award + 1 rejection = 5, versus 2 for a straight centralised assignment. Contract Net costs roughly 2n + 1 messages for n contractors. You pay that to get an allocation that reflects current load and capability instead of a static rota. On tasks that take seconds, that overhead is invisible; on tasks that take 40 milliseconds, it dominates.
Step 4 is the one people delete because it "does nothing". It matters: contractors typically reserve capacity when they bid, and without an explicit rejection they hold that reservation until a timeout, so the system's usable capacity quietly shrinks under load.
Voting and consensus
When several agents produce candidate answers and you need one, aggregate rather than pick arbitrarily.
1from collections import Counter23def majority_vote(votes: list[str], threshold: float = 0.5):4 counts = Counter(votes)5 winner, n = counts.most_common(1)[0]6 return winner if n / len(votes) > threshold else None78def weighted_vote(votes: dict[str, str], weights: dict[str, float]):9 scores = Counter()10 for agent, choice in votes.items():11 scores[choice] += weights.get(agent, 1.0)12 return scores.most_common(1)[0][0]A worked example. Five agents classify a filing as material or routine: three say material, two say routine. Plain majority gives 3/5 = 0.6 > 0.5, so material. Now weight by historical accuracy — the two "routine" voters have accuracy 0.94 and 0.91, the three "material" voters 0.62, 0.58 and 0.71. Weighted: material scores 0.62 + 0.58 + 0.71 = 1.91; routine scores 0.94 + 0.91 = 1.85. Material still wins, but by 1.91 to 1.85 rather than 3 to 2 — a margin of 3% instead of 50%. That thin margin is the signal to escalate to a human, and unweighted voting would have hidden it completely.
Quorum is the related idea for availability: require agreement from a majority of nodes rather than all of them, so the system keeps deciding while a minority is down. With 5 agents a quorum is 3, which means the system tolerates 2 failures. With 4 agents a quorum is still 3, so it tolerates only 1 — which is why odd numbers are conventional.
Named failure modes
The coordinator as accidental bottleneck. Symptom: adding workers stops helping and then starts hurting; worker utilisation drops as the fleet grows. Cause: n > T/D. Fix: batch assignments, or shard the coordinator by task type.
The claim without a TTL. Symptom: a handful of tasks never complete and never error; their claim keys exist with no owner alive. Cause: an agent crashed holding a lock that never expires. Fix: always set an expiry, and set it to a few multiples of the expected task duration, not to the exact duration — a task running slightly long must not have its lock stolen mid-flight.
Lost updates on shared state. Symptom: an agent's contribution is simply missing from the final output, with no error anywhere. Cause: two agents wrote the same blackboard key. Fix: key ownership, or append-only values, or compare-and-swap on a version number.
Distributed deadlock. Agent A holds the rate-limit token and waits for a page that B is fetching; B waits for the token. Neither yields. Symptom: throughput hits zero with all agents "busy". Fix: acquire shared resources in a globally fixed order, and put a timeout on every acquire so a cycle breaks itself.
Contract Net for trivial tasks. Symptom: coordination messages outnumber work messages ten to one. Cause: a bidding protocol applied to tasks shorter than the bidding round. Fix: use it only where allocation quality actually varies — long tasks, heterogeneous agents, uneven load.
What this means when you build
Start centralised, and measure D — the coordinator's decision time — on day one. It is one timer around the assignment function. That number, divided into your average task duration, is your worker ceiling, and knowing it before you scale saves you the experience of watching throughput fall as you add capacity.
Choose a hybrid before you choose pure peer-to-peer. In practice the shape that works is a coordinator that owns allocation and budget — the decisions that need a global view — while agents talk directly for everything else, such as one agent asking another a clarifying question. That keeps the coordinator's decision rate low because it is not in the path of every interaction.
Write termination detection deliberately, whichever architecture you pick. Centralised, it is not pending and not assigned. Distributed, it is an idle-broadcast protocol with in-flight message accounting. If you cannot point at the line of code that decides the system is finished, your system does not know it is finished, and you will discover this when a run costs eleven times its budget because two agents kept politely handing a task back and forth.
Finally, treat every shared resource as needing an explicit protocol. A rate-limited API, a file, a blackboard key, a database row: name it, decide who may write it, and decide what happens on conflict. The 487 fetches for 300 URLs were not a bug in any agent. They were the absence of a decision about who owns a URL.