Перейти к содержимому
PD
AI-автоматизация 8 мин чтения

Хватит запускать агентов внутри HTTP-запроса: архитектура очереди, которую я использую

AI-агент - это фоновая задача, а не веб-запрос. Разбираю архитектуру с очередью в Postgres, чекпоинтами по шагам, правилами retry, лимитами по стоимости и отменой прогона.

PD

Pavel Duglas

AI Automation & MVP Architect

Почти в каждом AI-продукте, который меня зовут спасать, баг лежит не в промпте и не в выборе модели. Он лежит на самом дне стека: агент выполняется внутри HTTP-запроса. Кто-то повесил POST /api/agent прямо на цикл tool-вызовов и надеется, что браузер подождёт. Браузер ждёт секунд тридцать, потом какой-нибудь балансировщик рвёт соединение, клиент делает retry, и вот две копии агента одновременно правят одну и ту же запись в CRM.

Агент - это не запрос. Это batch-задача, которая умеет разговаривать. Как только вы это принимаете, почти все больные места (таймауты, повторы, неконтролируемые расходы на токены, отмена, прогресс в UI) превращаются в обычные инфраструктурные задачи, которые вы уже сто раз решали.

Почему request-путь всегда проигрывает

HTTP-запрос оптимизирован под ровно противоположный профиль нагрузки. Запросы короткие, stateless, дешёвые для повтора и безопасные для потери. Прогон агента - длинный, stateful, дорогой для повтора и опасный, если оборвать его на середине.

Что конкретно убивает агентов в request-пути:

  • Таймауты прокси, которые вы не контролируете. Cloudflare, nginx, API Gateway, роутер хостинга. У каждого свой idle-лимит. Ваш четырёхминутный research-агент обязательно встретится с самым строгим.
  • Retry, о которых вы не просили. Обёртки над fetch, мобильная сеть, пользователь, который два раза нажал кнопку. Если у агента есть side effects, каждый retry - это дубль побочного эффекта.
  • Deploy. Любой деплой во время шестиминутного прогона убивает прогон, не оставив следа, на каком шаге он остановился.
  • Нет backpressure. Десять пользователей нажали одновременно - у вас десять параллельных агентов, каждый держит воркер веб-сервера и жжёт токены. Сервер умирает не от трафика, а от собственной логики.
  • Нет наблюдаемости. Когда всё падает, у вас одна строка в логе и 502. На вопрос “на каком шаге он был?” ответить нечем, потому что шагов не было, был стек-фрейм.

В serverless всё ещё жёстче: функция упирается в лимит по wall clock прямо посреди tool-вызова, платформа возвращает то ли успех, то ли абстрактный таймаут, а провайдер модели всё равно выставляет счёт за уже потраченные токены.

Двухслойная модель

Разрежьте систему пополам и никогда не размывайте границу.

Слой запроса. Валидирует вход, создаёт строку прогона, кладёт её в очередь, возвращает 202 Accepted и run_id. Цель - меньше 100 мс. Он никогда не вызывает модель. Вообще никогда.

Слой работы. Отдельный процесс-воркер (другой юнит деплоя, другие правила масштабирования), который забирает прогоны из очереди, выполняет шаги, сохраняет состояние после каждого и пишет события. Он может работать двадцать минут, и никому от этого не плохо.

Клиент дальше поллит GET /api/runs/:id. Вот и вся архитектура. Всё остальное ниже - детали о том, как заставить рабочий слой выжить при встрече с реальностью.

Минимально жизнеспособная очередь - это Postgres

Kafka вам не нужна. Пока у вас меньше нескольких тысяч прогонов в сутки, Postgres - правильный ответ, потому что состояние прогона и очередь живут в одной транзакции. Это убирает целый класс багов вида “задача в очереди есть, а строки в базе нет”.

Две таблицы: прогоны и шаги.

create table agent_runs (
  id            uuid primary key default gen_random_uuid(),
  tenant_id     uuid not null,
  kind          text not null,           -- 'research', 'enrich_lead', ...
  input         jsonb not null,
  status        text not null default 'queued',
                -- queued|running|waiting_human|done|failed|cancelled
  attempt       int  not null default 0,
  max_attempts  int  not null default 3,
  cursor        jsonb not null default '{}',  -- откуда продолжать
  result        jsonb,
  error         text,
  idempotency_key text unique,
  cost_cents    int not null default 0,
  locked_by     text,
  locked_until  timestamptz,
  run_after     timestamptz not null default now(),
  created_at    timestamptz not null default now(),
  updated_at    timestamptz not null default now()
);

create table agent_steps (
  id         bigserial primary key,
  run_id     uuid not null references agent_runs(id) on delete cascade,
  seq        int not null,
  name       text not null,
  input      jsonb,
  output     jsonb,
  tokens_in  int,
  tokens_out int,
  status     text not null,
  started_at timestamptz not null default now(),
  ended_at   timestamptz,
  unique (run_id, seq)
);

