Documentation
¶
Overview ¶
Package drainer реализует универсальный outbox-drainer для Kachō outbox-pattern.
Drainer слушает LISTEN/NOTIFY-канал Postgres, дренит pending rows на старте (catch-up), декодирует payload через caller-supplied Decoder[T] и применяет каждую row через caller-supplied Applier[T] к target-системе.
Свойства ¶
- **Идемпотентность**: Applier возвращает ErrAlreadyApplied → drainer трактует как success и mark'ит sent_at. Это позволяет at-least-once дренаж быть exactly-once на бизнес-уровне (caller возвращает ErrAlreadyApplied, когда адресат сообщил «уже есть»). Отказ, при котором адресат НЕ ЗАПИСАЛ НИЧЕГО, сюда не относится и обязан остаться transient: пометив такую строку отправленной, дренаж потерял бы намерение молча.
- **Exactly-once across HA replicas**: claim открывает транзакцию с `SELECT … FOR UPDATE SKIP LOCKED` и держит row-lock на время Apply. Другие реплики drainer-а SKIP'нут lock'нутый row до commit'а текущей. within-service-инварианты выражаются на DB-уровне.
- **Exp-backoff retry**: transient errors (все, что не ErrAlreadyApplied и не ErrPermanent) → retry с backoff [BackoffMin..BackoffMax] на каждый attempt; следующий NOTIFY/poll переклеймит row.
- **Poisoned-skip**: permanent errors (`errors.Is(err, ErrPermanent)`) или decoder-fail → force attempt_count = MaxAttempts, drainer больше не переклеймит. last_error содержит сообщение для debugging. Operator при необходимости reset вручную (UPDATE … SET attempt_count = 0).
- **Graceful shutdown**: ctx.Done() → дозавершает текущий in-flight Apply (использует detached ctx с ApplyTimeout grace), затем exit.
- **LISTEN reconnect**: на drop LISTEN-conn'а → reconnect с exp-backoff (1s → 30s cap), pool.Reset() форсит destroy idle pool-conn'ов (могут быть тоже мертвы при общем admin-shutdown / FATAL); после reconnect — catch-up wakeup (NOTIFY могли быть потеряны во время disconnect).
Архитектура ¶
Drainer параметризован generic-типом T (decoded payload). Один Drainer[T] — один outbox-table + один applier. Для двух outbox-таблиц (e.g. fga_outbox + subject_change_outbox в kaname) запускайте две независимые Drainer-инстанции.
LISTEN использует dedicated pgx-connection, hijacked из переданного pool'а. Hijack необходим, потому что LISTEN-state не выживает pool's idle-conn recycle. Conn закрывается при reconnect / shutdown.
Claim/apply/mark идут через pool в одной транзакции на batch (1..4 rows random для HA-fairness). Транзакция держит row-lock от claim до mark'а, исключая race-window «claim → apply done, но markSuccess не успел» при HA.
Пример (kaname fga_outbox) ¶
type FGAOutboxEvent struct {
User, Relation, Object string
}
d, err := drainer.New[FGAOutboxEvent](
pool,
drainer.Config{
Table: "kaname.fga_outbox",
Channel: "kacho_iam_fga_outbox",
BatchSize: 32,
PollFallback: 30 * time.Second,
MaxAttempts: 10,
BackoffMin: 1 * time.Second,
BackoffMax: 30 * time.Second,
},
func(p []byte) (FGAOutboxEvent, error) {
var e FGAOutboxEvent
if err := json.Unmarshal(p, &e); err != nil {
return e, errors.Join(drainer.ErrPermanent, err)
}
return e, nil
},
func(ctx context.Context, eventType string, e FGAOutboxEvent) error {
switch eventType {
case "fga.tuple.write":
err := openFGA.WriteTuples(ctx, ...)
if isConflict(err) { return drainer.ErrAlreadyApplied }
return err
case "fga.tuple.delete":
err := openFGA.DeleteTuples(ctx, ...)
if isMissing(err) { return drainer.ErrAlreadyApplied }
return err
default:
return errors.Join(drainer.ErrPermanent,
fmt.Errorf("unknown event_type %q", eventType))
}
},
logger,
)
if err != nil { return err }
return d.Run(ctx) // blocks until ctx.Done()
См. также ¶
- Writer-side (Emit): пакет github.com/PRO-Robotech/kacho/pkg/outbox
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ErrAlreadyApplied = errors.New("drainer: target reports already-applied (idempotent)")
ErrAlreadyApplied — applier возвращает, когда адресат сообщил «уже есть»: повторная запись того, что у него уже лежит, либо снятие того, чего у него уже нет. Drainer трактует как success.
ВНИМАНИЕ: отказ, при котором адресат НЕ ЗАПИСАЛ НИЧЕГО (транзакционный обрыв, конфликт с параллельным писателем), сюда НЕ относится, как бы он ни назывался на проводе. Пометив такую строку sent_at, дренаж молча потерял бы намерение — перманентная дыра в правах; такой отказ обязан остаться transient, его и лечит повтор. Различить их может только applier: словарь адресата знает он один.
var ErrPermanent = errors.New("drainer: permanent error, no retry")
ErrPermanent — applier wrap'ит в это, если retry бессмыслен (HTTP 4xx кроме conflict-класса, malformed payload, etc). Drainer mark'ит last_error и пропускает row через force attempt_count = MaxAttempts.
Functions ¶
This section is empty.
Types ¶
type Applier ¶
Applier[T] — применяет T к target-системе. Возвращает nil → success (drainer mark'ит sent_at). Возвращает ErrAlreadyApplied → idempotent success (drainer mark'ит sent_at). Возвращает любую другую error → transient (retry с exp backoff)
ИЛИ permanent (если errors.Is(err, ErrPermanent)).
CONCURRENCY: при Config.ApplyConcurrency>1 drainer вызывает Applier из НЕСКОЛЬКИХ горутин одновременно (разные строки одного claim-батча) — Applier ОБЯЗАН быть безопасен для конкурентного вызова. Канонические appliers это gRPC-клиенты поверх *grpc.ClientConn (конкурентно-безопасен) + чистое построение запроса — они уже удовлетворяют требованию. Applier, шарящий mutable-состояние без синхронизации, обязан либо синхронизироваться, либо оставить ApplyConcurrency=1.
type Class ¶
type Class int
Class is the outcome of classifying an applier (or decoder) error. It is the single, testable decision point that drives whether the drainer marks a row success, poisons it (no further retry) or retries it unbounded with backoff.
Transient-class no-poison rule: a long-but-transient IAM outage (gRPC Unavailable / DeadlineExceeded / connection-refused / timeout) — and a concurrency conflict (409) — must NEVER poison a row: it retries forever (with backoff) so the owner-tuple is never lost to a temporary peer outage or a racing writer. Only permanent errors (4xx other than the conflict class, decode-failure, malformed) poison.
const ( // ClassSuccess — nil error; the row is delivered. ClassSuccess Class = iota // ClassAlreadyApplied — the target reports already-applied; idempotent success, // the row is marked sent. What decides the class is not the wire code but what // the target DID: "it is already there, nothing left to do" is success, while an // abort that applied NOTHING is ClassTransient — marking such a row sent would // silently drop the intent. Only the applier can tell the two apart, because // only it knows its target's vocabulary. ClassAlreadyApplied // ClassPermanent — retry is pointless (ErrPermanent, gRPC InvalidArgument / // 4xx-non-409, decode-failure). The row is poisoned (attempt_count forced to // MaxAttempts) and surfaced for an operator. ClassPermanent // ClassTransient — a temporary failure (peer Unavailable / DeadlineExceeded / // connection-refused / timeout / any unclassified error). The row is retried // unbounded with backoff and is NEVER driven into the poison gate. ClassTransient )
func Classify ¶
Classify maps an applier error to a Class.
Decision order (most-specific first):
- nil → ClassSuccess
- errors.Is(err, ErrAlreadyApplied) → ClassAlreadyApplied
- errors.Is(err, ErrPermanent) → ClassPermanent
- gRPC InvalidArgument / PermissionDenied → ClassPermanent
- everything else (Unavailable, DeadlineExceeded, NotFound, FailedPrecondition, connection-refused/timeout, raw) → ClassTransient
Rationale for step 4: both codes describe a decision about the REQUEST — its content, or the caller's authority to make it — and an identical retry changes neither. See isPermanentGRPC for why calling a refusal "transient" produces a permanently wedged partition rather than an eventual success.
Rationale for step 5: the remaining codes describe peer STATE, which a retry can genuinely find changed. A raw, un-wrapped error of unknown shape is likewise treated as transient — fail-SAFE for delivery (retry rather than lose the tuple). Appliers that KNOW an error is permanent must wrap it in ErrPermanent.
type Config ¶
type Config struct {
// Table — полное имя outbox-таблицы (`<schema>.<table>`), e.g. "kaname.fga_outbox".
Table string
// Channel — имя LISTEN-канала, e.g. "kacho_iam_fga_outbox".
Channel string
// BatchSize — сколько rows клейм'ить за один catch-up SELECT (default 32).
BatchSize int
// PollFallback — интервал poll'а на случай missed NOTIFY (default 30s).
PollFallback time.Duration
// MaxAttempts — отметка «poisoned», после которой drainer перестает ретраить
// (default 10). Permanent-error → force attempt_count = MaxAttempts, drainer пропускает.
MaxAttempts int
// BackoffMin/BackoffMax — exp-backoff bounds (default 1s..30s).
BackoffMin time.Duration
BackoffMax time.Duration
// ApplyTimeout — таймаут на один Apply-вызов (default 5s). Используется
// также как graceful-grace при shutdown: in-flight Apply имеет
// non-cancellable inner ctx с этим deadline, чтобы row не осталась
// half-applied при ctx.Cancel parent-loop'а.
ApplyTimeout time.Duration
// ApplyConcurrency — сколько строк одного claim-батча применять ПАРАЛЛЕЛЬНО
// (default 1 = последовательно, историческое поведение). >1 разворачивает
// внешние Apply-вызовы батча по N горутинам, скрывая per-call latency пира:
// одиночный drainer с последовательными ApplyTimeout-bounded apply'ями
// упирается в ~1/apply_latency (при таймаутящем пире — ~1/ApplyTimeout, что
// катастрофически мало под write-burst). Claim-батч при ApplyConcurrency>1
// сайзится ровно в ApplyConcurrency, поэтому вся волна применяется за один
// проход параллельно; mark'и остаются ПОСЛЕДОВАТЕЛЬНЫМИ на единственной
// claim-транзакции. Exactly-once НЕ меняется: та же claim-tx держит
// FOR UPDATE SKIP LOCKED lock КАЖДОЙ заклейменной строки до commit'а, а
// Apply не трогает DB-состояние (только внешний вызов) → параллельные apply
// не создают ни второй tx, ни лишних conn'ов пула. Требование: Applier
// БЕЗОПАСЕН для конкурентного вызова (см. Applier godoc).
//
// ORDERING (важно): drainer НЕ гарантирует apply в id-порядке — ни при
// ApplyConcurrency>1 (строки батча применяются конкурентно), ни при
// ApplyConcurrency=1: claim сортирует `ORDER BY attempt_count, id`, поэтому
// transiently-bumped предшественник (attempt≥1) уступает свежему преемнику
// (attempt=0) и применяется ПОСЛЕ него — внутри одного батча и тем более между
// батчами. Порядок ВНУТРИ ключа даёт ТОЛЬКО PartitionColumn (см. ниже);
// ApplyConcurrency — это throughput-ручка, а НЕ источник (и не причина)
// reorder'а.
//
// Держать PartitionColumn пустым допустимо ТОЛЬКО когда финальное состояние
// target'а СХОДИТСЯ при любом порядке применения строк ОДНОГО ключа:
// (а) поток по ключу append-only/коммутативен — нет события, отменяющего
// предыдущее (напр. только register/label-update, без unregister); либо
// (б) target версионирует ОБЕ ветки — и upsert, и удаление (versioned
// tombstone), так что stale-apply любого вида no-op'ится.
//
// Канонические register-outbox'ы kacho (`<svc>.fga_register_outbox`) под (а) и
// (б) НЕ подпадают и ОБЯЗАНЫ задавать PartitionColumn (ключ = колонка ресурса,
// напр. `resource_id`): они несут register И unregister ОДНОГО объекта, а
// материализация в iam версионирована лишь ЧАСТИЧНО. source_version-LWW
// (`resource_mirror` UPSERT `WHERE source_version < EXCLUDED.source_version`,
// services/iam/.../resource_mirror/emitter.go) гейтит ТОЛЬКО ветку
// ON CONFLICT DO UPDATE, т.е. спасает лишь register↔register. Пара
// register(t1)→unregister(t2>t1) НЕ коммутативна и НЕ защищена: unregister
// делает ЖЁСТКИЙ DELETE без tombstone, поэтому переставленный stale register
// попадает в ветку INSERT (сравнивать не с чем) и ВОСКРЕШАЕТ mirror-строку
// удалённого ресурса; level-triggered reconciler читает mirror как источник
// истины и вечно ре-материализует owner-tuple, самоисцеления нет. Поведение
// закреплено Test_1_4_45_RegisterOutbox_UnregisterThenStaleRegister
// (register-outbox без PartitionColumn → resurrect; с PartitionColumn →
// корректное ABSENT).
//
// Тот же класс разбирался на журнале намерений iam `fga_outbox`: он несёт
// СЫРОЙ, вообще не версионированный owner-hierarchy tuple, а WRITE(grant) и
// DELETE(revoke) одного (user,relation,object) не коммутируют, поэтому там
// стоял PartitionColumn (`tuple_key` — материализованный триггером триплет)
// вместе с ApplyConcurrency>1. Дренажа у этого журнала больше нет — он снят
// вместе со своим адресатом, — но пример остаётся: правило про ширину ключа
// от него не зависит.
ApplyConcurrency int
// PartitionColumn — SQL-выражение над строкой outbox-таблицы, дающее ключ
// ПАРТИЦИИ порядка (обычная колонка: `resource_id` для register-outbox'ов).
// Пусто (default) = фича выключена,
// claim-запрос БАЙТ-в-БАЙТ прежний (нулевое изменение поведения для всех
// consumer'ов, которые её не задают).
//
// # Ширина ключа: САМЫЙ УЗКИЙ ключ, на котором события НЕ коммутируют
//
// Ключ выбирается не «по ресурсу», а по НЕКОММУТАТИВНОСТИ: партиция обязана
// совпадать с той единицей состояния target'а, которую события перезаписывают
// друг у друга. Всё, что шире, — чистая ПЕРЕсериализация: она ничего не
// добавляет к корректности, но ставит каждое событие в очередь за событиями, до
// которых ему нет дела (радиус wedge тоже растёт до всей широкой группы).
//
// Канонический промах — iam `fga_outbox` до миграции 0067: партиция была
// `payload->>'object'`, тогда как состояние адресата было МНОЖЕСТВОМ кортежей
// с ключом (user, relation, object). WRITE(user:a,v_get,P) и DELETE(user:b,v_get,P)
// трогают РАЗНЫЕ элементы и коммутируют; общий у них только object. Замер на
// живом стенде под пакетным e2e: 8 643 pending-строки, 1 439 объектов, но 8 641
// РАЗЛИЧНЫХ tuple'ов — то есть широкий ключ давал 1 439 claimable-голов вместо
// 8 641, один revoke ждал до 632 предшественников по объекту при ≤3 по своему
// tuple'у, и claim отбрасывал ~3 из 4 просмотренных строк как «не голова»
// (68 кандидатов на 16 заклеймленных против ровно 16 на узком ключе). Под
// нагрузкой это выносило видимость отзыва прав на ~30 c при бюджете пробы 15 c.
// Регрессию (обе стороны — и порядок, и отсутствие лишней сериализации) локают
// Test_1_4_46_PartitionKey_SameObject_DistinctSubjects_StayConcurrent и
// Test_1_4_47_NarrowKey_PreservesOrder_UnderSameObjectNoise.
//
// Сужение ключа безопасно ровно до этой границы и НЕ создаёт окна: предикат не
// меняется, меняется только его equi-ключ, поэтому пара «WRITE→DELETE одного
// ключа» по-прежнему упорядочена cross-batch и cross-replica
// (Test_1_4_48_NarrowKey_OrderHolds_AcrossReplicas). Ошибка в другую сторону —
// ключ УЖЕ реального (расщепить одну единицу состояния на две партиции) — вот
// это и есть потеря порядка; поэтому при сомнении в разделителе/рендеринге ключа
// выбирай тот, который может только СКЛЕИТЬ две единицы (пере-упорядочить), а не
// расщепить одну.
//
// # Зачем (order-preserving drain без cross-batch reorder-leak)
//
// Порядок apply НЕ сохраняется и БЕЗ конкурентности: claim
// `ORDER BY (attempt_count,id)` разносит bumped-предшественника (attempt≥1 после
// transient) и fresh-преемника (attempt=0) — преемник обгоняет предшественника
// уже внутри одного батча при ApplyConcurrency=1 и тем более попадает в БОЛЕЕ
// РАННИЙ батч. ApplyConcurrency>1 добавляет сверху ещё и intra-batch reorder
// (горутины конкурентны), но не является причиной проблемы. Для order-sensitive
// target'а это ломается двумя способами: iam raw-tuple (WRITE потом DELETE того
// же tuple) → delete-before-write → tuple ВЫЖИВАЕТ → authz over-grant /
// cross-account LEAK; register-outbox (register потом unregister того же
// ресурса) → unregister-before-register → воскрешённая mirror-строка удалённого
// ресурса (см. ORDERING в ApplyConcurrency).
//
// PartitionColumn закрывает это на CLAIM-уровне (не на apply-re-sort, который
// бессилен, если предшественник в другом батче): claim НЕ берёт строку t, если
// в её партиции существует ДОСТАВЛЯЕМЫЙ (sent_at IS NULL AND
// attempt_count < MaxAttempts) предшественник с меньшим id. Тогда:
// - преемник НИКОГДА не заклеймлен впереди unsent-предшественника (cross-batch
// И cross-replica: незакоммиченный claim соседа виден как sent_at IS NULL в
// снапшоте → его преемник остаётся заблокирован до commit'а) →
// per-partition FIFO держится by construction;
// - в любом снапшоте claimable ровно ОДНА строка партиции (её head), поэтому
// один claim-батч НИКОГДА не содержит две строки одной партиции → intra-batch
// конкурентный reorder двух строк одной партиции НЕВОЗМОЖЕН by construction
// (apply-level partition-grouping не нужен — было бы vestigial).
// Строки РАЗНЫХ партиций по-прежнему клеймятся и применяются конкурентно
// (пропускная способность ApplyConcurrency сохранена — предикат сериализует
// только внутри одной партиции).
//
// # Head-of-line wedge (осознанный trade-off leak-safety > per-partition liveness)
//
// Пока head партиции доставляем-но-ещё-не-доставлен (transient-stuck: peer down,
// attempt капнут на MaxAttempts-1, ретраится вечно с backoff) — его преемники
// ЖДУТ. Это ВРЕМЕННЫЙ wedge: transient-строка НИКОГДА не отравляется (cap ниже
// poison-gate) → применится, как только peer оживёт → wedge рассосётся. Границей
// wedge является восстановление peer'а, не бесконечность. ОТРАВЛЕННЫЙ (poisoned,
// attempt_count=MaxAttempts) предшественник из блокирующего набора ИСКЛЮЧЁН
// (предикат `attempt_count < MaxAttempts`), т.к. он НИКОГДА не применится —
// вечно wedge'ить партицию им нельзя; преемник «перепрыгивает» мёртвого
// предшественника. Leak это НЕ реинтродуцирует: отравленный WRITE не создаёт
// tuple → последующий DELETE (или отсутствие эффекта) не оставляет над-гранта;
// порядок между ДОСТАВЛЯЕМЫМИ строками партиции по-прежнему строгий. Wedge
// наблюдаем: (а) table-wide `outbox_oldest_pending_age_seconds{table}` (metrics
// Collector) растёт, пока любой head застрял → alertable; (б) per-partition
// атрибуция — опциональный WithWedgeObserver + WedgeWarnAfter (см. ниже).
//
// # Контракт выражения
//
// PartitionColumn — выражение над НЕквалифицированными колонками строки,
// начинающееся с идентификатора колонки, так что префикс `<alias>.` даёт
// валидный квалифицированный SQL. ПРЕДПОЧТИТЕЛЬНА обычная колонка (`tuple_key`,
// `resource_id`): claim'у нужен partial-индекс ровно по этому ключу, и колонка
// делает его простым btree, который планировщик не может «не сматчить»;
// составной ключ материализуй колонкой (миграция + BEFORE INSERT триггер как
// единственный источник истины для ВСЕХ writer'ов), а не встраивай выражение —
// разъехавшийся рендеринг ключа у одного из writer'ов расщепляет одну единицу
// состояния на две партиции и молча теряет порядок. Значение используется в
// equi-сравнении партиций в claim-предикате; в claim-пути оно не логируется, но
// per-partition wedge-observer (WithWedgeObserver) НАМЕРЕННО surfaces ключ в WARN
// для атрибуции застрявшей партиции — это id-based handle (напр.
// `user:<id> v_get vpc_network:<id>`), не PII/инфра-чувствительное. Rows с
// NULL-ключом партиции трактуются как независимые (NULL=NULL не true) — поэтому
// таблица обязана гарантировать непустой ключ (у iam — CHECK ... NOT VALID).
//
// # Обязательные индексы — РОВНО ДВА, и «ровно» здесь так же важно, как «два»
//
// Миграция сервиса ОБЯЗАНА нести ДВА partial-индекса (оба `WHERE sent_at IS
// NULL`, чтобы размер тянулся за PENDING-backlog'ом, а не за append-mostly
// таблицей) — и НИ ОДНОГО ТРЕТЬЕГО:
//
// 1. `(<PartitionColumn>, id) WHERE sent_at IS NULL` — коррелированный
// NOT EXISTS (поиск предшественника партиции);
// 2. `(attempt_count, id) WHERE sent_at IS NULL` — ВНЕШНИЙ упорядоченный
// скан `ORDER BY (attempt_count, id)`.
//
// (1) БЕЗ (2) — НЕДОСТАТОЧНО, и это не «чуть медленнее», а throughput-инверсия.
// Без упорядоченного access-path планировщик вообще не доходит до
// nested-loop anti-join, ради которого построен (1): он seq-скан'ит ВЕСЬ
// pending-backlog, сортирует его и hash-anti-join'ит со ВТОРЫМ полным
// seq-скан'ом того же backlog'а — (1) при этом не используется, claim
// становится O(backlog). Дальше включается положительная обратная связь: чем
// глубже очередь, тем медленнее claim → тем медленнее дренаж → тем глубже
// очередь; producer, обгоняющий consumer хотя бы на проценты, расходится
// неограниченно вместо небольшого постоянного лага. Замер на живом стенде
// (Postgres 16, kaname.fga_outbox), один claim: backlog 5k — 11.7ms без (2)
// против 0.81ms с (2); 20k — 61.6ms против 0.72ms; 80k — 327ms против 0.82ms
// (с (2) время от глубины не зависит). Регрессию локает
// Test_ClaimPlan_DoesNotScaleWithBacklogDepth.
//
// ЛЮБОЙ ТРЕТИЙ partial-индекс по `sent_at IS NULL` в ДРУГОМ порядке — decoy, и
// он возвращает ровно ту же инверсию, даже когда (1) и (2) на месте. Очередь
// почти всегда пуста, поэтому последний ANALYZE почти всегда снят на ПУСТОМ
// backlog'е, и в burst планировщик входит с оценкой `rows=1`; на ней Sort
// «бесплатен», и любой индекс, дающий дешёвый скан pending-строк в НЕ том
// порядке, выигрывает по стоимости. Итоговый план — «Index Scan(decoy) → Sort →
// Limit»: Sort обязан материализовать ВЕСЬ pending-набор прежде, чем сработает
// LIMIT, поэтому anti-join прогоняется по разу на КАЖДУЮ pending-строку, а (1)
// и (2) просто не используются. Замер на живом стенде, тот же claim, 8 600
// pending, разница только в наличии лишнего `(created_at) WHERE sent_at IS NULL`:
// 6 990 ms против 29 ms (в 240 раз). В логах drainer'а это выглядело как ровно
// один 16-строчный батч раз в ~11.5 c в течение 46 c (1.4 строки/с) с прыжком до
// ~600 строк/с в момент, когда autovacuum наконец переанализировал таблицу.
// Держи индексный набор МИНИМАЛЬНЫМ: не «добавь индекс — хуже не будет», а
// «каждый лишний порядок над pending-строками — это план, который планировщик
// выберет именно тогда, когда очередь глубокая». Per-table autovacuum-настройки
// (у iam — миграция 0064) сокращают ОКНО плохих статистик, но не устраняют
// плохой ПЛАН: autovacuum_naptime кластерный (1 мин), поэтому окно всё равно
// есть. Регрессию локает Test_ClaimPlan_StaleStatistics_DoesNotCollapse.
PartitionColumn string
// WedgeWarnAfter — если >0 И задан PartitionColumn И зарегистрирован
// WithWedgeObserver, drainer периодически (rate-limited) сообщает партиции, чей
// самый старый unsent-row старше этого порога (head-of-line wedge, см. выше).
// Сообщаются ТОЛЬКО N самых старых (wedgeReportLimit): задача observer'а —
// АТРИБУЦИЯ («какая партиция стоит»), а «сколько стоит» без кардинальности уже
// отвечают table-wide backlog-depth / oldest-pending-age метрики. Ключ узкий, и
// при выпадении peer'а под порог попадают ВСЕ партиции сразу — без cap'а это
// тысячи WARN-строк за скан, которые хоронят собственно атрибуцию.
// 0 (default) = per-partition wedge-репорт выключен (нулевая стоимость; table-wide
// oldest-pending-age метрика всё равно покрывает «доставка застряла»). Опрос
// использует тот же partial-index, что и claim.
WedgeWarnAfter time.Duration
// PermanentPolicy — что делать с ПОСТОЯННЫМ отказом ПРИМЕНЕНИЯ. Умолчание
// (нулевое значение) — травить, то есть сегодняшнее поведение.
//
// RetryPermanent законен ТОЛЬКО без PartitionColumn: обоснование политики
// целиком опирается на отсутствие партиции. Пара отвергается Validate, а не
// оговаривается комментарием. Полностью — у типа PermanentPolicy.
PermanentPolicy PermanentPolicy
}
Config — параметры конкретного экземпляра drainer-а.
type Decoder ¶
Decoder[T] — превращает payload JSONB в типизированный T. Ошибка decoder-а трактуется как permanent (poisoned row), drainer не вызывает applier и помечает row attempt_count = MaxAttempts + last_error = err.
type Disposition ¶
type Disposition int
Disposition — что полагается сделать со строкой.
Отделено от места, где строка помечается: пометка требует транзакции и проверяется интеграционно, а решение обязано быть проверяемо без базы. Иначе единственная точка принятия решения снова размазалась бы по ветвям, которые поодиночке защитимы.
const ( // DispositionDeliver — пометить доставленной. DispositionDeliver Disposition = iota // DispositionPoison — отравить: повтор не имеет шанса на успех, и партиция // обязана разблокироваться. DispositionPoison // DispositionRetry — повторить с отступом, не доводя до порога отравления. DispositionRetry )
func Decide ¶
func Decide(cls Class, policy PermanentPolicy) Disposition
Decide — единственная точка, отвечающая «что сделать со строкой» по классу отказа применения и политике очереди.
func DecideOutcome ¶
func DecideOutcome(decodeErr, applyErr error, policy PermanentPolicy) Disposition
DecideOutcome — ЕДИНСТВЕННАЯ точка, отвечающая «что сделать со строкой» по обоим исходам её обработки и политике очереди. Ровно её зовёт путь пометки, поэтому здесь нет ветви, которой нет в работе.
Почему отказ РАЗБОРА не читает политику ¶
Политика повтора относится к отказу ПРИМЕНЕНИЯ — к тому, что сказал сосед, и что способно измениться от внешнего события. Отказ разбора говорит о самой строке: её тело не станет разбираемым ни от какого события, поэтому повтор — вечная работа без единого шанса на успех.
Следствие надо назвать вслух, а не проглотить: очередь с политикой повтора всё ещё МОЖЕТ отравить строку — по разбору. «Травиться нечему» достигается не терпимостью к битой строке, а тем, что такую строку НЕЛЬЗЯ ЗАПИСАТЬ: каждое условие отказа разбора обязано быть закрыто ограничением схемы у владельца очереди. Для очередей kaname это сделано миграциями 0079/0080 и 0097.
func (Disposition) String ¶
func (d Disposition) String() string
String — для журналов и текстов отказа проб.
type Drainer ¶
type Drainer[T any] struct { // contains filtered or unexported fields }
Drainer[T] — экземпляр drainer-а для одного outbox-table + один applier.
Drainer работает по схеме:
- Слушает LISTEN-канал `cfg.Channel` через dedicated pgx.Conn (hijacked из pool).
- На старте — catch-up: SELECT pending rows (sent_at IS NULL, attempt_count < MaxAttempts) ORDER BY attempt_count, id LIMIT BatchSize → applies each (attempt_count-first ordering предотвращает starvation свежего intent транзитно-залипшим backlog'ом).
- Main loop: wake-up по NOTIFY (payload = row id) ИЛИ tick(PollFallback) → claim → apply → mark.
- Exactly-once: pre-claim атомарный UPDATE … RETURNING с CAS (sent_at IS NULL AND attempt_count < MaxAttempts) — две реплики не возьмут одну row.
- Graceful shutdown: при ctx.Done() — дозавершает текущий in-flight apply (отдельный inner ctx с ApplyTimeout grace), exit.
func New ¶
func New[T any]( pool *pgxpool.Pool, cfg Config, decoder Decoder[T], applier Applier[T], logger *slog.Logger, opts ...Option[T], ) (*Drainer[T], error)
New создает Drainer; не запускает (вызывайте Run).
pool — *pgxpool.Pool на БД сервиса (тот же pool, что используется для бизнес-логики; drainer Acquire().Hijack() один conn для LISTEN, остальные операции — через pool как обычно).
decoder/applier — пользовательские функции, см. Decoder[T] / Applier[T].
logger — slog.Logger; nil → slog.Default().
func (*Drainer[T]) Run ¶
Run — основной loop drainer-а. Блокирует до ctx.Done().
Поведение:
- Запускает LISTEN-loop в goroutine (own conn, reconnect on drop).
- Выполняет startup catch-up (drains all pending rows).
- Основной select: NOTIFY-wake-up ИЛИ PollFallback-tick ИЛИ ctx.Done(). Каждое wake-up → drainBatch() (claim + apply + mark, в loop пока не пусто).
- На ctx.Done() — дозавершает текущий drainBatch (с inner ApplyTimeout-grace на in-flight apply), exits.
Возвращает nil при clean shutdown.
РЕПЛИКИ: клейм — строки берутся клеймом `FOR UPDATE SKIP LOCKED` по голове партиции, поэтому репликам достаются непересекающиеся партии, а порядок внутри партиции держится даже между партиями и репликами.
type Option ¶
Option customises a Drainer at construction (functional-options pattern).
func WithClaimObserver ¶
WithClaimObserver registers a callback invoked once per claim query issued against the outbox table. Enables deterministic in-process observation of claim frequency (busy-poll guard) independent of pg_stat lag. nil is ignored.
func WithDeliveryObserver ¶
WithDeliveryObserver registers a callback invoked once per DELIVERED row, with that row's event type. Wire it to a metrics recorder to make «доставлено по направлению» a MONOTONIC counter. nil is ignored.
Почему счётчик здесь, а не скан живых строк ¶
Величина объявлена «за всё время» и служит ЕДИНСТВЕННЫМ способом отличить «отзывов не было» от «отзывы не проходят». Считать её `count(*)` по живым строкам можно ровно до тех пор, пока строки не убираются: уборка доставленных (#1361) обнулила бы её на исправной очереди, где отзыв редок, — то есть уничтожила бы именно тот сигнал, ради которого разбивка по направлениям заведена. Наблюдатель считает СОБЫТИЕ доставки, и уборка на него не влияет by construction (#1714).
Зовётся ПОСЛЕ того, как исход строки признан доставкой, и до возврата из пометки; последовательно в пределах партии — как и onPoison.
func WithPoisonObserver ¶
WithPoisonObserver registers a callback invoked once per poisoned row. Wire it to a metrics recorder's IncPoisoned to make poison events observable (outbox_poisoned_total). nil is ignored.
func WithWedgeObserver ¶
WithWedgeObserver registers a callback invoked (rate-limited) once per partition whose oldest unsent row exceeds Config.WedgeWarnAfter — the head-of-line wedge signal for a partition-ordered drain. Requires Config.PartitionColumn set and Config.WedgeWarnAfter>0 (otherwise the scan is skipped and this is a no-op). Wire it to a WARN log / gauge for per-partition attribution. nil is ignored.
type PermanentPolicy ¶
type PermanentPolicy int
PermanentPolicy — что делать с ПОСТОЯННЫМ отказом ПРИМЕНЕНИЯ.
Почему у этого вопроса вообще появился второй ответ ¶
Травление покупает РОВНО ОДНО: отравленная строка выбывает из блокирующего набора заявки, и партиция, стоявшая за ней, разблокируется. Комментарий выше это и обосновывает — и обосновывает целиком через партицию: «with PartitionColumn set every later row of its partition is never claimed».
У КОММУТАТИВНОГО потока партиции нет. Блокировать нечего, разблокировать нечего — травление не покупает ничего, а платит полной ценой: намерение выбывает навсегда. Тот же комментарий называет и условие правильности — «poisoning is therefore only correct in a service that also re-drives poisoned rows», — а возврат отравленных строк строится только вокруг ключа порядка, которого у коммутативной очереди нет by construction. То есть для такой очереди травление одновременно бесполезно и невосполнимо.
Отсюда второй ответ, и он не «мягче», а уместнее: повторять с отступом. Цена названа честно — вечный повтор заведомо безнадёжного вызова с частотой не чаще BackoffMax, видимый счётчиком незавершённых строк. Цена альтернативы — молча потерянное намерение, и наблюдать её нечем: «доставлено» и «потеряно» снаружи выглядят одинаково.
Выведено 2026-08-16 по kacho#455 из двух очередей kaname, у которых травление работало, а возврата не было ни у одной.
const ( // PoisonPermanent — отравить строку. Умолчание и сегодняшнее поведение: // нулевое значение обязано означать то, что уже провязанные очереди делают // сейчас, иначе введение поля сменило бы их поведение молча. // // Требует возврата отравленных строк. Что это свойство ДЕРЕВА, а не чьей-то // памяти, держит internal/repohygiene TestEveryPoisoningOutboxHasARedrive. PoisonPermanent PermanentPolicy = iota // RetryPermanent — повторять, как временный отказ. // // Законно ТОЛЬКО у коммутативной очереди (PartitionColumn пуст): пара с // ключом порядка отвергается Config.Validate, потому что там постоянный // отказ заклинил бы свою партицию навсегда. RetryPermanent )
func (PermanentPolicy) String ¶
func (p PermanentPolicy) String() string
String — для журналов и меток.