operations

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: 24 Imported by: 0

Documentation

Overview

retention.go — уборка терминальных строк таблицы операций.

Предмет — задача продукта #1360 и приёмка `docs/specs/sub-phase-RET-PLAT-1-platform-table-retention-acceptance.md` (APPROVED, круг 3). Строка таблицы операций заводится КАЖДОЙ мутацией платформы: контракт объявляет мутации асинхронными, и `Operation` возвращается вместо ресурса. Снятия строк не было ни у одного из восьми владельцев — рост монотонный и вечный, а темп задаёт внешний.

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

Package operations — Long-Running Operations primitive: bounded Worker для async-исполнения мутаций + Repo для durable-перехода `done=false → true`.

Operations.Run() — fire-and-trigger pattern: handler возвращает Operation клиенту сразу, фоновый worker делает реальную работу и durable-переводит строку в терминал.

Bounded worker-pool (Worker type):

  • Run() / RunWithWorker() кладут задачу в in-memory admission backlog и сразу возвращают управление (async-контракт). Dispatcher-loop разбирает backlog, ограничивая число одновременно ИСПОЛНЯЕМЫХ worker'ов семафором max-inflight (burst не порождает unbounded-горутины). Backlog ограничен max-backlog: при переполнении задача НЕ кладется в память — операция durable в БД и добирается reconciler'ом (backpressure без OOM, без потери данных).
  • Active() — текущее число исполняемых worker'ов; Wait(ctx) дренирует их на graceful-shutdown; Ready() отражает живость dispatcher-loop (readiness probe).
  • panic в fn перехватывается recover() → durable MarkError, процесс не падает.

Durable terminal-write:

  • финальные MarkDone/MarkError идут через retry+backoff (pkg/backoff) поверх CAS-on-`done`; transient DB-сбой ретраится, метрики retries/failures не проглатываются; при исчерпании budget строка остается done=false и добирается reconciler'ом.

Context-propagation (baggage):

  • Worker НЕ наследует deadline / cancel callerCtx — request-ctx cancel-ится сразу как handler возвращает Operation, а worker живет независимо.
  • Worker НАСЛЕДУЕТ observability-values callerCtx (OTel SpanContext, request-id, slog logger) через baggage.Extract — иначе worker-логи и trace-span'ы оторваны от исходного запроса.
  • Поверх worker-ctx накладывается собственный per-op deadline (WithOperationTimeout, дефолт 4m < OrphanGrace): зависший peer-вынов не удерживает слот семафора бесконечно — по timeout'у fn получает ctx.Done, операция durable-помечается DeadlineExceeded.

Index

Constants

View Source
const AnonymousPrincipalID = "anonymous"

AnonymousPrincipalID — зарезервированное слово, которым edge помечает запрос БЕЗ credential'а (`Principal{system, anonymous}`). Настоящего принципала с таким id не существует: реальные id — префиксованный crockford-base32 (`usr-…`/`sva-…`), их пространство с этим словом не пересекается.

Слово существует только как ярлык для аудита и логов. Оно НЕ является личностью, и любое сравнение личностей обязано это учитывать — см. Principal.IsAnonymous.

View Source
const ClockDrift = 5 * time.Minute

ClockDrift — слагаемое запаса на РАЗНИЦУ ИСТОЧНИКОВ ЧАСОВ.

`modified_at` пишется часами ПРОЦЕССА (`time.Now().UTC()` аргументом), а уборка судит часами БАЗЫ (`now()`) — чтобы источник был один на сервис, а не «столько источников, сколько реплик». Источники разные ⇒ разница входит в порог отдельным слагаемым, а не подразумевается нулём.

`pkg/tokenpolicy.RemovalSlack` здесь НЕ переиспользуется, и причина не в удобстве: та величина объявлена как запас на отставание ВНЕШНИХ читателей и связана размером с временем жизни токена через своего названного потребителя. Здесь закрывается пара «НАШ процесс против НАШЕЙ базы». Переиспользование сделало бы так, что сдвиг одной величины двигает и отсрочку снятия ключа подписи, и порог уборки в другом домене.

ВЕЛИЧИНА СЕГОДНЯ НЕ НЕСУЩАЯ, и это сказано числом, а не умолчано: порог удержания (7 суток) превосходит её на три порядка, поэтому сдвиг часов внутри кластера на минуты не меняет ни одного исхода. Несущей она станет, если `OperationRetention` когда-нибудь опустят к бюджету полла клиента (5 минут) — тогда сумма и станет тем, что сравнивают. Держит это TestOperationRetentionCoversTheDeclaredPollBudget: он утверждает ПРАВИЛО, а не сегодняшний запас.

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

OperationRetention — сколько терминальная операция остаётся доступной.

ВЕЛИЧИНА ВЫВЕДЕНА ИЗ ЧИТАТЕЛЯ, а не назначена числом. Читатель — оператор, разбирающий отказавшую мутацию: колонки `error_code`, `error_message`, `error_details` живут ТОЛЬКО здесь. Журнал аудита несёт актора и глагол мутации, но не текст отказа платформы, поэтому после снятия строки причина отказа не восстановима нигде. Предикат читателя — «строка ещё на месте, когда до неё дошли руки».

Интервал между отказавшей мутацией и разбором ограничен снизу нерабочими днями: отказ вечером пятницы разбирают утром понедельника. Порог короче двух суток делает такой разбор невозможным by construction; недельный запас покрывает его вместе с отпускной подменой и остаётся величиной, которую видно оператору, а не запасом, покрывающим произвольную ошибку в слагаемых.

ЧЕГО ВЕЛИЧИНА НЕ ОБЕЩАЕТ. Таблица операций не становится историей ресурса: история — предмет журнала аудита, чьё удержание объявлено намеренно. Списочное чтение после уборки отвечает «операции за последние OperationRetention», а не «все операции ресурса». Это смена наблюдаемого исхода, и она названа здесь, а не обнаруживается.

View Source
const SubjectOperations = "operations"

SubjectOperations — имя предмета уборки. Совпадает с именем таблицы: имя попадает в отчёт прохода и в метку метрики, и оператор обязан узнавать в нём таблицу.

Variables

View Source
var ErrAlreadyDone = errors.New("operation already completed")

ErrAlreadyDone возвращается из терминальных переходов, если строка уже завершена (CAS-on-`done` не совпал). Маппится в FAILED_PRECONDITION на Cancel-пути; на worker-пути трактуется как идемпотентный no-op.

View Source
var ErrNotFound = errors.New("operation not found")

ErrNotFound возвращается из Get, если операция не найдена.

View Source
var ErrWorkerStarted = errors.New("operations: worker already started; configure before Start")

ErrWorkerStarted — Configure вызван после старта dispatcher-loop. Опции (Recorder/Logger/MaxInflight/retry-budget) применимы только ДО запуска worker'а: после старта dispatcher читает поля Worker без блокировки, поэтому менять их под ним — data race.

Functions

func Active

func Active() int64

Active — pkg-level: число исполняемых workers default-registry.

func CheckRecordedOwnership

func CheckRecordedOwnership(caller, recorded Principal) error

CheckRecordedOwnership — вправе ли `caller` читать и отменять операцию, владелец которой прочитан из самой строки (`recorded`).