Забор работы с лизом, чтобы мёртвый воркер не блокировал прогон навсегда:

update agent_runs
set status = 'running',
    locked_by = $1,
    locked_until = now() + interval '5 minutes',
    attempt = attempt + 1,
    updated_at = now()
where id = (
  select id from agent_runs
  where status in ('queued')
    and run_after <= now()
  order by run_after
  for update skip locked
  limit 1
)
returning *;

for update skip locked - та самая фича Postgres, благодаря которой весь паттерн работает с несколькими воркерами без гонок. Пока воркер занят, он раз в 30 секунд сдвигает locked_until вперёд (heartbeat). Отдельный reaper переводит всё, у чего locked_until < now(), назад в queued, чтобы упавшие прогоны кто-нибудь подобрал.

Чекпоинт после каждого шага, а не в конце

Главное отличие от “агента в запросе” - у прогона появляется точка возобновления. Если ваш агент - это один большой while (true) в памяти, то любой рестарт теряет всё, включая самые дорогие части.

Поэтому цикл пишется как явные шаги с сохраняемым результатом:

async function runStep(run, step) {
  const existing = await findStep(run.id, step.seq);
  if (existing?.status === 'done') return existing.output; // реплей, бесплатно

  const out = await execute(step);                          // вызов модели или тула
  await saveStep(run.id, step.seq, out);
  await updateCursor(run.id, nextCursor(step, out));
  return out;
}

При повторе воркер проигрывает завершённые шаги из базы, а не вызывает модель заново. Прогон, который умер на седьмом шаге из девяти, стоит вам два шага, а не девять. На реальном клиентском проекте это срезало расходы на токены при повторах примерно на 70%, потому что большинство падений происходило поздно, в фазе записи, а не в фазе рассуждений.

Бонусом вы бесплатно получаете аудит. Когда клиент спрашивает, почему агент написал письмо не тому контакту, вы открываете agent_steps и читаете фактические входы tool-вызова. Без гадания по логам.

Правила retry, которые не делают хуже

Слепые повторы на агенте - это способ отправить один счёт четыре раза. Сначала классифицируем ошибку, потом решаем.

  • Повторяемые, без изменения состояния: 429, 5xx от провайдера модели, обрывы соединения, таймауты на read-only тулах. Экспоненциальный backoff с jitter, ставим run_after, возвращаем в queued.
  • Неповторяемые: ошибки валидации схемы, 400, отказ модели, превышен бюджет. Падаем сразу и показываем человеку понятную причину.
  • Неоднозначные ошибки с побочным эффектом: таймаут на POST в платёжку или CRM. Никогда не повторять слепо. Либо tool-вызов несёт idempotency key, который вендор реально уважает, либо шаг уходит в waiting_human. “Неоднозначная запись” - это ровно тот случай, где человек дешевле любого умного кода.

Лимит попыток - 3. Прогон, который упал три раза, это баг, а не неудача. А отравленный прогон в бесконечном цикле повторов с радостью съест месячный бюджет на модель за выходные.

Параллелизм, справедливость и жёсткий потолок по деньгам

Очередь - это ещё и место, где вы наконец можете управлять экономикой. В request-пути это невозможно физически.

  • Глобальный лимит воркеров. Начните с 5, не с 50. Rate limit провайдера накажет вас задолго до того, как кончится CPU.
  • Лимит на тенанта. Один корпоративный клиент, залив CSV на 3000 строк, не должен голодить всех остальных. Добавьте в claim-запрос условие вида and (select count(*) from agent_runs r2 where r2.tenant_id = agent_runs.tenant_id and r2.status = 'running') < 3 или держите простую табличку счётчиков.
  • Потолок стоимости на прогон. Инкрементируйте cost_cents после каждого шага. Как только перешли лимит - стоп и понятная ошибка. Меня это дважды спасло от цикла, в котором агент искал один и тот же запрос по кругу.
  • Разные очереди под разную форму задач. Дешёвая пятисекундная классификация и восьмиминутный research не должны делить один пул воркеров, иначе длинные заблокируют короткие и график latency будет выглядеть как кардиограмма.

Отмена, пауза и human-in-the-loop

Поскольку прогон - это строка в таблице, отмена выглядит так: update agent_runs set status = 'cancelled'. Воркер проверяет статус между шагами и выходит. Это самая дешёвая корректная реализация кнопки “стоп”, и пользователям она нужна намного сильнее, чем вам кажется.

