outbox

package
v1.3.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 10, 2026 License: Apache-2.0 Imports: 9 Imported by: 0

Documentation

Overview

Package outbox реализует транзакционный outbox-паттерн: каждое мутирующее действие на ресурс пишет строку в per-service outbox-таблицу в ТОЙ ЖЕ транзакции (см. Emit), а trigger pg_notify будит stream subscribers. Единый writer — outbox.Emit; drainer/reconciler читают backlog по sequence_no.

retention.go — УБОРКА ДОСТАВЛЕННЫХ СТРОК очереди дренажа: предикат и порог.

Предмет — задача продукта #1361. Дренаж помечает доставленную строку `sent_at` и НЕ УДАЛЯЕТ её никогда: и клейм, и применение, и повтор ключуются на `sent_at IS NULL`, а операторов снятия в этом пакете было ноль. Строка при этом заводится в writer-транзакции КАЖДОЙ мутации, то есть темп задаёт арендатор, и рост был монотонным и вечным у семи очередей шести владельцев.

Петля, партии, потолок проходов и наблюдаемость живут в `pkg/retention` и здесь не повторяются. Здесь — только ПРЕДИКАТ и его порог, потому что предикат есть свойство ЭТОЙ семьи таблиц: у очереди дренажа признак доставки — `sent_at`, и он принадлежит дренажу, живущему рядом.

СЕМЬ ОЧЕРЕДЕЙ, ОДИН ПРЕДИКАТ — и почему не семь копий

Семь расписаний об одном предмете разошлись бы молча и дали бы разный срок жизни доставленной строки у разных доменов при ОДНОМ контракте доставки. Имя таблицы и ключ партиции приезжают сюда ЗНАЧЕНИЯМИ, а оператор — один.

Index

Constants

View Source
const DeliveredRetention = 7 * 24 * time.Hour

DeliveredRetention — сколько доставленная строка остаётся в очереди.

ВЕЛИЧИНА ВЫВЕДЕНА ИЗ ЧИТАТЕЛЯ, а не назначена числом

Читателей у доставленной строки ДВА, и они требуют разного.

  1. `reconciler.RedrivePoisoned` — читает доставленных ПРЕЕМНИКОВ, решая, можно ли оживить отравленную строку. Его требование СТРУКТУРНОЕ, а не временно́е, и закрывается оно предикатом ниже, а не порогом: снять защищающую строку нельзя НИКОГДА, сколько бы ей ни было лет.

  2. ОПЕРАТОР, разбирающий «доехало ли снятие доступа». Вот его требование и задаёт величину: столько, сколько разбирают инцидент. Та же неделя, что у таблицы операций и у журнала подписки, и по той же причине — нерабочее время: упавший в пятницу вечером разбирается в понедельник утром.

ПОЧЕМУ ТА ЖЕ ВЕЛИЧИНА, ЧТО У СОСЕДЕЙ, И ПОЧЕМУ ОНА НЕ ПЕРЕИСПОЛЬЗУЕТСЯ

`pkg/operations.OperationRetention` и `pkg/subscription.JournalRetention` выведены из СВОИХ читателей и совпали числом. Совпадение чисел не есть один предмет: связать их константой значило бы, что сдвиг срока разбора отказов двигает срок жизни доставленной строки очереди. Величины объявляются порознь, и каждая — рядом со своим читателем.

У ПОРОГА НЕТ СЛАГАЕМОГО НА РАЗНИЦУ ЧАСОВ, и это замер, а не упущение

Отметку `sent_at` ставит сам дренаж оператором `UPDATE … SET sent_at = now()`, то есть часами БАЗЫ — теми же, которыми судит оператор ниже. Слагаемое, покрывающее разницу, которой нет, было бы запасом без предмета. У таблицы операций такое слагаемое ЕСТЬ (`pkg/operations.ClockDrift`) именно потому, что её отметку пишет процесс.

Variables

This section is empty.

Functions

func Emit

func Emit(ctx context.Context, tx pgx.Tx, table, kind, id, eventType string, payload map[string]any) error

Emit вставляет одну outbox-строку в произвольную таблицу с фиксированной схемой: (sequence_no BIGSERIAL PK, resource_kind TEXT, resource_id TEXT, event_type TEXT, payload JSONB, created_at TIMESTAMPTZ DEFAULT now()).