Возвращает nil, если вправе, и отказ полосы — если нет. Отказ ОДИН на все причины намеренно: он адресован арендатору, и различать в нём «твоя личность не извлеклась», «у строки нет владельца» и «владелец другой» значило бы рассказывать о чужой строке.

`recorded` с пустой парой — законный вход, а не ошибка вызывающего: строка может быть записана до появления учёта владельцев либо безымянным запросом. Такая строка НЕ читается арендатором (владелец неизвестен), и это решается здесь, а не у каждого вызывающего своим `if`.

func ConfigureDefault

func ConfigureDefault(opts ...WorkerOption) error

ConfigureDefault применяет опции к package-level default-registry ДО старта его dispatcher-loop. Composition root зовет ConfigureDefault(WithRecorder(promRec), WithLogger(logger)) на boot — это подключает live-worker метрики (terminal-write retries/failures, inflight) к Prometheus. Должен предшествовать Start()/первому Run; после старта — ErrWorkerStarted. Backward-compat: без вызова ConfigureDefault/ Start первый Run по-прежнему лениво стартует default-registry с NopRecorder.

func DefaultRetentionConfig

func DefaultRetentionConfig() retention.Config

DefaultRetentionConfig — величины ПЕТЛИ уборки таблицы операций.

Расписание общее для всех предметов платформы и объявлено один раз в `pkg/retention`; здесь оно только называется, а не переписывается числами.

func MetadataFor

func MetadataFor[T proto.Message](op *Operation) (T, error)

MetadataFor извлекает типизированные метаданные из операции. Возвращает ошибку, если Metadata nil или тип не совпадает.

func NewRetentionSweeper

func NewRetentionSweeper(s TerminalSweeper, cfg retention.Config, log *slog.Logger) (*retention.Sweeper, error)

NewRetentionSweeper собирает уборщика таблицы операций, НЕ поднимая петлю.

Отдельно от запуска, потому что «сборщик работает» обязано быть проверяемо без ожидания тикера: проба зовёт `Pass` методом и не спит.

func NotFoundStatus

func NotFoundStatus(id string) error

NotFoundStatus — ЕДИНСТВЕННЫЙ в дереве производитель отказа «нет такой операции».

Почему один, а не «два согласованных»

Текст — часть контракта (`api-conventions.md`), и он же несёт СОКРЫТИЕ: на адресе `/operations/{id}` этот 404 приходит из двух мест — от владельца, у которого строки нет или она не принадлежит вызывающему («есть, но не твоя» и «нет такой» намеренно неразличимы), и от края, которому известен префикс id, но не подключён его backend. Различие хоть в один байт отличает «нет доступа» от «не существует» и заодно называет, какие backend'ы край держит подключёнными (`security.md` §Hardening-инварианты, п. 6).

