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
- Variables
- func Active() int64
- func CheckRecordedOwnership(caller, recorded Principal) error
- func ConfigureDefault(opts ...WorkerOption) error
- func DefaultRetentionConfig() retention.Config
- func MetadataFor[T proto.Message](op *Operation) (T, error)
- func NewRetentionSweeper(s TerminalSweeper, cfg retention.Config, log *slog.Logger) (*retention.Sweeper, error)
- func NotFoundStatus(id string) error
- func Ready() bool
- func RetentionSubject(s TerminalSweeper) retention.Subject
- func Run(callerCtx context.Context, repo Repo, opID string, ...)
- func RunSync(ctx context.Context, repo Repo, op *Operation, ...) error
- func RunWithWorker(w *Worker, callerCtx context.Context, repo Repo, opID string, ...)
- func Start()
- func StartRetentionSweep(ctx context.Context, s TerminalSweeper, cfg retention.Config, log *slog.Logger) (*retention.Sweeper, error)
- func ValidateListPagination(filter ListFilter) error
- func Wait(ctx context.Context) error
- func WithPrincipal(ctx context.Context, p Principal) context.Context
- func WithoutPrincipal(ctx context.Context) context.Context
- type ClaimRecorder
- type FullRepo
- type ListFilter
- type MemRecorder
- func (m *MemRecorder) ExecutionClaimFailures() float64
- func (m *MemRecorder) ExecutionClaimRetries() float64
- func (m *MemRecorder) IncExecutionClaimFailures()
- func (m *MemRecorder) IncExecutionClaimRetries()
- func (m *MemRecorder) IncOrphansRecovered(outcome string)
- func (m *MemRecorder) IncReconcileErrors()
- func (m *MemRecorder) IncReconcileRuns()
- func (m *MemRecorder) IncTerminalWriteAlreadyResolved(op string)
- func (m *MemRecorder) IncTerminalWriteFailures(op string)
- func (m *MemRecorder) IncTerminalWriteRetries(op string)
- func (m *MemRecorder) Inflight() float64
- func (m *MemRecorder) MaxInflight() float64
- func (m *MemRecorder) OrphansRecovered(outcome string) float64
- func (m *MemRecorder) ReconcileErrors() float64
- func (m *MemRecorder) ReconcileRuns() float64
- func (m *MemRecorder) SetInflight(n float64)
- func (m *MemRecorder) TerminalWriteAlreadyResolved(op string) float64
- func (m *MemRecorder) TerminalWriteFailures(op string) float64
- func (m *MemRecorder) TerminalWriteRetries(op string) float64
- type MetadataFinalizer
- type NopRecorder
- func (NopRecorder) IncExecutionClaimFailures()
- func (NopRecorder) IncExecutionClaimRetries()
- func (NopRecorder) IncOrphansRecovered(string)
- func (NopRecorder) IncReconcileErrors()
- func (NopRecorder) IncReconcileRuns()
- func (NopRecorder) IncTerminalWriteAlreadyResolved(string)
- func (NopRecorder) IncTerminalWriteFailures(string)
- func (NopRecorder) IncTerminalWriteRetries(string)
- func (NopRecorder) SetInflight(float64)
- type Operation
- func ListForCaller(ctx context.Context, repo Repo, filter ListFilter) ([]Operation, string, error)
- func New(domainPrefix, description string, metadata proto.Message) (Operation, error)
- func NewFromContext(ctx context.Context, domainPrefix, description string, metadata proto.Message) (Operation, error)
- type OwnedOperationRepo
- type Owner
- type Principal
- type Reconciler
- type ReconcilerConfig
- type ReconcilerOption
- type Recorder
- type Repo
- type Resolver
- type ResolverResult
- type TerminalOutcome
- type TerminalSweeper
- type TerminalWriteConfig
- type Worker
- type WorkerOption
Constants ¶
const AnonymousPrincipalID = "anonymous"
AnonymousPrincipalID — зарезервированное слово, которым edge помечает запрос БЕЗ credential'а (`Principal{system, anonymous}`). Настоящего принципала с таким id не существует: реальные id — префиксованный crockford-base32 (`usr-…`/`sva-…`), их пространство с этим словом не пересекается.
Слово существует только как ярлык для аудита и логов. Оно НЕ является личностью, и любое сравнение личностей обязано это учитывать — см. Principal.IsAnonymous.
const ClockDrift = 5 * time.Minute
ClockDrift — слагаемое запаса на РАЗНИЦУ ИСТОЧНИКОВ ЧАСОВ.
`modified_at` пишется часами ПРОЦЕССА (`time.Now().UTC()` аргументом), а уборка судит часами БАЗЫ (`now()`) — чтобы источник был один на сервис, а не «столько источников, сколько реплик». Источники разные ⇒ разница входит в порог отдельным слагаемым, а не подразумевается нулём.
`pkg/tokenpolicy.RemovalSlack` здесь НЕ переиспользуется, и причина не в удобстве: та величина объявлена как запас на отставание ВНЕШНИХ читателей и связана размером с временем жизни токена через своего названного потребителя. Здесь закрывается пара «НАШ процесс против НАШЕЙ базы». Переиспользование сделало бы так, что сдвиг одной величины двигает и отсрочку снятия ключа подписи, и порог уборки в другом домене.
ВЕЛИЧИНА СЕГОДНЯ НЕ НЕСУЩАЯ, и это сказано числом, а не умолчано: порог удержания (7 суток) превосходит её на три порядка, поэтому сдвиг часов внутри кластера на минуты не меняет ни одного исхода. Несущей она станет, если `OperationRetention` когда-нибудь опустят к бюджету полла клиента (5 минут) — тогда сумма и станет тем, что сравнивают. Держит это TestOperationRetentionCoversTheDeclaredPollBudget: он утверждает ПРАВИЛО, а не сегодняшний запас.
const OperationRetention = 7 * 24 * time.Hour
OperationRetention — сколько терминальная операция остаётся доступной.
ВЕЛИЧИНА ВЫВЕДЕНА ИЗ ЧИТАТЕЛЯ, а не назначена числом. Читатель — оператор, разбирающий отказавшую мутацию: колонки `error_code`, `error_message`, `error_details` живут ТОЛЬКО здесь. Журнал аудита несёт актора и глагол мутации, но не текст отказа платформы, поэтому после снятия строки причина отказа не восстановима нигде. Предикат читателя — «строка ещё на месте, когда до неё дошли руки».
Интервал между отказавшей мутацией и разбором ограничен снизу нерабочими днями: отказ вечером пятницы разбирают утром понедельника. Порог короче двух суток делает такой разбор невозможным by construction; недельный запас покрывает его вместе с отпускной подменой и остаётся величиной, которую видно оператору, а не запасом, покрывающим произвольную ошибку в слагаемых.
ЧЕГО ВЕЛИЧИНА НЕ ОБЕЩАЕТ. Таблица операций не становится историей ресурса: история — предмет журнала аудита, чьё удержание объявлено намеренно. Списочное чтение после уборки отвечает «операции за последние OperationRetention», а не «все операции ресурса». Это смена наблюдаемого исхода, и она названа здесь, а не обнаруживается.
const SubjectOperations = "operations"
SubjectOperations — имя предмета уборки. Совпадает с именем таблицы: имя попадает в отчёт прохода и в метку метрики, и оператор обязан узнавать в нём таблицу.
Variables ¶
var ErrAlreadyDone = errors.New("operation already completed")
ErrAlreadyDone возвращается из терминальных переходов, если строка уже завершена (CAS-on-`done` не совпал). Маппится в FAILED_PRECONDITION на Cancel-пути; на worker-пути трактуется как идемпотентный no-op.
var ErrNotFound = errors.New("operation not found")
ErrNotFound возвращается из Get, если операция не найдена.
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 CheckRecordedOwnership ¶
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 ¶
DefaultRetentionConfig — величины ПЕТЛИ уборки таблицы операций.
Расписание общее для всех предметов платформы и объявлено один раз в `pkg/retention`; здесь оно только называется, а не переписывается числами.
func MetadataFor ¶
MetadataFor извлекает типизированные метаданные из операции. Возвращает ошибку, если Metadata nil или тип не совпадает.
func NewRetentionSweeper ¶
func NewRetentionSweeper(s TerminalSweeper, cfg retention.Config, log *slog.Logger) (*retention.Sweeper, error)
NewRetentionSweeper собирает уборщика таблицы операций, НЕ поднимая петлю.
Отдельно от запуска, потому что «сборщик работает» обязано быть проверяемо без ожидания тикера: проба зовёт `Pass` методом и не спит.
func NotFoundStatus ¶
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 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 WithPrincipal ¶
WithPrincipal кладет Principal в context. Используется auth-interceptor'ом api-gateway: после валидации JWT и резолва subject через kaname interceptor вызывает WithPrincipal и пробрасывает ctx дальше в handler.
Без auth ctx остается пустым и PrincipalFromContext возвращает SystemPrincipal().
func WithoutPrincipal ¶
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 ¶
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.
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 ¶
ListForCaller — единственная точка, через которую tenant-facing список операций попадает в сервисы: ключ владения выводится из ctx, предикат владения уходит внутрь SQL WHERE (ListOwned), несуженный Repo.List не вызывается никогда.
Зачем отдельная функция, а не «зовите ListOwned сами». Правильный вызов складывается из трёх шагов, и пропуск любого возвращает исходную дыру: получить ownership-апгрейд репозитория, вывести ключ владения именно из ctx (а не из PrincipalFromContext — тот на пустом контексте отдаёт системную личность, которая владеет каждой системно записанной строкой), и на отсутствие ключа отказать, а не продолжить. Собранные в одном месте, эти шаги проверяются одним гейтом; рассыпанные по пятнадцати вызывающим — пятнадцатью прочтениями.
Исходы ровно три:
- ключ владения есть → страница своих операций (полная проекция: строка принадлежит вызывающему, поэтому и содержимое, и личность на ней — его собственные);
- ключа нет (контекст без принципала, принципал снят недоверенным форвардером, именованная анонимность) → ПУСТАЯ страница. Не несуженная выдача и не «своих нет, значит покажем чужие»;
- репозиторий не несёт ownership-апгрейда (ошибка провязки) → отказ с фиксированным текстом. Откат на несуженный путь запрещён: молчаливый откат — это ровно тот исход, ради недопущения которого функция и написана.
Формат страницы проверяется ДО схлопывания в пустую: иначе мусорный курсор у неопознанного вызывающего вернул бы пустую страницу вместо отказа по формату (api-conventions: format-validate → authz → repo).
Несуженный Repo.List остаётся законным для доверенного внутреннего яруса (кластерный аудит, реконсайлер) — там вызывающий уже авторизован иначе, и именно поэтому запрет выражен гейтом с поимённым списком, а не удалением метода.
func New ¶
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 ¶
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 ¶
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 ¶
OwnerFromPrincipal строит Owner из УЖЕ УСТАНОВЛЕННОГО Principal'а. AccountID не заполняется (principal-ключ); IAM-сервис при необходимости выставляет его сам.
ВНИМАНИЕ: не скармливай сюда PrincipalFromContext(ctx). Тот на пустом (и на явно очищенном) ctx отдаёт SystemPrincipal() — backward-compat fallback для ЗАПИСИ операций фоновыми путями, а не удостоверение личности. Как owner-ключ для ЧТЕНИЯ он совпадает с ownership-предикатом на каждой операции, записанной системным принципалом, то есть анонимный вызывающий получает чужие операции. Для вывода owner-ключа из ctx есть OwnerFromContext.
func (Owner) IsAnonymous ¶
IsAnonymous — principal-ключ никого не называет. Такой ключ НИКОГДА не должен доезжать до ownership-предиката, и по двум разным причинам:
- пустые компоненты матчатся со строками, у которых колонки принципала пусты;
- именованная анонимность (`{system, anonymous}`) матчится со строками, записанными ЛЮБЫМ другим безымянным запросом — ключ у них общий по построению, их столько же, сколько имён, то есть одно на всех.
Условие делегировано Principal.IsAnonymous — единый предикат «этот принципал не называет никого» на весь дом.
Про AccountID: у предиката есть вторая, account-ветка, но её ключ сейчас НИКТО не заполняет — во всём дереве Owner строится исключительно из пары (PrincipalType, PrincipalID) через OwnerFromPrincipal/OwnerFromContext, а AccountID остаётся "" и ветка инертна. Поэтому проверка намеренно смотрит только на principal-пару. Если появится account-only владелец (Owner с пустым принципалом и непустым AccountID), этот предикат обязан быть расширен ВМЕСТЕ с ним — иначе такой владелец будет молча отвергнут; правку делать осознанно, а не «ослабить проверку, чтобы прошло».
type Principal ¶
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 ¶
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 ¶
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 ¶
IsAnonymous сообщает, что принципал не называет НИКОГО.
Таких случая два, и они обязаны трактоваться одинаково:
- пара пуста (принципала не было вовсе);
- пара несёт зарезервированное слово AnonymousPrincipalID — «неизвестно кто», выданное запросу без credential'а.
Второй случай — источник целого класса дефектов: как только у «неизвестно кто» появляется ИМЯ, любое сравнение с этим именем даёт истину. Два независимых безымянных запроса получают один и тот же ключ и оказываются друг для друга «своими»; гейт вида «личность извлеклась ⇒ аутентифицирован» на нём срабатывает. Поэтому именованная анонимность обязана вести себя как отсутствие принципала везде, где принимается решение о личности.
Тип при этом не проверяется намеренно: слово зарезервировано целиком, и «неизвестно кто» с любым заявленным типом остаётся неизвестно кем.
НЕ включает ЯВНО установленный системный принципал (`{system, bootstrap}`): это настоящая bootstrap-личность доверенного internal-пути, она остаётся владельцем своих операций. Сужается именно анонимность.
func (Principal) IsSystem ¶
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) 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 ¶
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.
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 терминальной записи.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package operationspb — арендаторская поверхность `OperationService`: перевод строки операции в контракт и обработчик, реализующий оба глагола.
|
Package operationspb — арендаторская поверхность `OperationService`: перевод строки операции в контракт и обработчик, реализующий оба глагола. |