table — имя таблицы (например "vpc_outbox"). kind — тип ресурса ("Network", "Subnet", "Address", "RouteTable", "SecurityGroup"). id — ID ресурса (TEXT, поддерживает любой формат — UUID, короткий непрозрачный id и т.д.). eventType — "CREATED" | "UPDATED" | "DELETED". payload — произвольная map (сериализуется в JSONB).

Должна вызываться внутри pgx.Tx, в которой выполняется INSERT/UPDATE/DELETE целевой ресурсной таблицы — это обеспечивает атомарность outbox-write.

На каждый INSERT срабатывает trigger pg_notify('<channel>', sequence_no), который будит подписанных stream subscribers.

func EmitAnchored

func EmitAnchored(ctx context.Context, tx pgx.Tx, table, kind, id, projectID, eventType string, payload map[string]any) error

EmitAnchored — то же, но с ПРОЕКТНЫМ ЯКОРЕМ в колонке `project_id`.

Почему отдельная функция, а не параметр у Emit

Якорь — свойство ЖУРНАЛА, а не вызова: колонка либо есть в таблице, либо нет, и вставка в несуществующую колонку — отказ хранилища, а не пустое значение. Единая сигнатура заставила бы владельцев журналов без колонки передавать что-нибудь, и «что-нибудь» уезжало бы в SQL, который у них не исполнится.

Зачем якорь нужен КОЛОНКОЙ

Общая форма подписки (`pkg/subscription`) несёт проектный якорь полем оболочки события и принимает по нему решение о показе — не обращаясь к предмету. Для события снятия это несущее: обращаться не к чему, а нагрузка снятия несёт один идентификатор. Владелец, у которого якорь только в нагрузке, отдаёт у снятий пустой якорь — то есть утверждение «предмет уровня аккаунта», — и подписка с осью проекта такие события молча не пропускает. Потребитель, снявший опрос, об удалении не узнаёт никогда.

Кто её зовёт — СПРАШИВАЕТСЯ У ДЕРЕВА, а не перечисляется здесь

Здесь стоял поимённый перечень зовущих, и он пережил свой предмет: журнал, названный в нём среди «колонки не несут и зовут Emit», к тому дню колонку нёс и звал ЭТУ функцию. Направление ошибки было худшим из возможных — перечень ЗАНИЖАЛ, то есть посылал следующего владельца журнала брать безъякорную форму, а это ровно тот тихий отказ, о котором предупреждает раздел выше.

Рукописный перечень в общей библиотеке протухает молча и читается как доказательство, поэтому он снят, а не переписан на сегодняшнее число. Спрашивайте дерево — обе команды дают ответ за секунду:

git grep -n 'outbox.EmitAnchored(' -- services/ ':!*_test.go'   # кто зовёт якорную
git grep -n 'outbox.Emit('         -- services/ ':!*_test.go'   # кто зовёт безъякорную

Несёт ли КОНКРЕТНЫЙ журнал колонку якоря — вопрос к его миграциям, а не к вызывающему:

git grep -rn '<таблица> ADD COLUMN IF NOT EXISTS project_id' -- services/

Вызов этой библиотеки — НЕ единственная форма записи журнала в дереве, и перечислять прочие здесь значило бы завести второй рукописный перечень взамен снятого. Их перепись ВЫВОДИТСЯ обходом в internal/repohygiene/journalwriteforms.go: она печатает по каждой форме и то, сколько её экземпляров распознаватель нашёл в дереве вообще, поэтому «ноль» там отличимо от «не искали».

Почему функций всё-таки две

Довод не зависит от того, кто их зовёт сегодня, и потому переживёт любой перечень: якорь — свойство ХРАНИЛИЩА, колонка либо есть, либо нет, и вставка в несуществующую колонку есть отказ базы. Две функции — не два языка об одном предмете, а две РАЗНЫЕ формы хранения, каждая названная вслух. Сведение журналов к одной форме заведено отдельным предметом; до него выбор функции диктует таблица, а не вкус вызывающего.

func SanitizeTable

func SanitizeTable(table string) string

