Supervisor против peer-схемы: проверяем две архитектуры мультиагентной системы

Введение: в чём проблема
Когда задачу делят между несколькими агентами, проблемы обычно начинаются не с качества модели, а с координации. Деплой запускается раньше миграции базы, два агента одновременно берут одну и ту же операцию, а после сбоя исполнителя никто не берёт задачу снова, потому что каждый думает, что ею займётся другой.
В статье мы соберём небольшой прототип на Python и сравним две архитектуры. Первая — Supervisor: координатор знает граф зависимостей, раздаёт задачи специализированным workers и сам решает, повторять задачу или останавливаться. Вторая — peer-схема: равноправные агенты сами выбирают работу с общей доски задач.
Что важно сразу оговорить:
- Workers в прототипе детерминированные заглушки, а не вызовы LLM. Так тесты воспроизводимы, и видно, что результат зависит от архитектуры, а не от случайности модели.
- Сбои задаются явно через «план отказов» (сколько раз worker упадёт на задаче), без генератора случайных чисел.
- Все числа ниже — вывод этого кода на этом сценарии. Это не бенчмарк и не замер на продакшене.
Что именно сравниваем
Сценарий — упрощённый релиз из трёх зависимых операций:
migrate_db— миграция базы (навыкdb);deploy_app— деплой, зависит от миграции (навыкdeploy);smoke_test— проверка, зависит от деплоя (навыкqa).
Исполнителей четыре: db-1, deploy-1, deploy-2 и qa-1. Два deploy-агента нужны специально: только так можно проверить конфликт за одну задачу. Считаем три метрики:
- нарушения порядка — задача выполнилась раньше своей зависимости или вовсе без неё;
- дубли — одна задача успешно выполнена больше одного раза;
- невыполненные задачи — то, что так и не завершилось.
Peer-схему проверяем в двух вариантах: наивном (каждый берёт, что видит) и с протоколом (проверка зависимостей, захват задачи, повтор после сбоя). Так сравнение честнее: peer-схема может работать корректно, но для этого все агенты должны соблюдать общий протокол.
Шаг 1. Подготовка окружения
Нужен Python 3.10 или новее. Внешних зависимостей нет, всё на стандартной библиотеке.
mkdir mas-lab
cd mas-lab
python3 --version
python3 -m venv .venv
. .venv/bin/activate
Виртуальное окружение здесь необязательно, но отделяет эксперимент от системного Python. Команды ничего не устанавливают и не обращаются к сети.
Шаг 2. Задачи, журнал и workers
Создайте файл mas.py. Начнём с общих структур: задачи, журнала событий (ledger) и исполнителя.
from __future__ import annotations
from dataclasses import dataclass, field
class TransientError(Exception):
"""Временный сбой исполнителя: задачу можно повторить."""
@dataclass(frozen=True)
class Task:
name: str
skill: str
deps: tuple[str, ...] = ()
@dataclass
class Ledger:
events: list[str] = field(default_factory=list)
done: set[str] = field(default_factory=set)
failed: set[str] = field(default_factory=set)
class Worker:
def __init__(self, name: str, skill: str, fail_plan: dict[str, int] | None = None):
self.name = name
self.skill = skill
# Сколько раз подряд worker упадёт на конкретной задаче.
self.fail_plan = dict(fail_plan or {})
def run(self, task: Task, ledger: Ledger) -> None:
if self.fail_plan.get(task.name, 0) > 0:
self.fail_plan[task.name] -= 1
ledger.events.append(f"fail:{task.name}:{self.name}")
raise TransientError(task.name)
ledger.events.append(f"run:{task.name}:{self.name}")
В реальной системе внутри Worker.run был бы вызов модели или инструмента. Контракт остаётся тем же: либо событие run, либо исключение. Журнал общий и только дописывается, поэтому по нему потом считаются метрики.
Шаг 3. Supervisor: порядок, маршрутизация, повторы
Supervisor делает три вещи: строит топологический порядок задач, выбирает worker по навыку и повторяет задачу при временном сбое, при необходимости на другом worker того же профиля. Если попытки кончились, он останавливается и не запускает задачи, которые зависят от упавшей. Добавьте в mas.py:
def topo_sort(tasks: list[Task]) -> list[Task]:
by_name = {t.name: t for t in tasks}
for t in tasks:
for dep in t.deps:
if dep not in by_name:
raise ValueError(f"{t.name}: неизвестная зависимость {dep}")
order: list[Task] = []
done: set[str] = set()
visiting: set[str] = set()
def visit(t: Task) -> None:
if t.name in done:
return
if t.name in visiting:
raise ValueError(f"цикл зависимостей на {t.name}")
visiting.add(t.name)
for dep in t.deps:
visit(by_name[dep])
visiting.discard(t.name)
done.add(t.name)
order.append(t)
for t in tasks:
visit(t)
return order
class Supervisor:
def __init__(self, workers: list[Worker], max_retries: int = 2):
self.by_skill: dict[str, list[Worker]] = {}
for w in workers:
self.by_skill.setdefault(w.skill, []).append(w)
self.max_retries = max_retries
def execute(self, tasks: list[Task]) -> Ledger:
ledger = Ledger()
for task in topo_sort(tasks):
if self._run_with_retries(task, ledger):
ledger.done.add(task.name)
else:
ledger.failed.add(task.name)
ledger.events.append(f"abort:{task.name}")
break # зависимые задачи не запускаем
return ledger
def _run_with_retries(self, task: Task, ledger: Ledger) -> bool:
candidates = self.by_skill.get(task.skill)
if not candidates:
raise LookupError(f"нет worker с навыком {task.skill!r}")
for attempt in range(self.max_retries + 1):
# Каждая следующая попытка уходит другому worker того же навыка.
worker = candidates[attempt % len(candidates)]
try:
worker.run(task, ledger)
except TransientError:
continue
return True
return False
Обратите внимание на break после исчерпания попыток. Это сознательное решение: лучше остановить релиз, чем выкатить приложение на базу без миграции.
Шаг 4. Peer-схема в двух вариантах
Peer-система работает раундами. В каждом раунде агенты по очереди смотрят на доску и выбирают задачу своего навыка. Затем все выбранные задачи выполняются: так моделируется параллельность, когда решение принято до того, как стал известен результат соседа. Три флага включают элементы протокола:
class PeerSystem:
def __init__(self, workers: list[Worker], *, check_deps: bool = False,
use_claims: bool = False, retry_failed: bool = False,
max_rounds: int = 10):
self.workers = workers
self.check_deps = check_deps
self.use_claims = use_claims
self.retry_failed = retry_failed
self.max_rounds = max_rounds
def execute(self, tasks: list[Task]) -> Ledger:
ledger = Ledger()
for _ in range(self.max_rounds):
picks = self._pick(tasks, ledger)
if not picks:
break
for worker, task in picks:
try:
worker.run(task, ledger)
except TransientError:
ledger.failed.add(task.name)
else:
ledger.done.add(task.name)
ledger.failed.discard(task.name)
return ledger
def _pick(self, tasks: list[Task], ledger: Ledger) -> list[tuple[Worker, Task]]:
picks: list[tuple[Worker, Task]] = []
claimed: set[str] = set()
for worker in self.workers:
for task in tasks:
if task.skill != worker.skill or task.name in ledger.done:
continue
if task.name in ledger.failed and not self.retry_failed:
continue # «пусть кто-нибудь другой разберётся»
if self.check_deps and not set(task.deps) <= ledger.done:
continue
if self.use_claims and task.name in claimed:
continue
claimed.add(task.name)
picks.append((worker, task))
break
return picks
Без флагов получаем все три симптома из постановки задачи: агенты не смотрят на зависимости, не договариваются о владении задачей, а упавшую задачу никто не подбирает.
Шаг 5. Метрики по журналу
Метрики считаем только по журналу, без доступа к внутреннему состоянию систем. Значит, тот же анализ подойдёт и для логов настоящих агентов, если привести их к такому формату. Добавьте в конец mas.py:
def analyze(tasks: list[Task], ledger: Ledger) -> dict:
runs = [e.split(":")[1] for e in ledger.events if e.startswith("run:")]
first: dict[str, int] = {}
for i, name in enumerate(runs):
first.setdefault(name, i)
violations = 0
for task in tasks:
if task.name not in first:
continue
for dep in task.deps:
if dep not in first or first[dep] > first[task.name]:
violations += 1
return {
"order_violations": violations,
"duplicates": len(runs) - len(first),
"missing": [t.name for t in tasks if t.name not in first],
}
Шаг 6. Сценарий и тесты
Вынесем сценарий в отдельный файл scenario.py. Задачи в списке намеренно перепутаны: порядок должна восстанавливать архитектура, а не автор списка. Workers создаются фабрикой, потому что план отказов меняется во время прогона, а каждому тесту нужен свежий набор.
from mas import Task, Worker
TASKS = [
Task("deploy_app", "deploy", ("migrate_db",)),
Task("smoke_test", "qa", ("deploy_app",)),
Task("migrate_db", "db"),
]
def make_workers(db_failures: int = 1, deploy1_failures: int = 0) -> list[Worker]:
return [
Worker("deploy-1", "deploy", {"deploy_app": deploy1_failures}),
Worker("deploy-2", "deploy"),
Worker("qa-1", "qa"),
Worker("db-1", "db", {"migrate_db": db_failures}),
]
Тесты в test_mas.py используют только unittest:
import unittest
from mas import PeerSystem, Supervisor, Task, analyze, topo_sort
from scenario import TASKS, make_workers
class PeerTests(unittest.TestCase):
def test_naive_peer_breaks_order_duplicates_and_drops_failed(self):
report = analyze(TASKS, PeerSystem(make_workers()).execute(TASKS))
self.assertEqual(report["order_violations"], 1)
self.assertEqual(report["duplicates"], 1)
self.assertEqual(report["missing"], ["migrate_db"])
def test_peer_with_protocol_is_clean(self):
system = PeerSystem(make_workers(), check_deps=True,
use_claims=True, retry_failed=True)
report = analyze(TASKS, system.execute(TASKS))
self.assertEqual(report, {"order_violations": 0, "duplicates": 0, "missing": []})
class SupervisorTests(unittest.TestCase):
def test_recovers_after_transient_failure(self):
ledger = Supervisor(make_workers()).execute(TASKS)
self.assertIn("fail:migrate_db:db-1", ledger.events)
report = analyze(TASKS, ledger)
self.assertEqual(report, {"order_violations": 0, "duplicates": 0, "missing": []})
def test_fails_over_to_another_worker(self):
ledger = Supervisor(make_workers(db_failures=0, deploy1_failures=1)).execute(TASKS)
self.assertIn("fail:deploy_app:deploy-1", ledger.events)
self.assertIn("run:deploy_app:deploy-2", ledger.events)
def test_stops_dependents_when_retries_exhausted(self):
ledger = Supervisor(make_workers(db_failures=5), max_retries=2).execute(TASKS)
self.assertEqual(ledger.failed, {"migrate_db"})
self.assertIn("abort:migrate_db", ledger.events)
report = analyze(TASKS, ledger)
self.assertEqual(report["order_violations"], 0)
self.assertEqual(report["duplicates"], 0)
def test_cycle_is_rejected(self):
cyclic = [Task("a", "x", ("b",)), Task("b", "x", ("a",))]
with self.assertRaises(ValueError):
topo_sort(cyclic)
if __name__ == "__main__":
unittest.main()
python3 -m unittest -v
Должны пройти все шесть тестов. Если хотя бы один упал, сначала сверьте код с листингами: результат детерминирован и от запуска к запуску меняться не должен.
Шаг 7. Сравнение архитектур одной командой
Файл compare.py прогоняет один сценарий через три системы:
from mas import PeerSystem, Supervisor, analyze
from scenario import TASKS, make_workers
SYSTEMS = {
"peer (наивная)": lambda: PeerSystem(make_workers()),
"peer (с протоколом)": lambda: PeerSystem(
make_workers(), check_deps=True, use_claims=True, retry_failed=True),
"supervisor": lambda: Supervisor(make_workers()),
}
for name, build in SYSTEMS.items():
report = analyze(TASKS, build().execute(TASKS))
print(f"{name:22} нарушений порядка={report['order_violations']} "
f"дублей={report['duplicates']} не выполнено={report['missing']}")
python3 compare.py
Если трассировать код вручную, вывод на этом сценарии должен быть таким:
peer (наивная) нарушений порядка=1 дублей=1 не выполнено=['migrate_db']
peer (с протоколом) нарушений порядка=0 дублей=0 не выполнено=[]
supervisor нарушений порядка=0 дублей=0 не выполнено=[]
Как проверить и правильно прочитать результат
Чтобы понять, откуда берутся цифры, распечатайте журнал наивной peer-системы:
python3 -c "from mas import PeerSystem; from scenario import TASKS, make_workers; print(PeerSystem(make_workers()).execute(TASKS).events)"
В первом же раунде оба deploy-агента берут deploy_app (дубль) и выполняют его до миграции (нарушение порядка). Миграция падает, и больше её никто не трогает (задача потеряна). Supervisor на тот же сбой отвечает повторной попыткой, а при исчерпании лимита останавливает зависимые задачи.
Главный вывод: чинит ситуацию не наличие Supervisor как такового, а явный протокол. Это проверка зависимостей перед запуском, единственный владелец задачи и правило повторов. В Supervisor протокол сосредоточен в одном месте и проверяется одним набором тестов. В peer-схеме его должен соблюдать каждый агент. Достаточно одного агента, обновлённого без флага check_deps, чтобы вернулись нарушения порядка. Это легко проверить: передайте в PeerSystem смешанный набор правил или отключите один флаг и перезапустите compare.py.
Подсказки по выбору архитектуры:
- Supervisor подходит, когда операции зависят друг от друга, побочные эффекты необратимы (миграции, платежи, деплой) и нужен один понятный журнал решений.
- Peer-схема подходит для независимых задач, где важнее пропускная способность и устойчивость к отказу координатора, а общий протокол можно обеспечить на уровне хранилища, например атомарным захватом задачи.
Типовые ошибки
- Доверить порядок операций LLM-координатору. Маршрутизацию по смыслу задачи можно поручить модели, но топологический порядок и лимит повторов должен соблюдать код. Иначе гарантия превращается в вероятность.
- Повторять неидемпотентные операции. Повтор безопасен, только если операция не повредит при двойном выполнении или если исполнитель сообщает, было ли действие применено. Разделяйте временные ошибки (
TransientError) и остальные: на постоянные повтор не нужен. - Захват задачи без атомарности. В прототипе
claimed— локальное множество одного процесса. С настоящими процессами захват должен быть атомарной операцией в общем хранилище, иначе дубли вернутся. - Общий worker-объект в нескольких тестах. План отказов меняется при прогоне, поэтому workers создаются фабрикой
make_workers()заново для каждого запуска. - Метрики по самоотчёту агентов. Считайте метрики по журналу событий, а не по тому, что агент написал о своей работе.
Ограничения прототипа
- Параллельность смоделирована раундами в одном потоке. Гонки настоящих процессов, сетевые таймауты и частичные отказы здесь не воспроизводятся.
- Между повторами нет задержки (backoff) и нет таймаутов. В рабочей системе нужны и то и другое, а ещё ограничение общего бюджета попыток.
- Supervisor — единая точка отказа. Его состояние (журнал и пройденные шаги) надо сохранять, чтобы после перезапуска продолжить, а не начать сначала.
- Supervisor выполняет задачи строго последовательно. Независимые ветки графа можно запускать параллельно, но это отдельное расширение со своими тестами.
- Сценарий из трёх задач и четырёх агентов показывает механизмы отказов, а не их частоту в реальных системах. Перенос выводов на ваш процесс требует собственных прогонов на ваших логах.
Что дальше
Следующий практичный шаг — заменить заглушку в Worker.run на реальный вызов модели или инструмента и оставить тесты без изменений. Если они по-прежнему проходят, протокол координации не зависит от исполнителя. Другие практические разборы собраны в разделе руководств, а определения терминов (worker, идемпотентность, топологическая сортировка и другие) — в глоссарии.