Course Content
Multi-Agent Systems and Collaboration
4 sections · 12 lessons
The Supervisor-Worker Pattern
An operations team had a nightly job: read every support ticket filed that day, classify it, draft a reply, and flag anything needing a human. One agent, one loop, 6.2 seconds per ticket. At 200 tickets a night that is 20.7 minutes, comfortably inside the four-hour window.
Then the company grew and the nightly volume hit 2,000 tickets. 2000 × 6.2 = 12,400 seconds — 3 hours 27 minutes — and any hiccup blew the window entirely. The obvious fix was to run eight copies of the agent. So they did, on eight machines, each reading the same ticket list from the database.
Each of the eight processed all 2,000 tickets. The job produced 16,000 draft replies, 8 duplicate replies per ticket, and an API bill eight times larger than the previous night. Nothing crashed. Every agent did exactly what it was told.
What was missing was not compute. It was someone to decide who does what. That role has a name.
The problem the pattern actually solves
Running n copies of an agent gives you n times the work, not n times the throughput, unless something partitions the work. The supervisor-worker pattern introduces one agent whose entire job is that partitioning: it decomposes work, assigns each unit to exactly one worker, tracks what is outstanding, and handles what fails.
Note carefully what it does not solve. It does not make a sequential task parallel. If ticket 2 genuinely cannot be classified until ticket 1's outcome is known, no supervisor helps. The pattern applies when the units of work are independent of each other, which for 2,000 unrelated support tickets they obviously are.
Adding workers without adding a supervisor does not divide the work. It multiplies it.
Anatomy
Three parts, and the third is the one people skip.
The task record
Work must be represented as data with a status, not as a function call. A function call has no memory: if it fails you have nothing to retry, and if the process dies you cannot tell which calls completed. A task record survives.
1from dataclasses import dataclass, field2from enum import Enum3import time, uuid45class Status(str, Enum):6 PENDING = "pending" # created, not yet given to anyone7 ASSIGNED = "assigned" # given to a worker, not yet started8 RUNNING = "running" # worker has begun9 COMPLETED = "completed"10 FAILED = "failed" # exhausted its retries11 DEAD = "dead" # parked for a human1213@dataclass14class Task:15 kind: str # which worker capability is needed16 payload: dict17 id: str = field(default_factory=lambda: str(uuid.uuid4()))18 status: Status = Status.PENDING19 assigned_to: str | None = None20 attempts: int = 021 max_attempts: int = 322 result: dict | None = None23 error: str | None = None24 heartbeat_at: float | None = None # last sign of life from the worker25 depends_on: list[str] = field(default_factory=list)Every field there exists because of a specific failure. attempts stops infinite retries. heartbeat_at lets the supervisor detect a worker that died without saying so. depends_on stops task B being dispatched before task A produced the thing B needs. DEAD is separate from FAILED so that "we gave up" is visible rather than being silently mixed into "it errored once".
The supervisor
It owns the task list and never executes domain work. The moment a supervisor starts drafting a reply itself, it is a worker with extra responsibilities, and it becomes the bottleneck.
1import threading23class Supervisor:4 def __init__(self, strategy="least_loaded", stale_after=90.0):5 self.tasks: dict[str, Task] = {}6 self.workers: dict[str, "Worker"] = {}7 self.load: dict[str, int] = {} # worker -> active task count8 self.strategy = strategy9 self.stale_after = stale_after10 self.lock = threading.Lock()11 self._rr = 01213 def register(self, worker):14 self.workers[worker.name] = worker15 self.load[worker.name] = 01617 def add_task(self, task: Task):18 self.tasks[task.id] = task1920 # ---- selection -------------------------------------------------21 def _eligible(self, task: Task) -> list[str]:22 return [n for n, w in self.workers.items() if task.kind in w.can_do]2324 def _pick(self, task: Task) -> str | None:25 candidates = self._eligible(task)26 if not candidates:27 return None28 if self.strategy == "round_robin":29 self._rr += 130 return candidates[self._rr % len(candidates)]31 if self.strategy == "least_loaded":32 return min(candidates, key=lambda n: self.load[n])33 if self.strategy == "capability":34 return max(candidates, key=lambda n: self.workers[n].skill[task.kind])35 raise ValueError(self.strategy)3637 def _ready(self, task: Task) -> bool:38 return all(self.tasks[d].status is Status.COMPLETED39 for d in task.depends_on)4041 # ---- the main loop ---------------------------------------------42 def dispatch_round(self):43 with self.lock:44 self._reclaim_stale()45 for task in self.tasks.values():46 if task.status is not Status.PENDING or not self._ready(task):47 continue48 worker = self._pick(task)49 if worker is None:50 continue # no capable worker right now51 task.status = Status.ASSIGNED52 task.assigned_to = worker53 task.heartbeat_at = time.time()54 self.load[worker] += 155 self.workers[worker].submit(task)5657 def _reclaim_stale(self):58 now = time.time()59 for task in self.tasks.values():60 if task.status in (Status.ASSIGNED, Status.RUNNING) \61 and now - (task.heartbeat_at or 0) > self.stale_after:62 self.load[task.assigned_to] = max(0, self.load[task.assigned_to] - 1)63 task.status = Status.PENDING # give it to someone else64 task.assigned_to = None6566 def on_complete(self, task_id, result):67 with self.lock:68 t = self.tasks[task_id]69 t.status, t.result = Status.COMPLETED, result70 self.load[t.assigned_to] -= 17172 def on_fail(self, task_id, error):73 with self.lock:74 t = self.tasks[task_id]75 self.load[t.assigned_to] -= 176 t.attempts += 177 t.error = error78 t.status = (Status.PENDING if t.attempts < t.max_attempts79 else Status.DEAD)80 t.assigned_to = None8182 def finished(self) -> bool:83 with self.lock:84 return all(t.status in (Status.COMPLETED, Status.DEAD)85 for t in self.tasks.values())The worker
A worker declares what it can do, executes one task, and reports back. It holds no view of the overall job.
1class Worker:2 def __init__(self, name, can_do: set[str], skill: dict[str, float],3 supervisor):4 self.name, self.can_do, self.skill = name, can_do, skill5 self.supervisor = supervisor67 def submit(self, task: Task):8 threading.Thread(target=self._run, args=(task,), daemon=True).start()910 def _run(self, task: Task):11 task.status = Status.RUNNING12 try:13 result = self.execute(task)14 self.supervisor.on_complete(task.id, result)15 except Exception as exc:16 self.supervisor.on_fail(task.id, f"{type(exc).__name__}: {exc}")1718 def execute(self, task: Task) -> dict:19 # A real worker calls a model, a tool, an API.20 raise NotImplementedErrorWiring it together is then unremarkable, which is the point:
1sup = Supervisor(strategy="least_loaded")2for i in range(8):3 sup.register(TicketWorker(f"w{i}", {"classify_ticket"}, {"classify_ticket": 1.0}, sup))4for ticket in todays_tickets: # 2,000 of them5 sup.add_task(Task(kind="classify_ticket", payload={"ticket": ticket}))67while not sup.finished():8 sup.dispatch_round()9 time.sleep(0.05)When to use it, and when not to
| Situation | Supervisor-worker? | Why |
|---|---|---|
| Many independent items of the same kind | Yes | Exactly what it is for; scales linearly until the supervisor saturates |
| Heterogeneous subtasks needing different tools | Yes | Capability routing sends each to the right specialist |
| Work whose size is unknown until runtime | Yes | The supervisor can add tasks as they are discovered |
| A strictly sequential chain (A then B then C) | No | Nothing overlaps; you pay coordination for zero parallelism |
| Two agents that negotiate as equals | No | There is no authority to centralise; use a peer protocol |
| Tasks under ~100 ms each | No | Dispatch overhead exceeds the work |
| You cannot tolerate any single point of failure | Not alone | The supervisor is one; you need failover or a claim-based design |
Letting a model do the decomposition
So far the supervisor received a ready-made list. Often the input is a goal in prose — "produce a competitive brief on the European payments market" — and decomposition itself is the hard part. A model can do it, provided you constrain the output hard enough to be executable.
1DECOMPOSE_PROMPT = """Break the goal into 3-8 independent subtasks.2Return ONLY JSON: a list of objects with keys3 "kind" : one of {allowed_kinds}4 "payload": object with the arguments that kind needs5 "depends_on": list of zero-based indices of subtasks that must finish first6Rules: no subtask may depend on a later index; prefer independence.78GOAL: {goal}"""910def decompose(goal: str, allowed_kinds: set[str], model) -> list[Task]:11 raw = model.generate(DECOMPOSE_PROMPT.format(12 goal=goal, allowed_kinds=sorted(allowed_kinds)))13 specs = json.loads(extract_json_block(raw))1415 tasks = [Task(kind=s["kind"], payload=s["payload"]) for s in specs]16 for i, s in enumerate(specs):17 if s["kind"] not in allowed_kinds: # hallucinated worker18 raise InvalidPlan(f"unknown kind {s['kind']!r}")19 for d in s.get("depends_on", []):20 if d >= i: # forward or self edge21 raise InvalidPlan(f"subtask {i} depends on {d}")22 tasks[i].depends_on.append(tasks[d].id)23 return tasksThe two validation checks are not optional. Models invent worker kinds that sound plausible ("kind": "verify_sources" when no such worker exists), and they produce dependency cycles. Both are silent disasters: an unknown kind means the task is never eligible for any worker and the run never finishes, and a cycle means _ready() is false forever. Validating at plan time turns a hang into an exception with a message.
A plan produced by a model is untrusted input. Validate it against the workers you actually have before a single task is dispatched.
Work distribution strategies, compared with numbers
Three tasks are long and six are short, arriving in this order: durations [10, 2, 2, 10, 2, 2, 10, 2, 2] seconds, three workers. Total work is 3×10 + 6×2 = 42 seconds, so a perfect split would finish in 14.
Round-robin
Cycle through workers regardless of load. Worker 1 gets tasks 1, 4, 7 — the three ten-second tasks — for 30 seconds. Workers 2 and 3 get 6 seconds each and then sit idle for 24. Makespan: 30 seconds. Two thirds of the fleet is idle for 80% of the run.
Least-loaded
Assign each task to whoever has least queued work. Tracing it: W1 takes the first 10 (loads 10/0/0), W2 and W3 take the two shorts (10/2/2), the second 10 goes to W2 (10/12/2), the next two shorts to W3 (10/12/6), the third 10 to W3 (10/12/16), the next short to W1 (12/12/16), and the last short to W1 or W2, which are tied (14/12/16). Makespan: 16 seconds — a 47% improvement over round-robin, and within 14% of the theoretical optimum of 14.
Capability-based
Route by who is best at this kind of work. Suppose a summarisation task takes the specialist worker 4 seconds and the generalist 16. Round-robin gives it to whoever is next, so it lands on the generalist half the time: expected duration 0.5×4 + 0.5×16 = 10 seconds. Capability routing always picks the specialist: 4 seconds, a 60% reduction in expected time — but it also concentrates all summarisation on one worker, which is exactly the round-robin pathology in a different costume.
| Strategy | Needs to know | Makespan on the example | Fails when |
|---|---|---|---|
| Round-robin | Nothing | 30 s | Task durations vary |
| Least-loaded | Active count per worker | 16 s | Count is a poor proxy for work (one huge task vs three tiny) |
| Capability | A skill score per kind | Depends on skill spread | One specialist becomes a queue |
| Capability then least-loaded | Both | Best in practice | Needs both signals maintained |
The production answer is the last row: filter to workers that can do the task, then among those pick the least loaded. That is two lines of code and it removes both pathologies. Better still, make "load" the sum of estimated durations rather than a count, so one 10-second task does not look identical to one 2-second task.
Handling worker failure
Workers fail in three distinct ways and each needs a different response.
| Failure | How the supervisor sees it | Correct response |
|---|---|---|
| Task raised an exception | on_fail called with an error | Retry with backoff up to max_attempts, then mark DEAD |
| Worker process died | Heartbeat goes stale; no callback ever arrives | Reclaim after stale_after, reassign to another worker |
| Worker alive but stuck | Heartbeats continue; task never completes | A per-task deadline independent of the heartbeat |
The second case is the one that silently breaks naive implementations. If the supervisor only reacts to callbacks, a worker that is killed mid-task produces no callback at all, and that task sits in ASSIGNED forever while finished() never returns true. The _reclaim_stale sweep exists solely for this.
Choosing stale_after matters. Set it to the expected task duration and you will reclaim tasks that are merely slow, so the work is done twice. Set it to three to five times the 99th-percentile duration. If p99 is 18 seconds, 90 seconds is a reasonable choice — long enough that a slow task is safe, short enough that a dead worker is noticed within a minute and a half.
Retries need exponential backoff with jitter, not immediate re-dispatch. If a downstream API is rate-limiting you, retrying instantly makes it worse:
1import random2def backoff_seconds(attempt: int, base: float = 2.0, cap: float = 60.0):3 # attempt 1 -> ~2 s, 2 -> ~4 s, 3 -> ~8 s, with +/-25% jitter4 delay = min(cap, base * (2 ** (attempt - 1)))5 return delay * random.uniform(0.75, 1.25)The jitter matters when many tasks fail at once — a rate limit typically fails dozens simultaneously — because without it they all retry at the same instant and trip the limit again. This is called the thundering herd, and jitter is the entire fix.
Reassignment is only safe if workers are idempotent: doing the task twice must have the same effect as doing it once. If classify_ticket also posts the reply to the customer, a reclaimed-but-actually-alive task sends two replies. Make the side effect conditional on a claim key derived from the task ID, or separate "compute the reply" (safely repeatable) from "send the reply" (done once, at the end, by the supervisor).
Where people get it wrong
"The supervisor should also do some of the work when it is idle." Tempting, and it destroys the pattern. While the supervisor is running a 6-second task it dispatches nothing, so every worker that finishes in that window sits idle. Supervisor decision time D sets your worker ceiling at T/D; with T = 6.2 s and a lean D = 0.15 s that ceiling is 41 workers, and giving the supervisor real work collapses it to roughly one.
"More workers is always faster." Only until you hit a shared limit. Eight workers against an API allowing 10 requests per second get 1.25 each, which is fine; forty workers get 0.25 each and spend most of their time being throttled, while your 429 rate climbs and your retry storm burns budget. Throughput is capped by the scarcest shared resource, not by worker count.
"Retries handle failures." Retries handle transient failures. A malformed payload fails identically three times and consumes three times the tokens on the way to the same outcome. Classify errors before retrying: network timeouts and 429s are retryable, schema violations and 400s are not, and marking the second class DEAD immediately saves both time and money.
"Every task is independent." Usually true, occasionally catastrophically false. If two tasks write the same output file or the same database row, running them concurrently corrupts it. The depends_on field is how you say so — and if the true relationship is mutual exclusion rather than ordering, you need a lock on the resource, not a dependency between tasks.
What this means when you build
Build the task record first, before the supervisor and long before the workers. Almost everything hard about this pattern — retries, reclaiming, dependency ordering, giving up gracefully — is a property of that record. If your task is a function call rather than a row of data, none of it is available to you, and you will find out on the night the supervisor restarts with 900 tasks in flight.
Persist the tasks. An in-memory dictionary is fine for a demo and wrong for the nightly ticket job: a supervisor restart at 02:00 loses every assignment and either redoes 2,000 tickets or drops them. A single database table with the same columns as the dataclass makes restart a matter of loading rows where status is not COMPLETED.
Instrument three numbers from the start: tasks by status over time, per-worker active count, and the age of the oldest non-completed task. The first tells you whether the run is converging; the second tells you whether distribution is actually balancing (round-robin's pathology is instantly visible as one worker pinned and two flat); the third catches the reclaim bug, because a task stuck in ASSIGNED shows up as an oldest-age that climbs forever while everything else looks healthy.
And keep the supervisor boring. It should contain no prompts, no tool calls, and no domain knowledge beyond "which kinds exist". Every piece of intelligence you move into the supervisor is a piece of intelligence you cannot scale, because there is exactly one of it.