SanitizeTable квотирует имя таблицы (опц. схема-квалифицированное "schema.table") через pgx.Identifier — идентификатор экранируется библиотекой независимо от дисциплины вызывающего. Даже при контрактe «caller передаёт literal» это defense-in-depth: имя таблицы больше не может стать вектором statement-injection при interpolation в `INSERT INTO %s`.

Единый source-of-truth для всех outbox-подпакетов (reconciler/metrics), чтобы политика квотирования имён не расходилась между ними.

func StartQueueRetentionSweep

func StartQueueRetentionSweep(
	ctx context.Context,
	db QueueExecer,
	cfg QueueRetentionConfig,
	rcfg retention.Config,
	log *slog.Logger,
) (*retention.Sweeper, error)

StartQueueRetentionSweep поднимает фоновую уборку доставленных строк очереди.

ОДНА строка в композиционном корне владельца — намеренно: шесть владельцев несут семь очередей одной формы с одним контрактом доставки, и всё, что у них могло бы разойтись, собрано здесь.

Отказывает, а не предупреждает: отказ означает негодные величины расписания либо объявление, при котором уборка небезопасна.

Types

type QueueExecer

type QueueExecer interface {
	Exec(ctx context.Context, sql string, args ...any) (pgconn.CommandTag, error)
}

QueueExecer — то, чем уборщик исполняет свой оператор.

type QueueRetentionConfig

type QueueRetentionConfig struct {
	// Table — полное имя таблицы очереди («<схема>.<таблица>» либо «<таблица>»).
	Table string
	// PartitionColumn — ключ партиции ПОРЯДКА: тот же, которым пользуются клейм
	// дренажа и анти-джойн реконсайлера. ОБЯЗАТЕЛЕН, см. NewQueueSweeper.
	PartitionColumn string
}

QueueRetentionConfig — объявление владельца очереди.

type QueueSweeper

type QueueSweeper struct {
	// contains filtered or unexported fields
}

QueueSweeper — уборщик доставленных строк ОДНОЙ очереди.

func NewQueueSweeper

func NewQueueSweeper(db QueueExecer, cfg QueueRetentionConfig) (*QueueSweeper, error)

NewQueueSweeper собирает уборщика по объявлению владельца.

ОТКАЗЫВАЕТ БЕЗ КЛЮЧА ПАРТИЦИИ, и это fail-closed, а не придирка к форме: без ключа предикат «пощадить защищающую строку» НЕВЫРАЗИМ, и собранный молча уборщик снимал бы защиту отзыва у каждой очереди. Отказ наступает на сборке, потому что тихо собранный уборщик обнаружился бы воскрешённым доступом.

func (*QueueSweeper) RetentionSubject

func (s *QueueSweeper) RetentionSubject() retention.Subject

RetentionSubject — запись реестра уборки для этой очереди.

Порог собирается ЗДЕСЬ, а не у каждого из шести владельцев: семь копий одного порога разошлись бы молча и дали бы разный срок жизни доставленной строки у разных доменов при ОДНОМ контракте доставки.

func (*QueueSweeper) Sweep

func (s *QueueSweeper) Sweep(ctx context.Context, grace time.Duration, batch int) (int64, bool, error)

Sweep — один проход партии.

Момент времени входом НЕ приходит: часы уборки — базы, те же, что ставят отметку доставки (см. DeliveredRetention). Проверяемость от этого не теряется — проход зовётся методом, поэтому пробе не приходится ни спать, ни подменять часы.

Возвращает число снятых строк и признак «партия ушла полной»; признак отличает «убрал всё, что было» от «упёрся в партию».

Directories

Path Synopsis
Package bootgate is the fail-closed boot gate for services that must have a working IAM-register delivery path before they accept mutating Creates.
Package bootgate is the fail-closed boot gate for services that must have a working IAM-register delivery path before they accept mutating Creates.
Package drainer реализует универсальный outbox-drainer для Kachō outbox-pattern.
Package drainer реализует универсальный outbox-drainer для Kachō outbox-pattern.
Package metrics exposes the outbox-delivery observability surface: backlog depth, oldest-pending age and poison count per outbox table/channel.
Package metrics exposes the outbox-delivery observability surface: backlog depth, oldest-pending age and poison count per outbox table/channel.
Package reconciler is the backstop layer over the outbox/drainer.
Package reconciler is the backstop layer over the outbox/drainer.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL