Multi-Agent Systems and Collaboration

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.

One supervisor, four workers, one queueSupervisorholds the task listWorker 1 — ticket batchWorker 2 — ticket batchWorker 3 — ticket batchWorker 4 — ticket batchFailed tasksreturn to the list
Workers stay stateless so a dead worker costs one task, not the run — the supervisor owns all the state worth keeping.

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.

Python
from dataclasses import dataclass, fieldfrom enum import Enumimport time, uuidclass Status(str, Enum):    PENDING   = "pending"      # created, not yet given to anyone    ASSIGNED  = "assigned"     # given to a worker, not yet started    RUNNING   = "running"      # worker has begun    COMPLETED = "completed"    FAILED    = "failed"       # exhausted its retries    DEAD      = "dead"         # parked for a human@dataclassclass Task:    kind: str                          # which worker capability is needed    payload: dict    id: str = field(default_factory=lambda: str(uuid.uuid4()))    status: Status = Status.PENDING    assigned_to: str | None = None    attempts: int = 0    max_attempts: int = 3    result: dict | None = None    error: str | None = None    heartbeat_at: float | None = None  # last sign of life from the worker    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.

Python
import threadingclass Supervisor:    def __init__(self, strategy="least_loaded", stale_after=90.0):        self.tasks: dict[str, Task] = {}        self.workers: dict[str, "Worker"] = {}        self.load: dict[str, int] = {}         # worker -> active task count        self.strategy = strategy        self.stale_after = stale_after        self.lock = threading.Lock()        self._rr = 0    def register(self, worker):        self.workers[worker.name] = worker        self.load[worker.name] = 0    def add_task(self, task: Task):        self.tasks[task.id] = task    # ---- selection -------------------------------------------------    def _eligible(self, task: Task) -> list[str]:        return [n for n, w in self.workers.items() if task.kind in w.can_do]    def _pick(self, task: Task) -> str | None:        candidates = self._eligible(task)        if not candidates:            return None        if self.strategy == "round_robin":            self._rr += 1            return candidates[self._rr % len(candidates)]        if self.strategy == "least_loaded":            return min(candidates, key=lambda n: self.load[n])        if self.strategy == "capability":            return max(candidates, key=lambda n: self.workers[n].skill[task.kind])        raise ValueError(self.strategy)    def _ready(self, task: Task) -> bool:        return all(self.tasks[d].status is Status.COMPLETED                   for d in task.depends_on)    # ---- the main loop ---------------------------------------------    def dispatch_round(self):        with self.lock:            self._reclaim_stale()            for task in self.tasks.values():                if task.status is not Status.PENDING or not self._ready(task):                    continue                worker = self._pick(task)                if worker is None:                    continue                    # no capable worker right now                task.status = Status.ASSIGNED                task.assigned_to = worker                task.heartbeat_at = time.time()                self.load[worker] += 1                self.workers[worker].submit(task)    def _reclaim_stale(self):        now = time.time()        for task in self.tasks.values():            if task.status in (Status.ASSIGNED, Status.RUNNING) \               and now - (task.heartbeat_at or 0) > self.stale_after:                self.load[task.assigned_to] = max(0, self.load[task.assigned_to] - 1)                task.status = Status.PENDING     # give it to someone else                task.assigned_to = None    def on_complete(self, task_id, result):        with self.lock:            t = self.tasks[task_id]            t.status, t.result = Status.COMPLETED, result            self.load[t.assigned_to] -= 1    def on_fail(self, task_id, error):        with self.lock:            t = self.tasks[task_id]            self.load[t.assigned_to] -= 1            t.attempts += 1            t.error = error            t.status = (Status.PENDING if t.attempts < t.max_attempts                        else Status.DEAD)            t.assigned_to = None    def finished(self) -> bool:        with self.lock:            return all(t.status in (Status.COMPLETED, Status.DEAD)                       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.

Python
class Worker:    def __init__(self, name, can_do: set[str], skill: dict[str, float],                 supervisor):        self.name, self.can_do, self.skill = name, can_do, skill        self.supervisor = supervisor    def submit(self, task: Task):        threading.Thread(target=self._run, args=(task,), daemon=True).start()    def _run(self, task: Task):        task.status = Status.RUNNING        try:            result = self.execute(task)            self.supervisor.on_complete(task.id, result)        except Exception as exc:            self.supervisor.on_fail(task.id, f"{type(exc).__name__}: {exc}")    def execute(self, task: Task) -> dict:        # A real worker calls a model, a tool, an API.        raise NotImplementedError

Wiring it together is then unremarkable, which is the point:

Python
sup = Supervisor(strategy="least_loaded")for i in range(8):    sup.register(TicketWorker(f"w{i}", {"classify_ticket"}, {"classify_ticket": 1.0}, sup))for ticket in todays_tickets:                       # 2,000 of them    sup.add_task(Task(kind="classify_ticket", payload={"ticket": ticket}))while not sup.finished():    sup.dispatch_round()    time.sleep(0.05)

When to use it, and when not to

SituationSupervisor-worker?Why
Many independent items of the same kindYesExactly what it is for; scales linearly until the supervisor saturates
Heterogeneous subtasks needing different toolsYesCapability routing sends each to the right specialist
Work whose size is unknown until runtimeYesThe supervisor can add tasks as they are discovered
A strictly sequential chain (A then B then C)NoNothing overlaps; you pay coordination for zero parallelism
Two agents that negotiate as equalsNoThere is no authority to centralise; use a peer protocol
Tasks under ~100 ms eachNoDispatch overhead exceeds the work
You cannot tolerate any single point of failureNot aloneThe 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.

Python
DECOMPOSE_PROMPT = """Break the goal into 3-8 independent subtasks.Return ONLY JSON: a list of objects with keys  "kind"  : one of {allowed_kinds}  "payload": object with the arguments that kind needs  "depends_on": list of zero-based indices of subtasks that must finish firstRules: no subtask may depend on a later index; prefer independence.GOAL: {goal}"""def decompose(goal: str, allowed_kinds: set[str], model) -> list[Task]:    raw = model.generate(DECOMPOSE_PROMPT.format(        goal=goal, allowed_kinds=sorted(allowed_kinds)))    specs = json.loads(extract_json_block(raw))    tasks = [Task(kind=s["kind"], payload=s["payload"]) for s in specs]    for i, s in enumerate(specs):        if s["kind"] not in allowed_kinds:              # hallucinated worker            raise InvalidPlan(f"unknown kind {s['kind']!r}")        for d in s.get("depends_on", []):            if d >= i:                                  # forward or self edge                raise InvalidPlan(f"subtask {i} depends on {d}")            tasks[i].depends_on.append(tasks[d].id)    return tasks

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

StrategyNeeds to knowMakespan on the exampleFails when
Round-robinNothing30 sTask durations vary
Least-loadedActive count per worker16 sCount is a poor proxy for work (one huge task vs three tiny)
CapabilityA skill score per kindDepends on skill spreadOne specialist becomes a queue
Capability then least-loadedBothBest in practiceNeeds 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.

FailureHow the supervisor sees itCorrect response
Task raised an exceptionon_fail called with an errorRetry with backoff up to max_attempts, then mark DEAD
Worker process diedHeartbeat goes stale; no callback ever arrivesReclaim after stale_after, reassign to another worker
Worker alive but stuckHeartbeats continue; task never completesA 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:

Python
import randomdef backoff_seconds(attempt: int, base: float = 2.0, cap: float = 60.0):    # attempt 1 -> ~2 s, 2 -> ~4 s, 3 -> ~8 s, with +/-25% jitter    delay = min(cap, base * (2 ** (attempt - 1)))    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.