Хватит запускать агентов внутри HTTP-запроса: архитектура очереди, которую я использую
AI-агент - это фоновая задача, а не веб-запрос. Разбираю архитектуру с очередью в Postgres, чекпоинтами по шагам, правилами retry, лимитами по стоимости и отменой прогона.
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 на сотни параллельных ветвей или многодневные процессы с настоящей компенсирующей логикой.
Чеклист миграции
Если агент прямо сейчас живёт в хендлере, я делаю так, по порядку:
- Добавить две таблицы. Хендлер создаёт прогон и возвращает
202с id. - Вынести цикл агента в
worker.ts, который запускается отдельным процессом. - Добавить claim-запрос со
skip locked, лиз и reaper. - Разбить цикл на именованные шаги и сохранять выход каждого.
- Добавить проверку реплея, чтобы повторы пропускали завершённые шаги.
- Классифицировать ошибки: повторяемые, фатальные, неоднозначные.
- Добавить лимит на тенанта и потолок стоимости на прогон.
- Добавить отмену и
waiting_human. - Перевести 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 и многодневных процессах.
Похожие статьи
Сделаю под ключ
Соберу ИИ-агента под реальную задачу
С инструментами, памятью и логами, чтобы он работал в проде, а не только в демо.
от 1 500 $ · 1-2 недели