Каноническая конвенция: что происходит с фоновой задачей от постановки до терминала — общая для всех лейнов и потребителей. Топология очередей и воркеров — на очередях и воркерах; откуда у замка корректность — на водоразделе Redis ⟂ Postgres.
Постановка идемпотентна по job_id: повтор
с тем же id не плодит вторую задачу. Ретрай — та же задача, тот же id,
новая попытка; не вторая строка в журнале.
enqueue → queued → running → succeeded / failed → reaping. Терминал и снятие замка — отдельный шаг (см. терминал и DLQ, единственность); упавший без очистки воркер задачу не освобождает — её снимает сторож (см. heartbeat).
Heartbeat — пульс живой задачи: долгий прогон каждые
~30с штампует heartbeat_at в свою строку в Postgres. По
этому пульсу сторож отличает работающую задачу от мёртвой.
heartbeat_at свежий — воркер работает; замок
единственности держится, следующая такая задача ждёт.
failed / stale, замок единственности отпускается,
слот свободен.
Без пульса упавшая задача вечно числилась бы running и
блокировала следующую — очистка закрывает этот тупик.
N реплик гоняли бы очистку наперегонки. Дом —
тот же выделенный singleton,
что тикает cron. Цена связки: пока singleton лежит, реапинг тоже стоит —
замок упавшей задачи не снимается, и этот тип задачи ждёт возврата
реплики. Простой ограничен: на возврате сторож сразу разгребает
протухшее, а каденсы грубые — приемлемо для v1.
Два уровня по длине паузы — короткую пережидаем в задаче, длинную возвращаем в очередь, чтобы не держать воркер спящим.
Retry-After — задача возвращается в очередь с
отложенным стартом, слот воркера освобождается. Не спим в задаче.
Взаимное исключение держит Postgres, не брокер: частичный уникальный индекс рядом с журналом плюс heartbeat. Замок корректностно-критичный — потому он там, где система записи.
UNIQUE … WHERE state IN ('queued','running') — пока задача
активна, вторую с тем же ключом индекс не пустит. Повторный запуск
(двойной клик, гонка планировщика, ретрай) упирается в индекс и
получает skipped, а не дубль. После терминала строка
выходит из-под условия — индекс снова свободен.
job_id держится в окне, чтобы повторный
enqueue при at-least-once-доставке не запустил задачу заново.
job_id отпускается быстро — оператор или
ретрай-политика вправе пересабмитить ту же работу.
Источник истины окна — журнал в Postgres; ключ в Redis (см. эфемерные ключи) лишь быстрый путь, дублирующий, не заменяющий журнал.
Терминал читает потребитель: итог прогона или доставки — статус в его собственной Postgres-строке, он не теряется в брокере. Сбойные задачи у каждого потребителя оседают в его доме, откуда доводочный проход их перезабирает.
reconcile
/ dlq-retry перезабирает.
failed с причиной и статистикой
шагов) в curation_runs; следующий проход доводки
перебирает граф заново.
failed с причиной) в строке прогона
agent_runs.
Унифицируем контракт. Скелет жизненного цикла один на всех потребителей; дом данных и смысл перезабора — доменные. Так монитор и сторож работают по любому прогону, не схлопывая разные домены в один стор.
queued · running · succeeded · failed (+ stale как под-итог сторожа)
attempts
счётчик попыток против потолка ретраев
heartbeat_at
пульс живого прогона; по нему сторож отличает мёртвого
last_error
причина терминала + статистика шагов
next_retry_at
отложенный старт ретрая; NULL, если не запланирован
Монитор читает каждую таблицу прогонов в её доменном доме и сводит их в
админ-сводке (module, job_id, state, attempts, age, error):
один вид «что и где упало» без центральной таблицы. Единая проекция
UNION ALL поверх таблиц прогонов — удобство
v2, когда таблиц прогонов станет много.