Перейти к содержимому

Multi-agent оркестрация встала намертво

Один агент завис — это agent-stuck: watchdog на цикл, и готово. Но когда агентов несколько (LangGraph / CrewAI / AutoGen), появляется класс поломок, которого у одиночки нет: каждый агент по отдельности жив и что-то делает, а система в целом не двигается.

  • researcher ждёт данных от fetcher-а, fetcher завис на падающем tool-е;
  • writer ждёт ревью от critic-а, critic ждёт финальный текст от writer-а — классический circular wait;
  • очередь сообщений пуста, все агенты в состоянии «жду входящее», и никто никому уже не напишет.

Локальные heartbeat-ы каждого агента тут молчат — формально все «живы». Нужен сигнал о глобальном прогрессе.

Заводим один глобальный счётчик завершённых шагов. Любой агент, доведя шаг до конца, дёргает mark_progress(). Отдельный watchdog смотрит: если счётчик не рос дольше STALL_SEC — оркестрация встала, шлём push.

import os, time, threading, requests
class ProgressWatchdog:
def __init__(self, stall_sec=300):
self.stall_sec = stall_sec
self.counter = 0
self.last_bump = time.time()
self._fired = False
self._lock = threading.Lock()
def mark_progress(self, agent, step):
# любой агент вызывает это, завершив реальный шаг
with self._lock:
self.counter += 1
self.last_bump = time.time()
self._fired = False # прогресс пошёл — сбрасываем алёрт
self._last = f"{agent}:{step}"
def watch(self):
while True:
time.sleep(15)
idle = time.time() - self.last_bump
if idle >= self.stall_sec and not self._fired:
self._fired = True
notify("🕸️🛑 Оркестрация встала",
f"Ни один шаг не завершился {int(idle)}с "
f"(всего шагов: {self.counter}, последний: {self._last}). "
f"Похоже на deadlock.",
priority=9)
def notify(title, message, priority):
requests.post(f"{os.environ['NOTIFLY_URL']}/message",
params={"token": os.environ["NOTIFLY_TOKEN"]},
json={"title": title, "message": message, "priority": priority},
timeout=5)

Подключаем к оркестратору и запускаем watcher в фоне:

wd = ProgressWatchdog(stall_sec=300)
threading.Thread(target=wd.watch, daemon=True).start()
# в узле графа / коллбэке агента — по завершении шага:
def on_step_complete(agent, step):
wd.mark_progress(agent, step) # глобальный прогресс сдвинулся
orchestrator.run(...) # закончит сам — daemon-поток умрёт с процессом

Важное отличие от agent-stuck: там heartbeat на одном цикле, тут — один счётчик на всю систему. Отдельный агент может честно молотить (а значит, его личный heartbeat пингует), но если это не двигает общий счётчик — оркестрация всё равно мертва, и watchdog это увидит.

Если оркестратор — короткоживущий процесс (Yandex Cloud function по расписанию, CI-джоба), свой while-поток не нужен: пусть heartbeat следит снаружи. Заведите heartbeat в админке с интервалом = ваш STALL_SEC, и пингуйте его только при реальном глобальном прогрессе:

PROGRESS_PING = os.environ["ORCH_PROGRESS_PING"] # https://.../heartbeat/ping/H...
def on_step_complete(agent, step):
# пингуем НЕ на каждом чихе агента, а на завершённом шаге оркестрации
requests.get(PROGRESS_PING, timeout=3)

Пока шаги завершаются — ping идёт, heartbeat зелёный. Оркестрация встала — ping прекращается, и по intervalSec + graceSec Notifly сам пришлёт алёрт. Плюс — сработает и при полном крахе процесса (OOM, kill), когда никакой внутренний watchdog уже не жив. Включите recovery-сообщение — когда система разблокируется, получите подтверждение.

Dead-man-switch скажет «стоим», но не скажет почему. Часто причина — circular wait: A ждёт B, B ждёт A. Ведём граф «кто кого ждёт» и ищем в нём цикл — это ловит deadlock мгновенно, не дожидаясь таймаута.

class WaitGraph:
def __init__(self):
self.waits = {} # agent -> на кого он сейчас ждёт (или None)
def set_waiting(self, agent, waiting_on):
self.waits[agent] = waiting_on
cycle = self._find_cycle(agent)
if cycle:
notify("🔁🛑 Циклическое ожидание агентов",
f"Deadlock: {''.join(cycle)}{cycle[0]}. "
f"Никто не разблокируется сам.",
priority=10)
def clear(self, agent):
self.waits[agent] = None # агент получил своё, пошёл дальше
def _find_cycle(self, start):
seen, path, cur = set(), [], start
while cur is not None and cur not in seen:
seen.add(cur); path.append(cur)
cur = self.waits.get(cur)
if cur is not None: # вернулись в уже пройденный узел
return path[path.index(cur):]
return None
wg = WaitGraph()
wg.set_waiting("writer", "critic") # writer ждёт critic-а
wg.set_waiting("critic", "writer") # critic ждёт writer-а → цикл, push!

Комбинируйте: детектор цикла ловит явный deadlock моментально, а dead-man-switch остаётся как страховка от тихих зависаний, которые в граф ожиданий не попали (агент завис на tool-е, а не на другом агенте).

  • сколько всего шагов завершено и когда был последний (5 мин назад);
  • имя последнего активного агента и шага;
  • если это цикл — сама цепочка ожидания writer → critic → writer;
  • глубину очереди сообщений (пустая очередь + все «ждут» = верный deadlock).