Тот же механизм бесплатно даёт согласования. Шаг, которому нужна подпись человека, пишет предлагаемое действие в свою строку и переводит прогон в waiting_human. Воркер отпускает лиз и уходит делать другую работу. Когда человек нажал “подтвердить” в UI, статус возвращается в queued, и прогон продолжается с курсора - возможно, через несколько часов. Ни одного открытого соединения, ни байта состояния в памяти. Попробуйте собрать такое внутри HTTP-хендлера.

Прогресс без собора из вебсокетов

Фаундеры просят streaming. Пользователю на самом деле нужно понимать, что процесс живой и примерно где он сейчас. Поллинг GET /api/runs/:id раз в 2 секунды, отдающий статус, имя текущего шага, количество завершённых шагов и частичный результат. Это двадцать строк кода, и оно работает на дырявом мобильном интернете, через корпоративные прокси и через ваш следующий деплой.

SSE добавите позже, если продукту действительно нужен стриминг по токенам. И даже тогда стримить надо из воркера в канал, а не из хендлера, который владеет циклом агента.

Нужен ли вам workflow-движок?

Когда-нибудь, возможно. Фреймворки durable execution решают эту задачу полнее, и если он у вас уже развёрнут, используйте его. Но паттерн на Postgres я вывожу в прод чаще всего остального: он делается за вечер, не требует новой инфраструктуры, а состояние лежит в таблицах, которые ваша команда умеет читать из psql в два часа ночи. За тяжёлым инструментом идите, когда появится fan-out на сотни параллельных ветвей или многодневные процессы с настоящей компенсирующей логикой.

Чеклист миграции

Если агент прямо сейчас живёт в хендлере, я делаю так, по порядку:

  1. Добавить две таблицы. Хендлер создаёт прогон и возвращает 202 с id.
  2. Вынести цикл агента в worker.ts, который запускается отдельным процессом.
  3. Добавить claim-запрос со skip locked, лиз и reaper.
  4. Разбить цикл на именованные шаги и сохранять выход каждого.
  5. Добавить проверку реплея, чтобы повторы пропускали завершённые шаги.
  6. Классифицировать ошибки: повторяемые, фатальные, неоднозначные.
  7. Добавить лимит на тенанта и потолок стоимости на прогон.
  8. Добавить отмену и waiting_human.
  9. Перевести UI на поллинг.

Пункты с 1 по 3 убирают весь класс багов с таймаутами. На пятом видна экономия денег. На восьмом клиенты начинают доверять системе настолько, что дают ей права на запись.

Слой запросов должен быть скучным, быстрым и глупым. Весь интеллект живёт в воркере, которому разрешено не спешить, запоминать своё место и быть прерванным.

Вопросы и ответы

Можно ли обойтись без отдельного процесса-воркера и просто запускать агента в background-таске того же приложения?

Технически да, но вы теряете половину смысла. Фоновая таска внутри веб-процесса умирает при деплое и рестарте, её нельзя масштабировать отдельно от HTTP-трафика, и она конкурирует за память с обработкой запросов. Минимум, что стоит сделать сразу: тот же кодовый образ, но запуск отдельной командой (`node worker.js`) и отдельный сервис в деплое. Тогда вы можете держать 2 веб-инстанса и 1 воркер, перезапускать их независимо, а при росте нагрузки просто поднять число воркеров, не трогая API.

Как правильно организовать idempotency, чтобы retry не создавал дубли в CRM или платёжке?

На двух уровнях. Первый - уровень прогона: `idempotency_key` с unique-индексом в `agent_runs`, чтобы двойной клик пользователя не создал два прогона. Ключ считайте из полезной нагрузки, например хеш от tenant_id, типа задачи и id сущности. Второй - уровень шага: для каждого шага с записью генерируйте детерминированный ключ вида `run_id:seq` и передавайте его вендору в заголовке Idempotency-Key, если он это поддерживает. Если вендор такого не умеет, помечайте шаг как ambiguous при таймауте и уводите прогон в `waiting_human` вместо повтора. Дешевле показать оператору одну кнопку, чем разбирать четыре отправленных счёта.

Когда Postgres как очередь перестаёт справляться и пора переходить на что-то другое?

По моему опыту, простая схема с `for update skip locked` спокойно живёт до нескольких десятков тысяч прогонов в сутки при десятках воркеров, если вы держите таблицу компактной. Первые признаки боли: рост bloat из-за частых UPDATE на heartbeat, длинные autovacuum-циклы, и claim-запрос, который начинает попадать в топ pg_stat_statements. Лечится в этом порядке: выносите heartbeat в отдельную маленькую таблицу лизов, архивируйте завершённые прогоны в холодную таблицу по расписанию, добавляйте партиционирование по дате. Только после этого имеет смысл смотреть на выделенный брокер или durable execution движок - и то обычно не из-за нагрузки, а из-за потребности в fan-out и многодневных процессах.

Похожие статьи