← Cache & Workers

Жизненный цикл задачи

cache-workers · workzone

Каноническая конвенция: что происходит с фоновой задачей от постановки до терминала — общая для всех лейнов и потребителей. Топология очередей и воркеров — на очередях и воркерах; откуда у замка корректность — на водоразделе Redis ⟂ Postgres.

Путь задачи

Постановка идемпотентна по job_id: повтор с тем же id не плодит вторую задачу. Ретрай — та же задача, тот же id, новая попытка; не вторая строка в журнале.

queued в очереди, ждёт воркера
enqueue
running воркер взял · штампует heartbeat
pick
succeeded итог в строке потребителя
failed ошибка · в DLQ-политику
terminal
reaping пост-терминал · замок отпущен, окно идемпотентности

enqueue → queued → running → succeeded / failed → reaping. Терминал и снятие замка — отдельный шаг (см. терминал и DLQ, единственность); упавший без очистки воркер задачу не освобождает — её снимает сторож (см. heartbeat).

Heartbeat и очистка

Heartbeat — пульс живой задачи: долгий прогон каждые ~30с штампует heartbeat_at в свою строку в Postgres. По этому пульсу сторож отличает работающую задачу от мёртвой.

пульс есть
задача жива
heartbeat_at свежий — воркер работает; замок единственности держится, следующая такая задача ждёт.
пульс протух
воркер мёртв
Нет пульса дольше порога → сторож очищает прогон в failed / stale, замок единственности отпускается, слот свободен.

Без пульса упавшая задача вечно числилась бы running и блокировала следующую — очистка закрывает этот тупик.

Сторож живёт на singleton-планировщике. Скан протухших heartbeat — периодический проход, а не работа каждого воркера: иначе N реплик гоняли бы очистку наперегонки. Дом — тот же выделенный singleton, что тикает cron. Цена связки: пока singleton лежит, реапинг тоже стоит — замок упавшей задачи не снимается, и этот тип задачи ждёт возврата реплики. Простой ограничен: на возврате сторож сразу разгребает протухшее, а каденсы грубые — приемлемо для v1.
Плановая остановка — мягкий слив, не реапинг. Сторож ловит крах; при штатном рестарте (деплой, SIGTERM) воркер сливается аккуратно: перестаёт брать новые задачи, доводит начатые, затем выходит. Не успел в срок — задача возвращается в очередь и подхватывается заново: доставка at-least-once плюс идемпотентность по job_id делают повтор безопасным.

Ретраи

Два уровня по длине паузы — короткую пережидаем в задаче, длинную возвращаем в очередь, чтобы не держать воркер спящим.

короткий · ≤60с
in-memory, очередь не трогаем
Пауза мелкая — задача сама ждёт и повторяет; брокер и слот не меняются. Дешевле, чем гонять задачу через очередь туда-обратно.
долгий · минуты–часы
defer / re-enqueue с задержкой
Большой Retry-After — задача возвращается в очередь с отложенным стартом, слот воркера освобождается. Не спим в задаче.
Дефолты надёжности — платформенные, override конфигом. Классификация ошибки решает, повторять ли вообще: transient (таймаут, 429, временный сбой) → ретрай; permanent (валидация, 4xx-логика) → сразу терминал, без попыток впустую. Бэкофф — экспоненциальный с full jitter (случайный разброс гасит thundering herd — синхронный навал ретраев после общего сбоя), потолок паузы ~300с, не более ~5 попыток. Конкретные числа — тюнинг под нагрузку, не догма.
classify transient → ретрай · permanent → терминал сразу backoff exp + full jitter · cap ~300с attempts ≤ 5, затем терминал → DLQ-политика

Единственность

Взаимное исключение держит Postgres, не брокер: частичный уникальный индекс рядом с журналом плюс heartbeat. Замок корректностно-критичный — потому он там, где система записи.

Замок = партиал-UNIQUE на активных состояниях. UNIQUE … WHERE state IN ('queued','running') — пока задача активна, вторую с тем же ключом индекс не пустит. Повторный запуск (двойной клик, гонка планировщика, ретрай) упирается в индекс и получает skipped, а не дубль. После терминала строка выходит из-под условия — индекс снова свободен.
succeeded · ~24ч
окно идемпотентности
Успешный job_id держится в окне, чтобы повторный enqueue при at-least-once-доставке не запустил задачу заново.
failed · коротко
пересабмит разрешён
Сбойный job_id отпускается быстро — оператор или ретрай-политика вправе пересабмитить ту же работу.

Источник истины окна — журнал в Postgres; ключ в Redis (см. эфемерные ключи) лишь быстрый путь, дублирующий, не заменяющий журнал.

Отвергнуто: распределённый лок-сервис (Redlock / SETNX). Взаимное исключение уже держит партиал-уникальный индекс плюс heartbeat — Redis-лок не нужен и без fencing-токенов опасен для корректности (зависший держатель не теряет лок). Та же ось, что и на водоразделе Redis ⟂ Postgres.

Терминал и DLQ

Терминал читает потребитель: итог прогона или доставки — статус в его собственной Postgres-строке, он не теряется в брокере. Сбойные задачи у каждого потребителя оседают в его доме, откуда доводочный проход их перезабирает.

таблица DLQ Harvester — сбойные элементы в таблицу dead-letter; reconcile / dlq-retry перезабирает.
статус «не доставлено» Email и Notifications — провал помечается на строке доставки, повтор по ретрай-политике.
строка прогона Knowledge Store — терминал прогона (failed с причиной и статистикой шагов) в curation_runs; следующий проход доводки перебирает граф заново.
строка прогона Agent Engine — терминал прогона (failed с причиной) в строке прогона agent_runs.
Отвергнуто: durable DLQ-очередь в Redis. Долговечный дом сбоя — Postgres у потребителя, не отдельная очередь в эфемерном слое. Redis-DLQ остаётся лишь транспортным бэкстопом на пути к этому дому, не местом, где разбирают разбор.

Платформенный контракт прогона

Унифицируем контракт. Скелет жизненного цикла один на всех потребителей; дом данных и смысл перезабора — доменные. Так монитор и сторож работают по любому прогону, не схлопывая разные домены в один стор.

state queued · running · succeeded · failed (+ stale как под-итог сторожа) attempts счётчик попыток против потолка ретраев heartbeat_at пульс живого прогона; по нему сторож отличает мёртвого last_error причина терминала + статистика шагов next_retry_at отложенный старт ретрая; NULL, если не запланирован
общее · платформа
скелет — один на всех
Поля выше — общий словарь (mixin), имена выровнены где применимо. Не каждый потребитель несёт все: счётчик попыток и отложенный старт — лишь у тех, кто ретраит счётчиком; KS перезабирает новым прогоном, не ретраем. Сторож heartbeat → stale и замок единственности — общие для любой таблицы с контрактом.
домашнее · потребитель
смысл — у каждого свой
Payload строки и смысл «переделать работу» доменные: reconcile / next pass / resend у каждого свой. Дом данных остаётся у потребителя — принцип близости.

Монитор читает каждую таблицу прогонов в её доменном доме и сводит их в админ-сводке (module, job_id, state, attempts, age, error): один вид «что и где упало» без центральной таблицы. Единая проекция UNION ALL поверх таблиц прогонов — удобство v2, когда таблиц прогонов станет много.