Зачем агенту очередь, если API отвечает сразу

Синхронный вызов удобен до первого длинного анализа, временного отказа модели или всплеска трафика. Очередь отделяет обещание «задача принята» от фактического выполнения. Веб-слой быстро валидирует запрос, создаёт job_id и возвращает статус, а worker забирает работу тогда, когда есть квота модели, память и доступные tool-сервисы.

Это не просто ускоритель. Очередь становится границей управления: здесь задаются приоритет, дедлайн, tenant, бюджет токенов, допустимые инструменты и политика повторов.

Контракт задания

Минимальный envelope содержит job_id, tenant_id, тип операции, версию workflow, ссылку на входные данные, дедлайн и idempotency key. Большие документы не кладите в сообщение: сохраняйте их в объектном хранилище, а в очередь передавайте неизменяемую ссылку и checksum.

{
  "job_id": "job_01J...",
  "workflow": "research_report:v3",
  "input_ref": "s3://bucket/input/sha256...",
  "deadline_at": "2027-03-29T12:00:00Z",
  "attempt": 1,
  "trace_id": "..."
}

At-least-once означает дубликаты

Надёжные брокеры обычно предпочитают не потерять сообщение, поэтому worker может получить его повторно: подтверждение потерялось, процесс упал после внешнего действия или истёк visibility timeout. Нельзя рассчитывать, что «ровно один раз» обеспечит бизнес-результат. Перед каждым необратимым действием проверяйте idempotency key и храните итог операции в транзакционной таблице.

Visibility timeout и heartbeat

После выдачи сообщение временно скрывается от других workers. Timeout должен покрывать нормальную длительность шага с запасом, но не быть настолько большим, чтобы падение зависло на часы. Для непредсказуемо долгих задач worker отправляет heartbeat и продлевает lease. Если heartbeat прекратился, сообщение возвращается в очередь.

Разделяйте общий дедлайн задания и lease одной попытки. Продление lease не должно бесконечно продлевать пользовательский SLA.

Retries без шторма

Повторяйте только временные ошибки: 429, часть 5xx, сетевые обрывы и истечение lease. Ошибки валидации, отсутствие прав и превышение бюджета обычно неретрайбельны. Используйте exponential backoff с jitter, ограничение попыток и отдельный retry budget. Иначе массовый сбой провайдера создаст retry storm и окончательно забьёт квоту.

  • 1-я попытка — сразу;
  • 2-я — через 5–15 секунд;
  • 3-я — через 30–90 секунд;
  • дальше — медленнее, но только до дедлайна.

DLQ — не кладбище

После исчерпания попыток сообщение уходит в dead-letter queue вместе с кодом ошибки, версией worker, последним checkpoint и ссылкой на trace. Для DLQ нужны владелец, алерт по росту, runbook разбора и безопасный redrive. Перед повторным запуском исправьте причину и отфильтруйте задания, которые уже просрочены или частично выполнили внешнее действие.

Приоритеты и честность

Одна FIFO-очередь заставляет срочные короткие задачи ждать за длинным отчётом. Практичнее несколько классов обслуживания: interactive, batch и maintenance. Внутри tenant применяйте квоты, чтобы один крупный клиент не вытеснил остальных. Приоритет не отменяет дедлайн и cost ceiling.

Метрики и план запуска

Смотрите на age of oldest message, time-in-queue, execution time, success rate по попыткам, DLQ rate, saturation workers, 429 провайдера и стоимость задания. Начните с одного типа workflow и небольшого concurrency, проведите тест падения worker в середине tool-call, затем включите autoscaling по возрасту очереди, а не только по её длине.

Архитектура потока от API до результата

Входной API не должен публиковать сообщение до фиксации задания: сохраните job и outbox event одной транзакцией, затем publisher передаст событие брокеру. Worker сначала захватывает lease, проверяет дедлайн и отмену, затем выполняет шаги и атомарно фиксирует результат. Клиент получает статус через polling, webhook или server-sent events, но ни один канал доставки статуса не считается источником истины.

СлойОтветственностьТипичный сбой
APIвалидация, job_id, authповтор запроса
Brokerбуфер и доставкадубликат
Workerworkflow и checkpointпадение процесса
Result storeстатус и артефактыконфликт записи

Backpressure и autoscaling

Масштабирование только по количеству сообщений обманчиво: десять двухчасовых отчётов тяжелее тысячи коротких классификаций. Оценивайте backlog в секундах работы, возраст p95 и доступную квоту tokens per minute. Новый worker бесполезен, если провайдер уже отвечает 429; в этом случае admission control должен замедлить приём или перевести batch в отложенный класс.

Отмена и дедлайн

Отмена — состояние, а не убийство одного процесса. Worker проверяет cancellation token между шагами, прекращает новые tool calls, отзывает временные полномочия и помечает незавершённые артефакты. Уже начавшийся внешний эффект нужно завершить, сверить или компенсировать. Результат должен различать cancelled, expired, failed и partially_completed.

С чего начать внедрение?
Выберите один реальный workflow, опишите его внешние эффекты и отказ в середине выполнения. Затем добавьте минимальные policy, idempotency и наблюдаемость до роста автономности.
Можно ли решить задачу одним system prompt?
Нет. Prompt задаёт поведение модели, но timeout, права, уникальность операций, изоляция и аудит должны обеспечиваться исполняющей инфраструктурой.
Какие ошибки нельзя повторять автоматически?
Ошибки прав и валидации, исчерпанный бюджет, просроченный дедлайн и необратимые действия с неизвестным исходом требуют остановки, сверки состояния или человека.
Что измерять в production?
Успешность и длительность по шагам, число повторов, возраст работы, внешние эффекты, ошибки policy, стоимость, насыщение лимитов и долю ручных вмешательств.
Как безопасно увеличить автономность?
Расширяйте полномочия по одному классу действий: shadow, затем canary с малым лимитом, обязательный аудит и kill switch, после чего анализируйте реальные инциденты.
← Все статьи блога