Две записи одного текста, сверяемые гейтом, — не то же самое, что одна. Сверка доказывает согласие в момент прогона; общий источник делает расхождение невыразимым. Здесь это не теория: производителей было два, и разошлись они регистром одной буквы — различие, которого не видно ни в обзоре изменения, ни в проверке, утверждающей код ответа (задача продукта #1370).

Вызывающие — обработчик владельца (`pkg/operations/operationspb`) и край (`gateway/internal/opsproxy`). Третьего быть не должно; держит это гейт `internal/repohygiene` `TestOperationNotFoundHasOneProducer`. Гейт судит ФОРМУ строкового выражения, а не его исходный текст, поэтому собрать тот же текст склейкой — не обход: перечень известных ему форм и границы этого перечня названы в шапке гейта.

От `ErrNotFound` отличается предметом: тот — внутреннее значение ошибки, по которому слой хранения узнают вызывающие внутри процесса, и id он не несёт by construction. Наружу уезжает только этот текст.

func Ready

func Ready() bool

Ready — pkg-level: живость dispatcher-loop default-registry.

func RetentionSubject

func RetentionSubject(s TerminalSweeper) retention.Subject

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

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

func Run

func Run(callerCtx context.Context, repo Repo, opID string, fn func(context.Context) (*anypb.Any, error))

Run — backward-compatible API: запускает worker в package-level default-registry.

callerCtx — request-context handler-а. Из него извлекаются observability-values (OTel SpanContext, request-id, slog logger) через baggage.Extract — они propagate'ятся в worker-ctx. worker НЕ наследует deadline/cancel callerCtx.

ВНИМАНИЕ (footgun — single dispatch/drain target): default-registry — это ГЛОБАЛЬНЫЙ Worker. Сервис ОБЯЗАН выбрать РОВНО ОДНУ цель dispatch/drain и придерживаться её на всех путях:

  • либо composition root владеет default-registry (ConfigureDefault → Start → Wait на shutdown) и ВЕСЬ код диспетчит через package-level Run;
  • либо composition root создаёт свой Worker (NewWorker) и ВЕСЬ код диспетчит через RunWithWorker(w, …), а на shutdown дренирует ИМЕННО его (w.Wait).

Смешивание (часть кода → Run на default-registry, а shutdown дренирует лишь свой NewWorker) приводит к тому, что in-flight операции на НЕ-дренируемом registry бросаются на перезапуске: терминальная запись (done=true) не происходит до тех пор, пока их не добёрет reconciler, и клиент, поллящий OperationService.Get, видит «застрявшую» операцию через рестарт. Предпочитайте явный NewWorker + RunWithWorker (DI из composition root).

func RunSync

func RunSync(ctx context.Context, repo Repo, op *Operation, fn func(context.Context) (*anypb.Any, error)) error

RunSync исполняет fn СИНХРОННО (op-in-response / statusless config-durable класс) и отражает терминальный результат в op — caller возвращает уже завершённую Operation (done=true) в том же ответе, без async-poll и follow-up GET. В отличие от Run (fire-and-trigger, worker в фоне), RunSync — для мутаций, чей предмет durable сразу после единственной writer-TX (config-INSERT, без саги).

Контракт зеркалит терминальную семантику async-worker'а (worker.execute):

  • fn success → repo.MarkDone(response); op.Done=true, op.Response=response.
  • fn error → маппится в терминальный google.rpc.Status (unwrapped/panic → фиксированный codes.Internal "internal worker error", raw driver-текст не течёт); repo.MarkError; op.Done=true, op.Error=st.Proto(). Бизнес-ошибка НЕ возвращается наружу — она живёт в Operation.result.error (durability = commit предмета мутации, ban #9; клиент читает исход из Operation).
  • panic в fn → recover → INTERNAL (как async-worker).

Sync-валидация (malformed id, формат, immutable, exactly-one placement) обязана произойти ДО op-create в Execute и вернуться там обычной gRPC-ошибкой (без Operation). RunSync отражает ТОЛЬКО исход worker-fn.

callerCtx применяется напрямую (в отличие от async-worker'а, где ctx detach'ится через baggage.Extract): клиент синхронно ждёт, поэтому его deadline/cancel уместны на fn. Сбой durable терминальной записи (сам MarkDone/MarkError вернул ошибку) → codes.Internal, чтобы Execute его вынес (строка op остаётся reconcilable). ErrAlreadyDone (строку уже разрешил Cancel/reconciler) — идемпотентный no-op.

func RunWithWorker

func RunWithWorker(w *Worker, callerCtx context.Context, repo Repo, opID string, fn func(context.Context) (*anypb.Any, error))

RunWithWorker — вариант с явным Worker registry (тесты / явное wiring). Семантика callerCtx — как у Run.

func Start

func Start()

Start явно запускает dispatcher-loop default-registry (идемпотентно). Composition root зовет Start() на boot ПОСЛЕ ConfigureDefault и подключает Ready() к readiness probe — тогда readiness lro-worker зеленый до первого трафика (нет boot-deadlock «NotReady → нет Run → worker не стартует → NotReady навсегда»).

func StartRetentionSweep

func StartRetentionSweep(ctx context.Context, s TerminalSweeper, cfg retention.Config, log *slog.Logger) (*retention.Sweeper, error)

StartRetentionSweep поднимает фоновую уборку таблицы операций.

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

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

func ValidateListPagination

func ValidateListPagination(filter ListFilter) error

ValidateListPagination проверяет page_size/page_token страницы операций по общей дисциплине List-RPC (api-conventions): page_size вне [0..Max] и мусорный курсор → InvalidArgument, а не «тихо поправим» и не INTERNAL.

Вынесена отдельно, потому что проверка обязана происходить ДО любого short-circuit'а выдачи (пустой ключ владения, пустой грант): иначе ошибка формата подменяется пустой страницей и клиент видит 200 там, где контракт обещает 400.

func Wait

func Wait(ctx context.Context) error

Wait — pkg-level: дренирует исполняемые workers default-registry.

func WithPrincipal

func WithPrincipal(ctx context.Context, p Principal) context.Context

WithPrincipal кладет Principal в context. Используется auth-interceptor'ом api-gateway: после валидации JWT и резолва subject через kaname interceptor вызывает WithPrincipal и пробрасывает ctx дальше в handler.

Без auth ctx остается пустым и PrincipalFromContext возвращает SystemPrincipal().

func WithoutPrincipal

func WithoutPrincipal(ctx context.Context) context.Context

WithoutPrincipal снимает principal с ctx (anonymous). После него PrincipalFromContextOK возвращает ok=false независимо от ранее установленного WithPrincipal. Используется transport-слоем для defense-in-depth-scrub'а на недоверенном peer'е.

Types

type ClaimRecorder

type ClaimRecorder interface {
	IncExecutionClaimRetries()
	IncExecutionClaimFailures()
	IncTerminalWriteAlreadyResolved(op string)
}

ClaimRecorder — ОПЦИОНАЛЬНОЕ расширение Recorder для наблюдаемости pre-execution подтверждения живости и поздней терминальной записи.

Почему расширение, а не новые методы базового Recorder: интерфейс реализуют шесть сервисных Prometheus-адаптеров, и дописывание в него ломает их сборку. Расширение подключается сервисом независимо; corelib-овые Nop/Mem реализуют его сразу, и это закреплено проверкой присваиванием ниже.

Канонические серии (имена — на стороне сервиса):

operations_execution_claim_retries_total   — повторы подтверждения живости
operations_execution_claim_failures_total  — отказ исполнить: живость не подтверждена
operations_terminal_write_already_resolved_total{op} — терминал не прошёл: строка уже разрешена

type FullRepo

type FullRepo interface {
	Repo
	OwnedOperationRepo
	MetadataFinalizer
	// TerminalSweeper — уборка терминальных строк по сроку (задача #1360).
	// Входит в FullRepo, а не берётся type-assert'ом: композиционный корень
	// каждого владельца обязан получить отказ КОМПИЛЯЦИИ, если уборщика у
	// репозитория нет, — иначе провязка молча выпадет у одного из восьми.
	TerminalSweeper
}

FullRepo — то, что действительно возвращает NewRepo: Repo плюс все опциональные апгрейды конкретной pgxpool-реализации. Метод-сет шире Repo, поэтому значение присваивается в любую переменную/поле/параметр типа Repo — расширение обратно совместимо, а use-case, которому апгрейд НУЖЕН, получает его проверку на компиляции, а не type-assert'ом в рантайме.

func NewRepo

func NewRepo(pool *pgxpool.Pool, schema string) FullRepo

NewRepo создает Repo для указанного пула и схемы. schema используется как квалификатор таблицы (schema.operations). Для схемы "public" передавайте "public".

Возвращает FullRepo (Repo + опциональные апгрейды реализации) — присваивается в любое место, ожидающее Repo; composition root, которому нужен апгрейд, получает его без type-assert'а.

type ListFilter

type ListFilter struct {
	ResourceID string // если непуст — фильтр по resource_id (денормализованное поле)
	AccountID  string // если непуст — фильтр по account_id (денормализованное поле, partial cursor-индекс)
	PageSize   int64
	PageToken  string
}

ListFilter — параметры фильтрации/пагинации для List.

type MemRecorder

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

MemRecorder — in-memory Recorder для тестов и как безопасный дефолт. Concurrency-safe.

func NewMemRecorder

func NewMemRecorder() *MemRecorder

NewMemRecorder — пустой in-memory Recorder.

func (*MemRecorder) ExecutionClaimFailures

func (m *MemRecorder) ExecutionClaimFailures() float64

ExecutionClaimFailures возвращает накопленные отказы исполнить работу без подтверждённой живости.

func (*MemRecorder) ExecutionClaimRetries

func (m *MemRecorder) ExecutionClaimRetries() float64

ExecutionClaimRetries возвращает накопленные повторы подтверждения живости.

func (*MemRecorder) IncExecutionClaimFailures

func (m *MemRecorder) IncExecutionClaimFailures()

IncExecutionClaimFailures инкрементит счётчик отказов исполнить работу из-за неподтверждённой живости операции.

func (*MemRecorder) IncExecutionClaimRetries

func (m *MemRecorder) IncExecutionClaimRetries()

IncExecutionClaimRetries инкрементит счётчик повторов подтверждения живости.

func (*MemRecorder) IncOrphansRecovered

func (m *MemRecorder) IncOrphansRecovered(outcome string)

IncOrphansRecovered инкрементит счетчик разрешенных reconciler'ом orphan'ов.

func (*MemRecorder) IncReconcileErrors

func (m *MemRecorder) IncReconcileErrors()

IncReconcileErrors инкрементит счетчик ошибок sweep-цикла.

func (*MemRecorder) IncReconcileRuns

func (m *MemRecorder) IncReconcileRuns()

IncReconcileRuns инкрементит счетчик прогонов sweep-цикла.

func (*MemRecorder) IncTerminalWriteAlreadyResolved

func (m *MemRecorder) IncTerminalWriteAlreadyResolved(op string)

IncTerminalWriteAlreadyResolved инкрементит счётчик терминальных записей, не прошедших сравнение-и-замену, потому что строка уже разрешена.

func (*MemRecorder) IncTerminalWriteFailures

func (m *MemRecorder) IncTerminalWriteFailures(op string)

IncTerminalWriteFailures инкрементит счетчик невосстановимых терминальных записей.

func (*MemRecorder) IncTerminalWriteRetries

func (m *MemRecorder) IncTerminalWriteRetries(op string)

IncTerminalWriteRetries инкрементит счетчик ретраев терминальной записи.

func (*MemRecorder) Inflight

func (m *MemRecorder) Inflight() float64

Inflight возвращает текущее значение gauge.

func (*MemRecorder) MaxInflight

func (m *MemRecorder) MaxInflight() float64

MaxInflight возвращает наблюденный пик inflight.

func (*MemRecorder) OrphansRecovered

func (m *MemRecorder) OrphansRecovered(outcome string) float64

OrphansRecovered возвращает счетчик разрешенных orphan'ов по outcome-лейблу.

func (*MemRecorder) ReconcileErrors

func (m *MemRecorder) ReconcileErrors() float64

ReconcileErrors возвращает счетчик ошибок sweep-цикла.

func (*MemRecorder) ReconcileRuns

func (m *MemRecorder) ReconcileRuns() float64

ReconcileRuns возвращает счетчик прогонов sweep-цикла.

func (*MemRecorder) SetInflight

func (m *MemRecorder) SetInflight(n float64)

SetInflight выставляет gauge числа запущенных worker'ов и трекает пик.

func (*MemRecorder) TerminalWriteAlreadyResolved

func (m *MemRecorder) TerminalWriteAlreadyResolved(op string) float64

TerminalWriteAlreadyResolved возвращает счётчик поздних терминальных записей по op-лейблу.

func (*MemRecorder) TerminalWriteFailures

func (m *MemRecorder) TerminalWriteFailures(op string) float64

TerminalWriteFailures возвращает накопленные failures по op-лейблу.

func (*MemRecorder) TerminalWriteRetries

func (m *MemRecorder) TerminalWriteRetries(op string) float64

TerminalWriteRetries возвращает накопленные ретраи по op-лейблу.

type MetadataFinalizer

type MetadataFinalizer interface {
	// MarkDoneWithMetadata переводит операцию в done=true, ЗАМЕНЯЯ metadata и
	// записывая response. Идемпотентна и durable — тот же CAS-on-`done`, что у
	// MarkDone: уже терминальная строка не перезаписывается (ErrAlreadyDone),
	// отсутствующая → ErrNotFound.
	MarkDoneWithMetadata(ctx context.Context, id string, metadata, response *anypb.Any) error
}

MetadataFinalizer — терминальная запись успеха, которая ВМЕСТЕ с response заменяет metadata операции.

Кому нужно: синхронной мутации, чей owning-ресурс становится известен только ПОСЛЕ записи — id строки выдаёт БД, а на идемпотентном повторе это id уже существующей строки, поэтому «сгенерировать заранее и вписать в Create» не работает by construction. Строку операции при этом ОБЯЗАНО создавать ДО мутации (иначе сбой Create оставляет закоммиченную мутацию без pollable операции — Get(id) отвечает NotFound навсегда), то есть на Create метаданные заведомо неполны. Без этого перехода объявленное контрактом поле метаданных не заполняется никогда.

Денормализованные индекс-колонки (resource_id / account_id) выставляются на Create и здесь НЕ пересчитываются: финальная metadata не обязана нести их поля, а перевычисление молча обнулило бы уже выставленный ключ.

type NopRecorder

type NopRecorder struct{}

NopRecorder — no-op Recorder, безопасный дефолт когда метрики не подключены.

func (NopRecorder) IncExecutionClaimFailures

func (NopRecorder) IncExecutionClaimFailures()

IncExecutionClaimFailures — no-op.

func (NopRecorder) IncExecutionClaimRetries

func (NopRecorder) IncExecutionClaimRetries()

IncExecutionClaimRetries — no-op.

func (NopRecorder) IncOrphansRecovered

func (NopRecorder) IncOrphansRecovered(string)

IncOrphansRecovered — no-op.

func (NopRecorder) IncReconcileErrors

func (NopRecorder) IncReconcileErrors()

IncReconcileErrors — no-op.

func (NopRecorder) IncReconcileRuns

func (NopRecorder) IncReconcileRuns()

IncReconcileRuns — no-op.

func (NopRecorder) IncTerminalWriteAlreadyResolved

func (NopRecorder) IncTerminalWriteAlreadyResolved(string)

IncTerminalWriteAlreadyResolved — no-op.

func (NopRecorder) IncTerminalWriteFailures

func (NopRecorder) IncTerminalWriteFailures(string)

IncTerminalWriteFailures — no-op.

func (NopRecorder) IncTerminalWriteRetries

func (NopRecorder) IncTerminalWriteRetries(string)

IncTerminalWriteRetries — no-op.

func (NopRecorder) SetInflight

func (NopRecorder) SetInflight(float64)

SetInflight — no-op.

type Operation

type Operation struct {
	ID          string
	Description string
	CreatedAt   time.Time
	CreatedBy   string
	ModifiedAt  time.Time
	Done        bool
	Metadata    *anypb.Any     // специфичные метаданные (CreateInstanceMetadata и т.д.)
	Error       *status.Status // заполнен если done && ошибка
	Response    *anypb.Any     // заполнен если done && успех — финальное состояние ресурса

	// ResourceID — id ресурса-владельца операции для денормализованного индекса
	// operations.resource_id (фильтр List(ListFilter{ResourceID})). Если задан
	// use-case'ом явно — используется как есть; если пуст — repo падает на
	// reflection-fallback (первое `*_id`-поле Metadata). Явное значение НАДЁЖНЕЕ:
	// reflection-угадывание «первое _id == owning resource» ошибётся, если в
	// *Metadata первым объявлено не-owning поле (project_id/parent_id/…). Ставьте
	// его при конструировании операции.
	ResourceID string

	// Principal — кто инициировал операцию (kaname-resolved). Без auth
	// заполняется SystemPrincipal(); при наличии auth-ctx — из него через
	// PrincipalFromContext.
	Principal Principal
}

Operation — domain-тип, зеркалит proto Operation. Используется repo / service-слоями внутри Kachō-сервисов.

func ListForCaller

func ListForCaller(ctx context.Context, repo Repo, filter ListFilter) ([]Operation, string, error)

ListForCaller — единственная точка, через которую tenant-facing список операций попадает в сервисы: ключ владения выводится из ctx, предикат владения уходит внутрь SQL WHERE (ListOwned), несуженный Repo.List не вызывается никогда.

Зачем отдельная функция, а не «зовите ListOwned сами». Правильный вызов складывается из трёх шагов, и пропуск любого возвращает исходную дыру: получить ownership-апгрейд репозитория, вывести ключ владения именно из ctx (а не из PrincipalFromContext — тот на пустом контексте отдаёт системную личность, которая владеет каждой системно записанной строкой), и на отсутствие ключа отказать, а не продолжить. Собранные в одном месте, эти шаги проверяются одним гейтом; рассыпанные по пятнадцати вызывающим — пятнадцатью прочтениями.

Исходы ровно три:

  • ключ владения есть → страница своих операций (полная проекция: строка принадлежит вызывающему, поэтому и содержимое, и личность на ней — его собственные);
  • ключа нет (контекст без принципала, принципал снят недоверенным форвардером, именованная анонимность) → ПУСТАЯ страница. Не несуженная выдача и не «своих нет, значит покажем чужие»;
  • репозиторий не несёт ownership-апгрейда (ошибка провязки) → отказ с фиксированным текстом. Откат на несуженный путь запрещён: молчаливый откат — это ровно тот исход, ради недопущения которого функция и написана.

Формат страницы проверяется ДО схлопывания в пустую: иначе мусорный курсор у неопознанного вызывающего вернул бы пустую страницу вместо отказа по формату (api-conventions: format-validate → authz → repo).

Несуженный Repo.List остаётся законным для доверенного внутреннего яруса (кластерный аудит, реконсайлер) — там вызывающий уже авторизован иначе, и именно поэтому запрет выражен гейтом с поимённым списком, а не удалением метода.

func New

func New(domainPrefix, description string, metadata proto.Message) (Operation, error)

New создает Operation с 20-char id (3-char prefix + crockford-base32), текущим временем, done=false. domainPrefix — 3-символьный префикс из ids.PrefixOperationRM / ids.PrefixOperationVPC, по которому api-gateway opsproxy маршрутизирует Operation.Get/Cancel в нужный backend.

Если prefix пуст — id без prefix (legacy/internal use); такая операция не маршрутизируется через opsproxy и доступна только локально внутри сервиса.

metadata — proto-сообщение специфичное для типа RPC (например, CreateInstanceMetadata{instance_id: uid}).

func NewFromContext

func NewFromContext(ctx context.Context, domainPrefix, description string, metadata proto.Message) (Operation, error)

NewFromContext — то же что New, но также заполняет op.Principal из ctx (через PrincipalFromContextOK). Удобный shortcut для use-case'ов, которые получают Principal в ctx от auth-interceptor api-gateway.

Без ctx-Principal (нет auth-interceptor'а, ИЛИ ctx был явно scrub'нут через WithoutPrincipal — defense-in-depth снятие forwarded-principal на недоверенном peer'е) — op.Principal остается zero (Create в repo сделает fallback к SystemPrincipal). С явно установленным (и не-scrub'нутым) ctx-Principal — он переносится в op.Principal, и Create / CreateWithPrincipal будут использовать его как источник правды.

БЕЗЫМЯННЫЙ ctx-Principal (Principal.IsAnonymous — пустая пара либо ярлык анонима) переносом НЕ считается: имя, общее для всех безымянных запросов, стало бы общим ключом ВЛАДЕНИЯ на записанных строках. Такую строку правильнее не создавать, чем рассчитывать, что читатель её отсеет; владелец у неё остаётся неизвестен (SystemPrincipal-fallback в Create), и tenant-чтение fail-closed по общему правилу.

type OwnedOperationRepo

type OwnedOperationRepo interface {
	// GetOwned возвращает операцию по id ТОЛЬКО если она принадлежит owner.
	// 0 строк (нет такой ИЛИ не владелец) → ErrNotFound.
	GetOwned(ctx context.Context, id string, owner Owner) (*Operation, error)
	// CancelOwned атомарно отменяет операцию owner'а (CAS WHERE done=false AND
	// ownership-предикат, RETURNING терминальное состояние без reload-Get).
	// Идемпотентно на уже-CANCELLED (→ OK с тем же Operation); на терминале
	// SUCCESS/ERROR → ErrAlreadyDone; чужая/нет → ErrNotFound.
	CancelOwned(ctx context.Context, id string, owner Owner) (*Operation, error)
	// ListOwned возвращает страницу операций owner'а: ownership-предикат внутри
	// SQL WHERE (симметрично GetOwned) комбинируется AND-ом с фильтрами
	// ListFilter (ResourceID/AccountID/keyset-пагинация). Чужие операции не
	// попадают в выдачу — tenant-facing OperationService.List обязан идти этим
	// путём, а не unscoped Repo.List. Возвращает (страница, next_page_token, err).
	ListOwned(ctx context.Context, filter ListFilter, owner Owner) ([]Operation, string, error)
}

OwnedOperationRepo — узкий ownership-scoped порт чтения/отмены операции. Ownership-предикат — внутри SQL WHERE / атомарного CAS (within-service инвариант на DB-уровне; software-side Get→check исключен). Чужой owner → ErrNotFound (no-leak, неотличимо от «нет такой»).

Это ОТДЕЛЬНЫЙ от Repo интерфейс: новые методы не добавляются в общий operations.Repo, чтобы не ломать его mock-реализации в других сервисах (iam/compute/nlb/geo). Реализуется конкретным pgRepo; consumer получает его через AsOwned(repo).

func AsOwned

func AsOwned(r Repo) (OwnedOperationRepo, bool)

AsOwned проверяет, что Repo поддерживает ownership-scoped доступ, и возвращает узкий порт. Для pgRepo (NewRepo) — true. Composition root vpc-handler'а:

owned, ok := operations.AsOwned(opsRepo)
if !ok { /* fail-fast при wiring'е */ }

type Owner

type Owner struct {
	PrincipalType string
	PrincipalID   string
	AccountID     string
}

Owner — ключ владельца операции для ownership-scoped доступа (Get/Cancel). Operative-ключ — пара (PrincipalType, PrincipalID) создателя операции (не project_id): видимость операции creator-principal-scoped. AccountID — дополнительный IAM-only ключ: gate ДОПОЛНИТЕЛЬНО матчит по account там, где колонка account_id NOT NULL; для сервисов без account-metadata (vpc/compute/ nlb) AccountID == "" и эта ветка инертна.

func OwnerFromContext

func OwnerFromContext(ctx context.Context) (Owner, bool)

OwnerFromContext выводит owner-ключ из ctx и отдельно сообщает, есть ли он.

ok=false означает «принципала не было» (ctx без auth-интерсептора, принципал явно снят transport-слоем на недоверенном форвардере), «принципал без id» ЛИБО «принципал — именованная анонимность» (`{system, anonymous}`, которую edge выдаёт запросу без credential'а). В этом случае возвращается НУЛЕВОЙ Owner, а не ключ: личность по умолчанию не может быть владельцем чего-либо, и «неизвестно кто» — тоже. Вызывающий обязан на ok=false отказать (для tenant-facing чтения — тем же no-leak NotFound, что и «чужая операция»), а не продолжать с полученным ключом.

Именованная анонимность отсекается здесь по той же причине, что и пустая: пара «система/аноним» непуста, поэтому как ключ владения она выглядит настоящей и совпадает сама с собой — любые два безымянных запроса делили бы один ключ, то есть один читал бы и отменял операции другого.

ЯВНО установленный системный принципал (WithPrincipal(SystemPrincipal()) на доверенном internal-пути) даёт ok=true и остаётся владельцем своих операций — сужается именно анонимность, а не bootstrap-личность как таковая. Ровно то же различение делает authz-слой (pkg/authz defaultSubjectExtractor).

func OwnerFromPrincipal

func OwnerFromPrincipal(p Principal) Owner

OwnerFromPrincipal строит Owner из УЖЕ УСТАНОВЛЕННОГО Principal'а. AccountID не заполняется (principal-ключ); IAM-сервис при необходимости выставляет его сам.

ВНИМАНИЕ: не скармливай сюда PrincipalFromContext(ctx). Тот на пустом (и на явно очищенном) ctx отдаёт SystemPrincipal() — backward-compat fallback для ЗАПИСИ операций фоновыми путями, а не удостоверение личности. Как owner-ключ для ЧТЕНИЯ он совпадает с ownership-предикатом на каждой операции, записанной системным принципалом, то есть анонимный вызывающий получает чужие операции. Для вывода owner-ключа из ctx есть OwnerFromContext.

func (Owner) IsAnonymous

func (o Owner) IsAnonymous() bool

IsAnonymous — principal-ключ никого не называет. Такой ключ НИКОГДА не должен доезжать до ownership-предиката, и по двум разным причинам:

  • пустые компоненты матчатся со строками, у которых колонки принципала пусты;
  • именованная анонимность (`{system, anonymous}`) матчится со строками, записанными ЛЮБЫМ другим безымянным запросом — ключ у них общий по построению, их столько же, сколько имён, то есть одно на всех.

Условие делегировано Principal.IsAnonymous — единый предикат «этот принципал не называет никого» на весь дом.

Про AccountID: у предиката есть вторая, account-ветка, но её ключ сейчас НИКТО не заполняет — во всём дереве Owner строится исключительно из пары (PrincipalType, PrincipalID) через OwnerFromPrincipal/OwnerFromContext, а AccountID остаётся "" и ветка инертна. Поэтому проверка намеренно смотрит только на principal-пару. Если появится account-only владелец (Owner с пустым принципалом и непустым AccountID), этот предикат обязан быть расширен ВМЕСТЕ с ним — иначе такой владелец будет молча отвергнут; правку делать осознанно, а не «ослабить проверку, чтобы прошло».

type Principal

type Principal struct {
	Type        string
	ID          string
	DisplayName string
}

Principal — кто инициировал операцию. Заполняется auth-interceptor'ом api-gateway. Без auth — stub {Type: "system", ID: "bootstrap", DisplayName: "System"}.

Семантика полей:

  • Type: "user" | "service_account" | "system"
  • ID: "usr-..." | "sva-..." | "bootstrap"
  • DisplayName: human-readable email / SA-name / "System"

Каждая запись в operations-таблице хранит эти поля в колонках principal_type / principal_id / principal_display_name (миграция migrations/common/0002_operations_principal.sql).

func PrincipalFromContext

func PrincipalFromContext(ctx context.Context) Principal

PrincipalFromContext извлекает Principal из ctx. Если ctx пустой (нет auth-interceptor'а, фоновый job, тест) — возвращает SystemPrincipal().

Use-case вызывает PrincipalFromContext в начале обработки мутации и передает результат в repo.CreateWithPrincipal:

p := operations.PrincipalFromContext(ctx)
if err := repo.CreateWithPrincipal(ctx, op, p); err != nil { ... }

func PrincipalFromContextOK

func PrincipalFromContextOK(ctx context.Context) (Principal, bool)

PrincipalFromContextOK извлекает Principal из ctx и отдельно сообщает, был ли он явно установлен (WithPrincipal). При ok=false ctx не нес Principal'а (anonymous / нет auth-interceptor'а) и возвращается SystemPrincipal()-fallback.

Этот вариант нужен authz-слою, чтобы отличить АНОНИМНЫЙ запрос (ctx без Principal'а) от ЯВНО установленного system-principal: первый обязан fail-closed, второй может быть пропущен опцией AllowSystemPrincipal. PrincipalFromContext такого различения не дает (оба → SystemPrincipal()).

func SystemPrincipal

func SystemPrincipal() Principal

SystemPrincipal — stub principal для control-plane bootstrap'а / системных операций (миграции, фоновый retry, тесты без auth). Также используется как дефолт в CreateWithPrincipal-обертке legacy-Create и в PrincipalFromContext для пустого ctx.

func (Principal) IsAnonymous

func (p Principal) IsAnonymous() bool

IsAnonymous сообщает, что принципал не называет НИКОГО.

Таких случая два, и они обязаны трактоваться одинаково:

  • пара пуста (принципала не было вовсе);
  • пара несёт зарезервированное слово AnonymousPrincipalID — «неизвестно кто», выданное запросу без credential'а.

Второй случай — источник целого класса дефектов: как только у «неизвестно кто» появляется ИМЯ, любое сравнение с этим именем даёт истину. Два независимых безымянных запроса получают один и тот же ключ и оказываются друг для друга «своими»; гейт вида «личность извлеклась ⇒ аутентифицирован» на нём срабатывает. Поэтому именованная анонимность обязана вести себя как отсутствие принципала везде, где принимается решение о личности.

Тип при этом не проверяется намеренно: слово зарезервировано целиком, и «неизвестно кто» с любым заявленным типом остаётся неизвестно кем.

НЕ включает ЯВНО установленный системный принципал (`{system, bootstrap}`): это настоящая bootstrap-личность доверенного internal-пути, она остаётся владельцем своих операций. Сужается именно анонимность.

func (Principal) IsSystem

func (p Principal) IsSystem() bool

IsSystem сообщает, что принципал — ЯВНАЯ bootstrap-личность внутреннего пути (`{system, bootstrap}`, см. SystemPrincipal).

Пара сверяется с SystemPrincipal(), а не с двумя написанными здесь словами: иначе имя внутренней личности объявлялось бы в дереве второй раз, и второе объявление разошлось бы с первым молча. Именно так оно и было — край перечислял эту пару голыми строками дважды, в двух соседних условиях.

НЕ то же, что IsAnonymous, и путать их нельзя. Аноним не называет НИКОГО; bootstrap-личность называет вполне определённого внутреннего вызывающего и остаётся владельцем своих операций. Различение делает и authz-слой, и вывод ключа владения (OwnerFromContext).

DisplayName не сверяется намеренно: он косметический, приходит из разных мест и решением о личности быть не может.

type Reconciler

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

Reconciler — startup + периодический backstop: разрешает осиротевшие in-flight операции (умершего процесса) в терминал по committed-реальности ресурса. Claim через FOR UPDATE SKIP LOCKED — конкурентные reconciler'ы реплик партиционируют множество (exactly-once), а не дерутся. Терминальная запись — через тот же CAS-on-`done` (идемпотентна с live-worker'ом).

func NewReconciler

func NewReconciler(pool *pgxpool.Pool, resolver Resolver, cfg ReconcilerConfig, opts ...ReconcilerOption) *Reconciler

NewReconciler конструирует Reconciler. pool/resolver обязательны.

func (*Reconciler) RecoverAll

func (rc *Reconciler) RecoverAll(ctx context.Context) error

RecoverAll прогоняет Sweep до тех пор, пока очередной прогон не разрешит 0 операций (backlog осиротевших исчерпан / остались только нерезолвимые в этом проходе). Зовется на старте ДО приема трафика.

func (*Reconciler) Run

func (rc *Reconciler) Run(ctx context.Context)

Run — периодический backstop: Sweep на каждом тике до отмены ctx. Ошибки отдельного прогона логируются (loop не умирает на transient-сбое).

РЕПЛИКИ: клейм — проход клеймит строки `FOR UPDATE SKIP LOCKED` в собственной транзакции, поэтому конкурирующим репликам достаются непересекающиеся партии.

func (*Reconciler) Sweep

func (rc *Reconciler) Sweep(ctx context.Context) (int, error)

Sweep — один прогон: claim пачки осиротевших операций (FOR UPDATE SKIP LOCKED) и разрешение каждой через Resolver в терминал. Терминальная запись — внутри той же транзакции, что держит claim-lock (exactly-once между репликами). Возвращает число разрешенных.

type ReconcilerConfig

type ReconcilerConfig struct {
	// Schema — schema-квалификатор таблицы operations ("public" / "kacho_vpc").
	Schema string
	// OrphanGrace — минимальный возраст (по modified_at) кандидата-orphan'а.
	// Должен превышать максимальную ожидаемую длительность операции, чтобы не
	// разрешать преждевременно еще-живого worker'а. Дефолт 5m.
	OrphanGrace time.Duration
	// BatchSize — размер пачки claim'а за один sweep. Дефолт 100.
	BatchSize int
	// Interval — период периодического sweep'а (Run). Дефолт 30s.
	Interval time.Duration
	// ResolveTimeout — жёсткий потолок времени одного Resolver.Resolve внутри
	// claim-транзакции. Sweep держит claim-tx (FOR UPDATE SKIP LOCKED + пул-коннект)
	// открытой на всё время резолва пачки; без потолка зависший доменный Resolve
	// (напр. peer-вызов без deadline) оставил бы tx idle-in-transaction на
	// неограниченное время — блокируя VACUUM и удерживая коннект. Потолок гарантирует
	// прогресс sweep'а. Дефолт 10s; ≤0 → default.
	ResolveTimeout time.Duration
	// SweepBudget — жёсткий потолок СУММАРНОЙ длительности claim-транзакции одного
	// Sweep'а. ResolveTimeout ограничивает лишь ОДИН резолв; без агрегатного потолка
	// claim-tx могла бы жить до BatchSize×ResolveTimeout (~1000s при outage), удерживая
	// пул-коннект, FOR UPDATE row-locks и xmin-горизонт operations-таблицы против
	// VACUUM на всё окно. При исчерпании бюджета Sweep коммитит
	// уже разрешённое и выходит; неразрешённые orphan'ы остаются durable (done=false)
	// и добираются следующим Sweep'ом. Дефолт 30s; ≤0 → default.
	SweepBudget time.Duration
}

ReconcilerConfig параметризует Reconciler.

type ReconcilerOption

type ReconcilerOption func(*Reconciler)

ReconcilerOption — функциональная опция Reconciler.

func WithReconcilerLogger

func WithReconcilerLogger(l *slog.Logger) ReconcilerOption

WithReconcilerLogger подключает структурированный логгер.

func WithReconcilerRecorder

func WithReconcilerRecorder(r Recorder) ReconcilerOption

WithReconcilerRecorder подключает sink метрик reconcile (runs/errors/orphans).

type Recorder

type Recorder interface {
	IncTerminalWriteRetries(op string)
	IncTerminalWriteFailures(op string)
	SetInflight(n float64)
	IncOrphansRecovered(outcome string)
	IncReconcileRuns()
	IncReconcileErrors()
}

Recorder — sink метрик durability-слоя LRO. Конкретный Prometheus-клиент сюда НЕ импортируется (corelib остается dependency-light): сервис подключает Prometheus-backed Recorder в composition root, тесты используют MemRecorder. Канонические серии (имена — на стороне сервиса, лейблы фиксированы здесь):

operations_terminal_write_retries_total{op}   — ретраи MarkDone/MarkError на transient-сбое
operations_terminal_write_failures_total{op}  — терминальная запись не удалась после max-ретраев
operations_inflight                           — gauge числа запущенных worker'ов (<= max)
operations_orphans_recovered_total{outcome}   — orphan'ы, разрешенные reconciler'ом (done|error)
operations_reconcile_runs_total               — прогоны sweep-цикла
operations_reconcile_errors_total             — ошибки sweep-цикла

op-лейбл — "MarkDone"/"MarkError"; outcome-лейбл — "done"/"error".

type Repo

type Repo interface {
	// Create сохраняет новую операцию (done=false). Principal записывается
	// как SystemPrincipal — backward-compat shim. Use-case с auth-ctx должен
	// вызывать CreateWithPrincipal (или передавать op.Principal заранее).
	Create(ctx context.Context, op Operation) error
	// CreateWithPrincipal сохраняет новую операцию с явно переданным
	// principal'ом. Вызывается из use-case'ов сервисов после
	// PrincipalFromContext(ctx).
	CreateWithPrincipal(ctx context.Context, op Operation, p Principal) error
	// Get возвращает операцию по ID. Возвращает ErrNotFound если операции нет.
	//
	// БЕЗОПАСНОСТЬ (IDOR, CWE-639): Get — UNSCOPED, без ownership-предиката. Строка
	// Operation несёт metadata_data/response_data (сериализованный ресурс, resource_id,
	// account_id). Для tenant-facing OperationService.Get НЕЛЬЗЯ отдавать результат
	// Get напрямую — id-format публичен и перечислим, иначе tenant A прочитает
	// операцию tenant'а B. Используйте ownership-scoped OwnedOperationRepo.GetOwned
	// (конкретный pgRepo его реализует). Unscoped Get — только для доверенных
	// internal-вызовов (reconciler, worker), уже авторизованных иначе.
	Get(ctx context.Context, id string) (*Operation, error)
	// List возвращает список операций с постраничной навигацией.
	List(ctx context.Context, filter ListFilter) ([]Operation, string, error)
	// MarkDone переводит операцию в done=true, записывает финальный ресурс (response).
	//
	// ВНИМАНИЕ: третий параметр — RESPONSE, не metadata. metadata остаётся
	// такой, какой её вписал Create; этот переход её не трогает. Если
	// owning-ресурс становится известен только ПОСЛЕ мутации (id приходит из
	// БД), объявленное поле метаданных так и останется пустым — для этого
	// случая есть MarkDoneWithMetadata (см. MetadataFinalizer).
	MarkDone(ctx context.Context, id string, response *anypb.Any) error
	// MarkError переводит операцию в done=true, записывает ошибку (google.rpc.Status).
	MarkError(ctx context.Context, id string, err *status.Status) error
	// Cancel переводит операцию в done=true со статусом CANCELLED.
	//
	// БЕЗОПАСНОСТЬ (IDOR, CWE-639): Cancel — UNSCOPED, без ownership-предиката
	// (UPDATE ... WHERE id=$1 AND done=false). id-format публичен и перечислим,
	// поэтому для tenant-facing OperationService.Cancel НЕЛЬЗЯ вызывать Cancel
	// напрямую — иначе principal B отменит in-flight операцию principal'а A
	// (cross-owner LRO corruption). Используйте ownership-scoped
	// OwnedOperationRepo.CancelOwned (конкретный pgRepo его реализует). Unscoped
	// Cancel — только для доверенных internal-вызовов (reconciler, worker),
	// уже авторизованных иначе. Симметрично IDOR-warning на Get выше.
	Cancel(ctx context.Context, id string) error
}

Repo — интерфейс для хранения и обновления Operations. Каждый сервис создает таблицу operations через migration (см. migrations/common/0001_operations.sql + 0002_operations_principal.sql).

type Resolver

type Resolver interface {
	Resolve(ctx context.Context, op Operation) (ResolverResult, error)
}

Resolver — доменный порт (реализуется сервисом): по метаданным осиротевшей операции определяет ее терминальный исход, сверяясь с committed-реальностью ресурса (repo.Get по resource_id из metadata). Движок reconciler'а — в corelib, resolver — в сервисе (знает типы метаданных и таблицы ресурсов).

Контракт диспетчеризации (vpc-style):

  • Create/Update-метаданные: ресурс присутствует → {OutcomeDone, current}, отсутствует → {OutcomeInterrupted}.
  • Delete-метаданные: ресурс отсутствует → {OutcomeDone, nil(Empty)}, присутствует → {OutcomeInterrupted}.
  • transient-ошибка чтения ресурса → возврат (ResolverResult{}, err): движок инкрементит reconcile_errors и пропускает orphan до следующего sweep'а.

type ResolverResult

type ResolverResult struct {
	Outcome  TerminalOutcome
	Response *anypb.Any // используется при OutcomeDone (nil → google.protobuf.Empty-семантика)
}

ResolverResult — решение Resolver'а по одной осиротевшей операции.

type TerminalOutcome

type TerminalOutcome int

TerminalOutcome — терминальный исход, который доменный Resolver вычислил по committed-реальности ресурса.

const (
	// OutcomeSkip — committed-реальность не позволяет уверенно разрешить операцию
	// в этом прогоне (resolver не смог прочитать ресурс / неоднозначно); строка
	// остается done=false, sweep повторится позже.
	OutcomeSkip TerminalOutcome = iota
	// OutcomeDone — работа фактически закоммичена → MarkDone(Response). Для Create/
	// Update Response — текущий ресурс; для Delete Response может быть nil (Empty).
	OutcomeDone
	// OutcomeInterrupted — работа не дошла до commit (ресурс отсутствует для
	// Create/Update либо еще жив для Delete) → MarkError(interrupted).
	OutcomeInterrupted
)

type TerminalSweeper

type TerminalSweeper interface {
	SweepTerminal(ctx context.Context, grace time.Duration, batch int) (int64, bool, error)
}

TerminalSweeper — порт уборщика таблицы операций.

Объявлен здесь, а не у вызывающего: реестр уборки собирается этим же пакетом (RetentionSubject), и порт нужен ему, а не владельцу композиционного корня.

type TerminalWriteConfig

type TerminalWriteConfig struct {
	InitialInterval time.Duration
	MaxInterval     time.Duration
	MaxElapsed      time.Duration // общий budget на ретраи; после — failure-метрика, done=false
}

TerminalWriteConfig параметризует durable retry-loop терминальной записи.

type Worker

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

Worker — bounded координатор async worker-горутин.

Назначение: graceful-shutdown сервиса не должен терять in-flight операции, а burst мутаций — не порождать unbounded-горутины. Dispatcher-loop разбирает in-memory backlog, ограничивая исполняемые worker'ы семафором max-inflight; терминальная запись durable (retry+CAS); readiness отражает живость loop'а.

func NewWorker

func NewWorker(opts ...WorkerOption) *Worker

NewWorker — новый изолированный Worker. Опции настраивают max-inflight, метрики, логгер, retry-budget. Без опций — разумные дефолты (max-inflight 64). Dispatcher-loop стартует лениво на первом Run или явно через Start().

func (*Worker) Active

func (w *Worker) Active() int64

Active — текущее число исполняемых worker'ов (для observability / drain).

func (*Worker) Configure

func (w *Worker) Configure(opts ...WorkerOption) error

Configure применяет опции к еще-не-стартованному Worker. Composition root зовет Configure ДО Start, чтобы подключить Prometheus-Recorder и логгер к worker'у (default-registry создается с NopRecorder — без Configure live-worker метрики мертвы). После Start — ErrWorkerStarted. Потокобезопасно: проверка «не стартован» и мутация полей идут под configMu, сериализуясь с ensureStarted.

func (*Worker) Ready

func (w *Worker) Ready() bool

Ready отражает живость dispatcher-loop. NotReady, пока loop не запущен или остановлен — deployment не должен слать мутирующий трафик на под без живого dispatcher'а.

func (*Worker) Start

func (w *Worker) Start()

Start явно запускает dispatcher-loop (идемпотентно). Composition root зовет Start() на boot и подключает Ready() к readiness probe. Для тестов/back-compat первый Run стартует loop сам.

func (*Worker) Stop

func (w *Worker) Stop()

Stop останавливает dispatcher-loop (идемпотентно): новые задачи перестают диспетчеризоваться, Ready() → false. Уже исполняемые worker'ы дренируются отдельно через Wait(). Backlog не-стартовавших задач durable → reconciler.

func (*Worker) Wait

func (w *Worker) Wait(ctx context.Context) error

Wait блокируется пока вся принятая работа (backlog + исполняемые) не завершится, либо ctx истечет. Возвращает nil при штатном drain'е, ctx.Err() при таймауте. Задачи, не уложившиеся в окно drain'а, остаются durable (done=false) → добираются reconciler'ом на следующем старте.

type WorkerOption

type WorkerOption func(*workerConfig)

WorkerOption — функциональная опция конфигурации Worker.

func WithLogger

func WithLogger(l *slog.Logger) WorkerOption

WithLogger подключает структурированный логгер (terminal-write failures, panic).

func WithMaxBacklog

func WithMaxBacklog(n int) WorkerOption

WithMaxBacklog ограничивает размер in-memory admission backlog'а. При переполнении новая задача не enqueue'ится в память (операция остается durable в БД и добирается reconciler'ом) — backpressure без OOM под перегрузкой. n<=0 → без ограничения. Дефолт — defaultMaxBacklog.

func WithMaxInflight

func WithMaxInflight(n int) WorkerOption

WithMaxInflight ограничивает число одновременно исполняемых worker'ов (burst сверх лимита ждет слот в in-memory backlog, ограниченном WithMaxBacklog). Дефолт — 64.

func WithOperationTimeout

func WithOperationTimeout(d time.Duration) WorkerOption

WithOperationTimeout задает верхнюю границу исполнения одной operation-fn. fn получает ctx с этим deadline; по истечении — DeadlineExceeded → durable MarkError, слот семафора освобождается. n<=0 игнорируется (остается дефолт defaultOpTimeout). Значение ДОЛЖНО быть строго меньше Reconciler.OrphanGrace (см. defaultOpTimeout) — иначе живая долгая операция может быть преждевременно помечена reconciler'ом как orphan.

func WithRecorder

func WithRecorder(r Recorder) WorkerOption

WithRecorder подключает sink метрик (inflight gauge, terminal-write retries/ failures). nil → NopRecorder.

func WithTerminalWriteConfig

func WithTerminalWriteConfig(tw TerminalWriteConfig) WorkerOption

WithTerminalWriteConfig настраивает retry-budget durable терминальной записи.

Directories

Path Synopsis
Package operationspb — арендаторская поверхность `OperationService`: перевод строки операции в контракт и обработчик, реализующий оба глагола.
Package operationspb — арендаторская поверхность `OperationService`: перевод строки операции в контракт и обработчик, реализующий оба глагола.

Jump to

Keyboard shortcuts

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