subscription

package
v1.7.0 Latest Latest
Warning

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

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

Documentation

Overview

Package subscription — ОБЩИЙ сервер потока изменений ресурсов. Один на платформу, по экземпляру на владельца журнала.

Зачем пакет существует

Сервер потока был написан ДВАЖДЫ и предметно, в двух разных сервисах: у одного 380 строк с выделенным соединением, курсором и пределом потоков, у другого свой обработчик того же предмета. Потребителя не оказалось ни у одного, а провязать подписку в остальных модулях копированием значило бы завести седьмую копию одного механизма.

Решение владельца: формат подписки один на всех как стандарт, определённый в фундаменте как переиспользуемый механизм; остальное снимается и переводится на него. Этот пакет — тот механизм.

Что приносит владелец, подключаясь

Journal — и ничего больше: где журнал лежит, каким каналом он будит, и как его строка становится событием общей формы. Всё прочее — курсор, горизонт устоявшегося, пределы, сужение по правам, порядок отказов — принадлежит серверу и владельцу не выдаётся. Это и есть предикат качества обобщения: появись у владельца возможность принести своё вместо любого из них, механизм перестал бы быть общим, оставшись общим по имени.

КОНТРАКТ МЕХАНИЗМА

Он записан здесь дословно, потому что подписчик ведёт по нему СВОЁ состояние, и каждое из четырёх утверждений ниже — обязательство, а не подробность.

## Доставка: КАЖДОЕ устоявшееся событие ровно один раз на поток

Событие, попавшее в журнал и ставшее устоявшимся, будет отдано подписчику, чьим осям оно отвечает и на чей взгляд у вызывающего есть право. В пределах ОДНОГО потока повтора не бывает: курсор двигается только вперёд.

Между потоками — «хотя бы раз»: клиент, возобновившийся с позиции `P`, получает всё, что после `P`, и не получает того, что `P` покрывает. Если он сохранил позицию, но не успел обработать событие, повтор произойдёт — и это нормальный исход, к которому подписчик обязан быть готов.

**Событие, НЕ попавшее в журнал, не существует для подписки.** Журнальная строка пишется в той же транзакции, что и ресурсная, поэтому при отказе мутации события нет ВОВСЕ: наблюдатель ресурсов не отличит «ещё делается» от «упало». Полл операции подписка не заменяет.

## Порядок: НЕУБЫВАЮЩИЙ по позиции, в пределах одного журнала

События одного журнала приходят в порядке неубывания позиции; событие, пришедшее позже, никогда не относится к состоянию, предшествующему уже полученному. Это верно и для одного предмета, и для журнала целиком.

**Порядок между журналами разных владельцев НЕ определён** и определён быть не может: у них независимые счётчики и независимые транзакции.

Порядок держится не сортировкой, а ГРАНИЦЕЙ УСТОЯВШЕГОСЯ (см. `watermark.go`): номер выдаётся на вставке, а видимость наступает на фиксации, поэтому строка отдаётся только когда все меньшие номера уже определены — видимы либо потеряны откатом. Без этого писатель, закоммитивший позже с меньшим номером, терялся бы навсегда: молча, без отказа и без пропуска в нумерации, видимого клиенту.

## Разрыв: поток обрывается, ПОЗИЦИЯ переживает

Позиция принадлежит КЛИЕНТУ; сервер её не помнит и остаётся stateless. Поэтому возобновление работает и после перезапуска сервера, и к ДРУГОЙ реплике.

Обрыв — штатное событие, а не отказ. Сервер закрывает поток чисто по истечении срока жизни (Config.StreamBudget), и клиент возобновляется со своей позиции. Отдельно и явно назван исход, при котором возобновиться нельзя: владелец, переставший удерживать часть журнала, отвечает `OUT_OF_RANGE` и называет позицию, с которой подписка ещё возможна, — а НЕ начинает молча с ближайшего удержанного места. Молчаливое начало клиент записал бы как «изменений не было».

Недоступность источника — тоже НЕ «событий нет»: она отвечает `UNAVAILABLE`. Пустой поток есть утверждение о мире, и делать его при неотвеченном чтении сервер не вправе.

## Медленный потребитель: держит СВОЙ поток, не чужие

Отправка идёт синхронно в поток вызывающего: не читающий клиент упирается в окно транспорта, и курсор ЕГО потока перестаёт двигаться. Он не задерживает ни других подписчиков (у каждого свой курсор и своё соединение), ни писателей журнала (сервер только читает).

Плата названа: медленный потребитель занимает слот и выделенное соединение всё время своей жизни. Поэтому число одновременных потоков ограничено (Config.MaxStreams), и превышение отвечает `RESOURCE_EXHAUSTED` — ОТКАЗОМ, а не молчаливой очередью: очередь превратила бы исчерпание в неограниченное ожидание, неотличимое для клиента от «событий нет». Сверху стоит срок жизни потока: слот освобождается даже у клиента, который не читает и не уходит.

**Буфера на подписчика здесь нет намеренно.** Он превратил бы отставание в потерю (переполнился — что-то выбросили) либо в неограниченную память. Вместо этого отставание становится задержкой, и она видна: курсор не двигается, позиция в служебном сообщении при следующем открытии это покажет.

Чего этот пакет НЕ делает

  • НЕ фильтрует по меткам: они мутабельны, ресурс входит в выборку и выходит из неё, и без предыдущего состояния выход неотличим от удаления. Метки отбирает КЛИЕНТ — и ровно там, где владелец состояние ПРОИЗВОДИТ (kacho#1025 заводит серверный отбор, когда объём событий будет измерен).

    **Оговорка несущая, и без неё эта строка лгала.** Прежняя редакция обосновывала клиентский отбор тем, что «событие несёт полное состояние», — утверждением обо ВСЕЙ форме, тогда как состояние производит владелец, а не она. Контракт (`SubscriptionRequest`) оговорку нёс, эта строка — нет: два места об одном предмете, из которых верно одно, и неверным было то, которое читает владелец, подключая новый журнал.

    Кто из владельцев состояние производит — НЕ выписывается ни здесь, ни в клиентской странице по памяти: число владельцев растёт, а выписанное устаревает молча. Согласие двух мест держит гейт `internal/repohygiene` `TestSubscriptionOwnerPageSaysWhatItsJournalDoes`: он обходит объявления журналов, судит УЗЕЛ (доходит ли функция состояния до `anypb.New`) и требует, чтобы таблица владельцев обещала то же. Судить подстрокой он не вправе: имя упаковщика стоит и в объяснениях тех, кто состояния не производит. Перепись он печатает по обеим сторонам.

    Спрашивать об этом документ вообще не нужно: событие без состояния несёт причину `NOT_PRODUCED`, машинно отличимую от сбоя сборки, и она верна на день запроса, а не на день правки страницы. Клиент, которому метки нужны у такого владельца, читает предмет по `resource_id`;

  • НЕ фильтрует по имени: имя мутабельно, и подписка по нему МОЛЧА перестаёт совпадать после переименования. Адресный отбор делает ось `ids`;

  • НЕ подписывает на операции: полл операции дёшев, а поток операций тиражировал бы подписчикам идентификатор, присвоенный до выполнения;

  • НЕ чистит журнал: удержание — свойство владельца, и он его объявляет;

  • НЕ принимает ЖУРНАЛ СМЕНЫ СУБЪЕКТА (`kaname.subject_change_outbox`), и это решение, а не объём работы. Он читается унарным глаголом (`InternalIAMService.PollSubjectChanges`, потребитель — `pkg/subjectchange` на крае), и попытка подключить его сюда была бы не обобщением, а подменой предмета.

    Предметы РАЗНЫЕ, и различие видно по вопросу, на который отвечает событие: этот поток отвечает арендатору «что стало с твоими ресурсами», тот — слою авторизации «твои решения устарели». Отсюда противоположная модель доступа: здесь построчное сужение по правам есть ЗАЩИТА (глагол объявлен `scope_filtered`), там оно fail-open — пропущенная строка означает, что кэш вердиктов не погас, то есть край продолжает отвечать по отозванному праву. Сужение сняло бы ровно тот эффект, ради которого журнал заведён.

    И оно замкнуло бы круг: чтобы отдать строку, владелец спросил бы модель прав, а предмет строки — обесценить кэшированные ответы на такие вопросы. Под fail-closed это даёт «модель недоступна ⇒ край не узнаёт об отзывах», то есть хуже прямого чтения, а не лучше.

    Два расхождения помельче названы здесь же, чтобы их не приняли за единственные: сужатель этой формы типизирован конкретным `*listnarrow.Narrower`, а у владельца прав своя реализация по решению (видимость строки определяется НАБОРОМ отношений, а не одним действием); и у того журнала нет колонки вида предмета, которой Storage требует, — тип субъекта лежит в нагрузке, а `resource_type` там необязательная подсказка, которую читатель не берёт.

    Разбор целиком, включая два отвергнутых исхода и их цену, — `docs/architecture/subject-change-journal-is-not-a-resource-stream.md`; здесь он не пересказывается, чтобы два места об одном предмете снова не разошлись (kacho#1397).

КТО ВЛАДЕЛЕЦ, А КТО НЕТ — и почему двух доменов здесь не будет

Глагол служат пять доменов из семи, и разница — РЕШЕНИЕ, а не остаток работы. Записано здесь, потому что снаружи «решено не заводить» и «ещё не сделано» неразличимы by construction: и то и другое выглядит как отсутствие провязки.

Считать домены по имени таблицы нельзя, и это измерено, а не предположено. Предикат «сервис, у которого есть таблица с именем на `_outbox`» неверен В ОБЕ СТОРОНЫ: он засчитывает очереди прав, аудита, реконсиляции и компенсации (у iam таких четыре, ни одна лентой изменений не является) и НЕ засчитывает настоящий журнал registry — `registry_resource_journal`, названный не так. Единица счёта здесь одна и принадлежит самому механизму: журнал служится тогда, когда его виды суть ТИПЫ МОДЕЛИ ПРАВ, — этого требует [Mapping.validate], и требует ровно потому, что иначе вопрос о видимости строки задать нечем.

По этой единице дерево считается так: лент изменений ресурсов ШЕСТЬ (compute, geo, nlb, registry, storage, vpc), доменов с собственными типами модели прав тоже шесть (те же минус geo, плюс iam), а пересечение — ПЯТЬ, и все пять глагол служат. Пять из пяти, а не пять из семи.

## geo — сужать нечем

`geo_outbox` по форме лента изменений и есть: номер, вид, идентификатор, род, состояние. Но виды у неё `Region` и `Zone`, а типов `geo_*` в модели прав НОЛЬ — каталог размещения намеренно вынесен из пообъектной авторизации: его обязан читать всякий аутентифицированный арендатор, и вопрос «вправе ли этот вызывающий видеть эту зону» имеет один ответ для всех.

Отсюда следует не «трудно», а «невозможно назвать законный словарь»: [Mapping.validate] отвергает пустой словарь видов и вид без типа объекта, то есть владелец geo не поднялся бы. Провязать его можно было бы, лишь заведя `geo_region`/`geo_zone` ради того, чтобы сужатель отвечал «да» всем, — проверка, форма которой есть, а содержания нет.

## iam — предмета нет

У iam ленты изменений ресурсов НЕТ ни одной. Четыре его очереди — кортежей прав, реконсиляции зеркала, смены субъекта и аудита — суть внутренняя механика: они отвечают не арендатору «что стало с твоими ресурсами», а слою авторизации «твои решения устарели».

Предметы разные, и разница не в объёме работы: у ресурсного потока построчное сужение — ЗАЩИТА, а для потребителя журнала субъектов оно fail-open (строка, не отданная из-за прав, означает непогашенный кэш вердикта, то есть край продолжает отвечать по отозванному). Разобрано в kacho#1397, и развилка «чем читать журнал субъектов» РЕШЕНА 2026-08-29 в пользу чтения курсором через унарный `InternalIAMService.PollSubjectChanges`. Здесь решение не принимается и не пересказывается: оно объявлено единственным местом — `docs/architecture/subject-change-journal-is-not-a-resource-stream.md`.

## Как эти два решения истекают

Не памятью. Гейт дерева (`internal/repohygiene`, `TestEveryDomainEitherServesSubscriptionOrRecordsWhyNot`) требует от каждого домена одного из двух — служить глагол либо нести запись решения, — и роняет прогон, когда запись пережила предмет: домен стал владельцем, домен исчез, либо у geo появился свой тип в модели прав.

grant.go — ОКНО МАТЕРИАЛИЗАЦИИ ГРАНТА: почему «нет» от модели прав не всегда окончательно, и чем ограничено ожидание.

ПРЕДМЕТ

Строка журнала и строка ресурса коммитятся ОДНОЙ транзакцией — это несущее свойство формы подписки. Кортеж владения кладётся ПОСЛЕ фиксации: сначала синхронным докладом владельцу прав, а если тот не дошёл — дренажем намерения из той же writer-транзакции. Пробуждение потока приходит НА ФИКСАЦИИ, то есть в самый ранний возможный миг: строка уже видна, гранта ещё нет.

Пообъектный вопрос в этот миг получает «нет» ЗАКОННО. До задачи #2264 такой ответ был ОКОНЧАТЕЛЬНЫМ: курсор идёт по ПРОЧИТАННОЙ строке (см. [Server.drain], и это верно — иначе невидимая партия перечитывалась бы вечно), поэтому строка уезжала под курсор, и открытый поток не отдавал её больше НИКОГДА.

Наблюдаемое следствие — то, ради чего подписка заведена, и оно ТИХОЕ: у клиента нет ни пропуска в нумерации, ни отказа, поэтому «изменений не было» и «изменение было и не доехало» для него неразличимы. Возобновление с той же позиции предмет ПРИНОСИЛО (грант к тому времени уже лежал) — то есть строка была и в журнале, и видна вызывающему, а открытый поток её пропустил.

РЕШЕНИЕ: ОТСРОЧИТЬ СУЖДЕНИЕ, А НЕ РАСШИРИТЬ ЕГО

Рассмотрен и отвергнут симметричный ход: судить создание ПО ЯКОРЮ, как уже судится снятие (см. `removalsAllowedByAnchor`). Довод там — «предмета уже нет в модели прав»; здесь он звучал бы «предмета ещё нет», и симметрия соблазнительна. Отвергнут по цене: у снятия расширение касается ОДНОГО рода изменения и предмета, которого больше нет, а у создания оно означало бы, что всякий, кому разрешён проект, видит КАЖДОЕ создание в нём — постоянно, а не в окне. Это расширение поверхности доступа, и принимать его молча нельзя.

Отсрочка расширением не является BY CONSTRUCTION: наружу не уходит ни одна строка, которой модель не сказала «да». Платится ЛАТЕНТНОСТЬЮ, и она названа ниже числом.

ЧЕГО ОТСРОЧКА НЕ ОБЕЩАЕТ

Она покрывает ДВУХ производителей гранта, разбуженных той же фиксацией: синхронный доклад и дренаж намерения. Реконсайлер — третий производитель, и его шаг измеряется десятками секунд; держать поток столько значило бы платить этой задержкой за КАЖДУЮ строку, которую вызывающему не разрешат никогда. Поэтому окно его не покрывает, а остаток называется ПЕРЕПИСЬЮ (см. [delivery.census]): строка, снятая по истечении окна, попадает в счётчик, и расхождение «прочитано против отдано» перестаёт быть немым.

retention.go — УБОРКА ЖУРНАЛА ПОДПИСКИ: предикат и его пороги.

Предмет — задача продукта #1666. Журнал подписки заведён у пяти владельцев (`polyrepo.md` §runtime-edges), строка в него пишется НА КАЖДОЙ мутации ресурса, а снятия строк не было ни на одном пути ни у одного из них: рост монотонный и вечный, темп задаёт арендатор.

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

ПОЧЕМУ ПРЕДИКАТ ЗДЕСЬ, А НЕ У КАЖДОГО ВЛАДЕЛЬЦА

`pkg/retention` разводит петлю и предикат ровно потому, что предикат — свойство СХЕМЫ: у операций признак `done`, у очереди дренажа `sent_at`. У пяти журналов подписки схема одна и та же по существу — возрастающая позиция плюс отметка времени, — и объявлена она уже: `Storage.PositionColumn` и `Storage.AgeColumn`. Пять копий одного оператора разошлись бы молча и дали бы РАЗНЫЙ срок жизни события у разных доменов при ОДНОМ контракте подписки.

ЦЕНА РЕШЕНИЯ, НАЗВАННАЯ ВСЛУХ

Имя таблицы приезжает сюда значением (`Storage.Table`), поэтому текста имени в исходнике НЕТ, и гейт роста (`internal/repohygiene` TestLiveTablesNameTheirGrowthLimit) этот оператор снятия у себя не резолвит — он уходит в НАЗВАННУЮ САМИМ ГЕЙТОМ слепую зону `RemovalsUnresolved`. Следствие обязано быть записано, а не умолчано: у журнала подписки остаётся запись реестра роста, и её вердикт — «предел» с причиной, называющей этот уборщик.

Обратный размен рассматривался и отвергнут: чтобы гейт видел имя, оператор пришлось бы написать ПЯТЬ раз, по разу у каждого владельца, различая их одним токеном — именем таблицы. Это тот самый разъезд, ради устранения которого заведён `pkg/retention`. Взамен машинную сторону закрывает ДРУГОЙ гейт — `TestSubscriptionJournalLanesAgreeOnRetention`: он сверяет полосы попарно и требует, чтобы объявленное удержание и провязанный уборщик были ОДНИМ решением.

Index

Constants

View Source
const GrantSettlingWindow = 10 * time.Second

GrantSettlingWindow — сколько поток ждёт кортеж владения, прежде чем счесть «нет» окончательным.

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

Ждут здесь ровно одного: чтобы владелец прав узнал о ресурсе, созданном мигом раньше. Производителей, разбуженных той же фиксацией, ДВА — синхронный доклад сразу после коммита и дренаж намерения, лежащего в той же writer-транзакции, — и каждый есть ОДИН вызов к соседу под сроком запроса. Срок такого вызова в этом дереве — секунды (пять у рёбер на пути запроса), поэтому окно взято как два таких срока: первая полоса плюс её замена, когда первая не дошла.

ПОЧЕМУ НЕ ДОЛЬШЕ

Окно есть ПОТОЛОК ЗАДЕРЖКИ, которую платит вызывающий за строку, которую ему не разрешат никогда: такая строка удерживает партию ровно окно и лишь потом снимается. Растянуть его до шага реконсайлера значило бы обменять редкую потерю на постоянную задержку — размен в неверную сторону, потому что потеря теперь ВИДНА переписью, а задержка была бы не видна ничем.

ПОЧЕМУ НЕ КОРОЧЕ

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

ЭТО НЕ ВЕЛИЧИНА ПОСАДКИ

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

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

JournalRetention — сколько событие остаётся возобновимым.

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

Читатель журнала подписки — ПОДПИСЧИК, возобновляющийся с сохранённой позиции. Его предикат: «моя позиция ещё выше нижней возобновимой». Значит порог есть ответ на вопрос «сколько подписчику позволено отсутствовать, не перечитывая каталог целиком».

Снизу его ограничивает не поток запросов, а НЕРАБОЧЕЕ ВРЕМЯ потребителя: потребитель, упавший вечером пятницы, поднимается утром понедельника, и порог короче двух суток делает возобновление невозможным by construction. Недельный запас покрывает это вместе с отпускной подменой у оператора потребителя.

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

`pkg/operations.OperationRetention` выведена из ДРУГОГО читателя (оператор, разбирающий отказавшую мутацию) и совпала числом. Совпадение чисел не есть один предмет: связать их константой значило бы, что сдвиг срока разбора отказов двигает срок возобновления подписки. Величины объявляются порознь, и каждая — рядом со своим читателем.

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

У таблицы операций такое слагаемое есть (`pkg/operations.ClockDrift`): её `modified_at` пишется часами ПРОЦЕССА, а уборка судит часами БАЗЫ, и разница источников входит в порог отдельно. Здесь источник ОДИН: отметка времени строки журнала ставится умолчанием колонки (`DEFAULT now()`), то есть теми же часами базы, которыми судит оператор ниже. Слагаемое, покрывающее разницу, которой нет, было бы запасом без предмета.

Это утверждение о СХЕМЕ ВЛАДЕЛЬЦА, а не о коде здесь, поэтому его держит гейт `TestSubscriptionJournalLanesAgreeOnRetention`: у чистящей полосы колонка срока обязана быть объявлена с `DEFAULT now()` в миграции, заводящей таблицу.

ЧЕГО ВЕЛИЧИНА НЕ ОБЕЩАЕТ

Журнал не становится историей ресурса. Подписчик, отсутствовавший дольше, получает ЯВНЫЙ отказ codes.OutOfRange с названной возобновимой позицией (см. `positionLost` в `server.go`) и перечитывает каталог — а не молча получает неполное. Молчаливое начало с ближайшего удержанного места клиент записал бы как «изменений не было».

Variables

This section is empty.

Functions

func StartJournalRetentionSweep

func StartJournalRetentionSweep(
	ctx context.Context,
	db Execer,
	j Journal,
	cfg retention.Config,
	log *slog.Logger,
) (*retention.Sweeper, error)

StartJournalRetentionSweep поднимает фоновую уборку журнала подписки.

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

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

Types

type Config

type Config struct {
	// Journal — объявление владельца: где журнал, каким каналом будит, как
	// строка становится событием.
	Journal Journal

	// DSN — строка соединения для ВЫДЕЛЕННОГО соединения вне пула. Подписка
	// держит его всё время своей жизни: `LISTEN` требует своей сессии, а
	// сессия из пула вернулась бы в него вместе с подпиской.
	DSN string

	// Narrower — сужатель по правам вызывающего, тот же, что у списков.
	//
	// Обязателен и обязан СУЖАТЬ. За этим методом нет пообъектной проверки на
	// крае (он `scope_filtered`), поэтому откатываться не на что: сужатель,
	// подвешенный и не сужающий, отдал бы весь журнал молча — под кодом,
	// который выглядит фильтрующим.
	Narrower *listnarrow.Narrower

	// ProjectGate — страж оси `project_id`. Обязателен у владельца, чей журнал
	// проектное измерение имеет; у [ProjectAbsent] запрещён — сторожить нечего.
	ProjectGate ProjectGate

	// MaxStreams — потолок числа ОДНОВРЕМЕННЫХ потоков процесса.
	//
	// Каждый поток держит своё соединение вне пула, поэтому потолок — не вкус, а
	// арифметика: число реплик × потолок + непуловые соединения обязаны
	// помещаться в предел владельца базы. Превышение отвечает ОТКАЗОМ, а не
	// молчаливой очередью: очередь превратила бы исчерпание в неограниченное
	// ожидание, неотличимое для клиента от «событий нет».
	MaxStreams int

	// StreamBudget — срок жизни одного потока. Приезжает из объявления сервиса
	// (`servicecontract.Spec.StreamBudget`) и здесь не имеет умолчания: величина
	// посадки, которую никто не выбирал, не обсуждаема и не сужаема.
	//
	// По истечении поток закрывается ЧИСТО (`OK`), а не ошибкой: клиент
	// возобновляется со своей позиции. Обрыв ошибкой читался бы как сетевой сбой.
	StreamBudget time.Duration

	// IdlePoll — холостой перепрос: как часто поток перечитывает журнал, не
	// дождавшись пробуждения.
	//
	// Он не «на всякий случай»: ОТКАТИВШИЙСЯ писатель уведомления не шлёт, и
	// подтверждение горизонта приезжает именно этим перепросом.
	IdlePoll time.Duration

	// Logger — журнал процесса. Ноль резолвится в [slog.Default].
	Logger *slog.Logger
	// contains filtered or unexported fields
}

Config — что приносит КОМПОЗИЦИОННЫЙ КОРЕНЬ владельца, поднимая сервер.

Здесь стоят величины, принадлежащие ПОСАДКЕ, а не журналу: они приезжают из объявления сервиса (`pkg/servicecontract`), и потому не зашиты.

type Execer

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

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

Querier ВЛОЖЕН, а не приставлен рядом вторым параметром: вызывающему нечего выбирать между двумя источниками, и провязать половину пары он не может.

Одного снимка это НЕ даёт и не обязано: наблюдение и оператор идут разными запросами, а на пуле — вообще разными сессиями. Верхняя граница безопасна не снимком, а МОНОТОННОСТЬЮ: устоявшаяся позиция только растёт, поэтому величина, наблюдённая до оператора, к моменту его исполнения остаётся НИЖНЕЙ оценкой — снимется не больше, чем позволено, и в худшем случае меньше. Ошибка в эту сторону стоит одного лишнего прохода; в обратную — потерянного события.

type Filter

type Filter struct {
	// Kinds — виды предметов В НАПИСАНИИ ПРОВОДА (типы объекта модели прав),
	// ровно как их назвал вызывающий. Пусто означает «все виды словаря
	// владельца».
	//
	// Слова хранилища здесь НЕТ намеренно: перевод в них принадлежит объявлению
	// журнала и делается на пути чтения ([Journal.journalWords]). Положи мы сюда
	// переведённое, принятое сужение перестало бы совпадать с тем, что назвал
	// вызывающий, — и перечень честно отобранных осей рядом говорил бы о другом.
	Kinds []string
	// ProjectID — проектный якорь. Пусто означает «проектом не сужаем».
	ProjectID string
	// IDs — идентификаторы предметов. Пусто означает «идентификатором не сужаем».
	IDs []string
	// Honored — оси, которые этот владелец отобрал ЧЕСТНО, именами полей
	// контракта. Ось, задание которой владелец принять не может, сюда не
	// попадает — она отвергает подписку ещё в [Journal.Accept].
	Honored []string
}

Filter — принятое сужение: что задал вызывающий и что владелец отобрал честно.

Заданные оси сужают ВМЕСТЕ (конъюнкция); незаданная не сужает ничем — и это её ЕДИНСТВЕННОЕ значение.

type Journal

type Journal struct {
	// Storage — где журнал лежит и чем он удерживается.
	Storage Storage

	// Channel — имя канала `LISTEN`, на который пишет триггер журнала.
	//
	// Он назван ОТДЕЛЬНО от таблицы, а не выведен из её имени: у трёх владельцев
	// из четырёх они совпадают, у четвёртого таблица схемо-квалифицирована, а
	// канал — нет. Вывод одного из другого работал бы у большинства и молча
	// ошибался у меньшинства — то есть ровно там, где ошибку не ищут.
	Channel string

	// Mapping — отображение строки журнала в событие общей формы.
	Mapping Mapping
}

Journal — ВСЁ, что владелец приносит, подключаясь к общему серверу потока: где журнал лежит, каким каналом он будит, и как его строка становится событием общей формы.

Трёх объявлений достаточно, и это ПРЕДИКАТ КАЧЕСТВА ОБОБЩЕНИЯ, а не эстетика: всё, чего здесь нет — курсор, горизонт устоявшегося, пределы, сужение по правам, порядок отказов, — есть механизм СЕРВЕРА. Появись у владельца возможность принести своё вместо любого из них, механизм перестал бы быть общим, оставшись общим по имени.

Владелец объявляет ЗНАЧЕНИЯ. Ни одного поля функционального типа, кроме отображения строки, здесь нет намеренно: восьмой владелец, желающий свой порядок отказов или свой курсор, не «нарушит правило» — ему не на чем его записать.

func (Journal) Accept

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

Почему отказ, а не пустой поток

Пустой поток есть УТВЕРЖДЕНИЕ «событий нет», а его сервер не вправе делать про вход, которого он не понял. Поэтому каждый негодный вход отвергается ДО открытия потока, с именем поля и названным значением.

Что здесь НЕ проверяется

Существование предмета. Идентификатор правильной формы, которому не отвечает ни один предмет, отказом НЕ является: подписка открывается и просто не приносит по нему событий. «Негодная форма» и «такого нет» — разные полосы, и подписка их не смешивает.

func (Journal) KindDictionary

func (j Journal) KindDictionary() []string

KindDictionary — ЗАКРЫТЫЙ СЛОВАРЬ ВИДОВ владельца в том написании, в котором его видит клиент: типы объекта модели прав, по одному разу, в возрастающем лексикографическом порядке.

ОДИН ВЫЗОВ — ТРИ ЧИТАТЕЛЯ, И ИМЕННО ПОЭТОМУ ОНИ НЕ РАЗОЙДУТСЯ

Отсюда берут словарь все, кому он нужен: служебное сообщение открытия (subscriptionv1.SubscriptionOpened.KnownKinds), отказ на неизвестном виде (Journal.Accept) и пробы владельцев. Второго перечня не существует, поэтому «объявленное» и «то, чем сервер на самом деле судит» — один объект, а не два похожих.

ПОЧЕМУ СЛОВО ПРОВОДА — ТИП ОБЪЕКТА, А НЕ КЛЮЧ СЛОВАРЯ

Ключ — слово ХРАНИЛИЩА: как владелец записал строку в своей таблице. Оно частное и разное у всех: перепись двух живых владельцев на день заведения этого метода дала `Instance` у одного и `nlb_load_balancer` у другого — два написания одного рода предмета, ни одного производителя. Тип объекта, наоборот, уже есть единственное платформенное имя предмета: им же сервер спрашивает модель прав о видимости строки, им же его знает каталог прав и аннотации контрактов домена.

Следствие, ради которого выбрано именно оно: ВТОРОЕ НАПИСАНИЕ ЗАВЕСТИ НЕКУДА. У владельца нет поля, в которое он мог бы положить своё клиентское слово, — не «нельзя по правилу», а негде.

ПОРЯДОК — ЧАСТЬ КОНТРАКТА

Словарь живёт отображением, обход отображения в Go случаен by construction. Неупорядоченный перечень клиент, ведущий состояние по индексу, читал бы как смену словаря на каждом открытии.

ПОВТОР СНИМАЕТСЯ

Два слова хранилища об одном предмете (историческое и нынешнее) — законный случай, и клиенту он не виден: он спрашивает про предмет, а не про то, чем его записали.

func (Journal) Validate

func (j Journal) Validate() error

Validate судит объявление владельца.

Он судит ОБЪЯВЛЕНИЕ, а не дерево: существует ли таблица и та ли у неё форма — вопрос подъёма, и на него отвечает первый же запрос. Здесь закрывается то, что после подъёма закрыть уже нечем: незаявленная ось, разошедшиеся половины одного решения и имя, которое нельзя безопасно поставить в запрос.

type JournalSweeper

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

JournalSweeper — уборщик журнала ОДНОГО владельца.

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

func NewJournalSweeper

func NewJournalSweeper(db Execer, j Journal, log *slog.Logger) (*JournalSweeper, error)

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

Отказывает на объявлении, при котором уборка невыразима: владелец, не назвавший чистку, уборщика не получает вовсе — иначе оператор снимал бы строки у журнала, чей контракт обещает подписчику, что отказ «позиция утрачена» не наступает никогда.

func (*JournalSweeper) RetentionSubject

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

RetentionSubject — запись реестра уборки для журнала этого владельца.

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

func (*JournalSweeper) Sweep

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

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

Момент времени входом НЕ приходит: часы уборки — базы, те же, что ставят отметку строки (см. JournalRetention, §«У порога нет слагаемого на разницу часов»). Проверяемость от этого не теряется — проход зовётся методом (`retention.Sweeper.Pass`), поэтому пробе не приходится ни спать, ни подменять часы.

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

type Kind

type Kind struct {
	ObjectType string
	Action     string
}

Kind — как авторизовать событие этого вида: тип объекта модели прав и действие.

Действие несётся В ТОЙ ЖЕ записи, а не берётся одно на всех: иначе добавленный позже вид унаследовал бы чужой глагол, и модель судила бы о нём по неверному действию, оставаясь при этом «зелёной».

type Mapping

type Mapping struct {
	// Kinds — ЗАКРЫТЫЙ словарь видов владельца.
	//
	// Он закрыт в обе стороны. Наружу: вид вне словаря отвергается
	// `INVALID_ARGUMENT`, а не открывает поток, в который никогда ничего не
	// придёт. Внутрь: строка журнала с видом вне словаря не имеет типа объекта
	// модели прав, поэтому вопрос «вправе ли вызывающий её видеть» задать
	// НЕЛЬЗЯ — и строка не доставляется.
	Kinds map[string]Kind

	// Changes — словарь рода изменения: слово владельца → род платформы.
	//
	// Отдельный словарь, а не работа отображения, потому что род изменения
	// нужен серверу ДО того, как состояние собрано: событие, чьё состояние
	// собрать не удалось, всё равно обязано назвать род — подписчик по нему
	// ведёт своё состояние.
	Changes map[string]subscriptionv1.SubscriptionEvent_Change

	// Anchor — проектный якорь строки. Заполняется тогда и только тогда, когда
	// [Storage.Project] равен [ProjectFromMapping].
	//
	// Отказ здесь — НЕ «состояние недоступно», а невозможность авторизовать:
	// без якоря решение о показе принять не из чего. Такая строка не
	// доставляется, и это громко.
	Anchor func(Row) (string, error)

	// State — состояние предмета. Возвращает ТРИ вещи, и у каждой позиции ровно
	// одно значение: собранное состояние · НАЗВАННАЯ причина его отсутствия ·
	// отказ сборки.
	//
	// Отказ событие НЕ отменяет: оно доставляется с признаком «состояния не
	// будет». Пустая нагрузка вместо состояния запрещена формой by construction.
	//
	// # ПОЧЕМУ ПРИЧИНУ НАЗЫВАЕТ ВЛАДЕЛЕЦ, А НЕ ВЫВОДИТ СЕРВЕР
	//
	// Причин отсутствия по контракту четыре, и они означают РАЗНОЕ ДЕЙСТВИЕ
	// подписчика: на «не удалось собрать» разумно перечитать, на «владелец такого
	// не производит» — сразу идти за предметом, на «не удерживается» и «не
	// показывается» — не ходить вовсе. Знание, какая из них верна, есть у
	// ВЛАДЕЛЬЦА и больше ни у кого: сервер видит `nil` и о его происхождении
	// сказать не может.
	//
	// До этой пары возвращаемых значений владельцу было НЕЧЕМ назвать причину:
	// пара «состояние, отказ» имеет две ветви на четыре смысла, и сервер сводил
	// обе к «не удалось сериализовать». То есть у одного владельца КАЖДОЕ событие
	// объясняло отсутствие состояния неудавшейся попыткой, которой не было. Вывод
	// причины из бедности нагрузки был рассмотрен и отвергнут: он сработал бы и
	// на настоящем отказе сборки, и два противоположных исхода стали бы
	// неразличимы — ровно та потеря различения, ради недопущения которой словарь
	// причин и объявлен закрытым.
	//
	// Соблазн назвать причину СЕНТИНЕЛЬНОЙ ОШИБКОЙ отвергнут по существу:
	// отсутствие состояния — не сбой, и класть его в канал ошибок значит требовать
	// от каждого читателя помнить, что здесь `error` иногда означает «всё в
	// порядке».
	//
	// Владелец, вернувший `nil` без названной причины, получает
	// `REASON_UNSPECIFIED` — контракт держит это значение ровно для такого случая
	// («владелец обязан назвать причину») — и громкую строку в журнале процесса.
	// Корзины «прочее» здесь нет: неназванное названо неназванным, а не подшито к
	// соседней записи.
	State func(Row) (*anypb.Any, StateAbsence, error)
}

Mapping — отображение строки журнала в событие общей формы.

type ProjectDimension

type ProjectDimension uint8

ProjectDimension — откуда берётся проектный якорь события. Состояния ТРИ, и каждое названо: умолчания у якорной оси не бывает.

const (
	// ProjectDimensionUnset — владелец не сказал ничего. [Journal.Validate]
	// отвергает: «пусто» и «не применимо» здесь неразличимы, а решения по ним
	// противоположны.
	ProjectDimensionUnset ProjectDimension = iota

	// ProjectInColumn — якорь лежит колонкой, и ось отбирается запросом.
	ProjectInColumn

	// ProjectFromMapping — колонки нет, якорь даёт отображение. Ось отбирается
	// сервером ПОСЛЕ отображения строки — то есть по-прежнему СЕРВЕРОМ, а не
	// клиентом. Цена названа: прочитано будет больше, чем отдано.
	ProjectFromMapping

	// ProjectAbsent — у журнала проектного измерения НЕТ вовсе (предмет уровня
	// аккаунта или кластера). Подписка, назвавшая эту ось, ОТВЕРГАЕТСЯ: открыть
	// её значило бы отдать поток, который молчит навсегда.
	ProjectAbsent
)

type ProjectGate

type ProjectGate struct {
	// ObjectType — тип объекта проекта в модели прав.
	ObjectType string
	// Action — действие, которым спрашивается доступ.
	Action string
	// Relations — отношения, любого из которых достаточно.
	Relations []string
	// NotFoundFormat — форма отсутствия владельца с ОДНИМ `%s` под
	// идентификатор (`"Project %s not found"`). Ею отвечает и отказ доступа.
	NotFoundFormat string
}

ProjectGate — как спросить модель прав про ДОСТУП К ПРОЕКТУ, названному осью подписки.

Почему отдельный вопрос, если каждая строка и так сужается

Без него вызывающий, назвавший недоступный ему проект, получил бы ОТКРЫТЫЙ поток, молчащий вечно: ни одна строка не прошла бы построчное сужение, и это выглядело бы как «изменений нет». Молчание — утверждение о мире, и делать его про проект, которого вызывающий не вправе видеть, сервер не может.

Почему форма отказа приносится владельцем

Отказ обязан быть НЕОТЛИЧИМ от «такого проекта нет»: различимый текст превращает подписку в способ узнать существование чужого проекта. Форма отсутствия принадлежит владельцу — он же отвечает ею на обычное чтение, — поэтому она приносится сюда, а не сочиняется здесь второй раз.

type Querier

type Querier interface {
	QueryRow(ctx context.Context, sql string, args ...any) pgx.Row
}

Querier — то, чем наблюдение задаёт свой единственный вопрос.

Соединение и пул отвечают ему ОБА, и это не удобство, а условие переиспользуемости: у механизма подписки наблюдение живёт на выделенном соединении потока, а у возобновимого чтения, отвечающего на запрос (`InternalIAMService.PollSubjectChanges`), своего соединения нет вовсе — есть пул.

Исключение СЕБЯ по идентификатору обслуживающего процесса остаётся верным и на пуле: наблюдатель в журнал не пишет, поэтому его собственная блокировка — не та, которую наблюдение ищет.

type Retention

type Retention uint8

Retention — что владелец обещает про удержание журнала. Ноль объявлением не является: «удерживаю всё» — свойство владельца, а не умолчание формы.

const (
	// RetentionUnset — владелец не сказал ничего.
	RetentionUnset Retention = iota

	// RetainsEverything — журнал не чистится; отказ «позиция утрачена» не
	// наступает никогда, и служебное сообщение говорит это прямо.
	RetainsEverything

	// RetainsFromEarliestRow — журнал чистится; нижняя возобновимая позиция
	// выводится из самой ранней удержанной строки.
	RetainsFromEarliestRow
)

type Row

type Row struct {
	// Position — номер строки в журнале владельца.
	Position int64
	// Kind — вид предмета в словаре владельца.
	Kind string
	// ID — идентификатор предмета.
	ID string
	// ProjectID — якорь из колонки; пуст, если якорь даёт отображение.
	ProjectID string
	// Change — род изменения словом владельца.
	Change string
	// Payload — состояние предмета, как оно лежит в журнале.
	Payload []byte
}

Row — строка журнала, как её прочитал сервер.

type Server

type Server struct {
	subscriptionv1.UnimplementedInternalSubscriptionServiceServer
	// contains filtered or unexported fields
}

Server — ОБЩИЙ сервер потока изменений. Один на платформу, по экземпляру на владельца журнала.

Он реализует `subscriptionv1.InternalSubscriptionServiceServer`, и владелец регистрирует ЕГО на своём внутреннем слушателе — не свою обёртку вокруг него.

func NewServer

func NewServer(cfg Config) (*Server, error)

NewServer собирает сервер и судит объявление.

Отказ здесь — отказ ПОДЪЁМА, а не запроса: величина посадки, о которой никто не сказал, не должна обнаруживаться первым запросом в бою.

func (*Server) Subscribe

Subscribe — единственный глагол подписки.

Порядок отказов, и он ЕСТЬ ПРЕДМЕТ

  1. ЛИЧНОСТЬ — безусловно и первой. Ответ безымянному не зависит от посадки: его никто не назвал, и это верно при любой конфигурации;
  2. СУЖАТЕЛЬ — есть ли он и сужает ли;
  3. ФОРМА ЗАПРОСА — виды, идентификаторы, начало. Терминальные отказы обязаны наступать ДО обращения к базе: иначе вызывающий получает повторяемый код на ввод, который валидным не станет никогда, а отказ начинает зависеть от доступности базы;
  4. ДОСТУП К ПРОЕКТУ — прежде чем занимать слот и соединение;
  5. СЛОТ — расход ограниченного ресурса;
  6. СОЕДИНЕНИЕ и `LISTEN`.

type Start

type Start struct {
	// FromBeginning — с начала журнала: всё, что владелец ещё удерживает.
	FromBeginning bool
	// Position — с этой позиции. Ноль указателя означает, что позиция не
	// названа: отсутствие представимо отдельно от всякого значения, потому что
	// нулевая позиция — законная величина («ничего ещё не устоялось»).
	Position *pagetoken.SubscriptionPosition
}

Start — с какого места отдавать. Три состояния, и все три различимы.

func AcceptStart

func AcceptStart(req *subscriptionv1.SubscriptionRequest) (Start, error)

AcceptStart разбирает начало потока.

ИСХОД НЕЗАДАННОГО НАЗВАН, а не подразумевается: `start` не задан означает ровно `CURRENT_END`. Молчаливая полная выдача журнала превратила бы подписку в выгрузку — тем более дорогую, чем дольше журнал не чистили, — и узнать об этом вызывающий мог бы только по счёту за неё.

Выбрав ВЕТВЬ, вызывающий обязан её назвать: умолчание живёт у незаданной ветви, а не у незаданного значения внутри выбранной. Иначе одно и то же намерение выражалось бы двумя способами.

type StateAbsence

type StateAbsence uint8

StateAbsence — ПОЧЕМУ состояния нет. Словарь закрыт и отвечает словарю контракта один в один: заводя здесь значение, которого нет там, владелец получил бы причину, которую некуда положить.

Тип отдельный, а не булев признак: «состояния нет» и «почему его нет» — разные вопросы, и подписчик действует по второму.

const (
	// StateAbsenceUnnamed — владелец причину НЕ НАЗВАЛ.
	//
	// Законным исходом не является: значение нулевое ради того, чтобы забытая
	// причина была ВИДНА, а не подшита к соседней. Сервер отдаёт по нему
	// `REASON_UNSPECIFIED` и пишет громко.
	StateAbsenceUnnamed StateAbsence = iota

	// StateNotProduced — владелец состояния для этого предмета НЕ ПРОИЗВОДИТ.
	//
	// Свойство журнала, а не сбой: собирать было нечего, попытки не было, повтор
	// ничего не изменит. Два живых случая: журнал, не несущий состояния ни у
	// одного вида, и событие снятия — предмета больше нет by construction.
	StateNotProduced

	// StateNotRetained — состояние на эту позицию владельцем больше НЕ
	// УДЕРЖИВАЕТСЯ.
	//
	// Производителя в дереве сегодня НЕТ, и это утверждение о ВСЕХ владельцах, а не
	// о тех, что были видны заводившей полосе: удержание объявляют предикатом
	//
	//	git grep -c 'Retention: subscription\.' -- 'services/*/internal/subscriptionjournal/journal.go'
	//
	// и на сегодня каждая найденная строка называет [RetainsEverything]. Прежняя
	// редакция говорила «оба владельца» — верно для двух журналов, приведённых к
	// этой форме первыми, и неверно для остальных трёх, приехавших соседними
	// полосами: число мерилось на популяции полосы, а не дерева.
	//
	// Значение объявлено потому, что объявлено контрактом, и владельцу, заведшему
	// чистку журнала, будет чем назвать этот исход, — а не «на будущее»: без него
	// он назвал бы соседнее.
	StateNotRetained

	// StateWithheld — состояние ЕСТЬ, но вызывающему не показывается.
	StateWithheld
)

type Storage

type Storage struct {
	// Table — таблица журнала, при необходимости со схемой (`kacho_vpc.vpc_outbox`).
	Table string

	// PositionColumn — колонка возрастающего номера (`bigint`), выдаваемого на
	// вставке. Номер НЕ есть позиция подписки: позицию производит сервер как
	// границу устоявшегося (см. `watermark.go`).
	PositionColumn string

	// KindColumn — колонка вида предмета.
	KindColumn string

	// IDColumn — колонка неизменяемого идентификатора предмета.
	IDColumn string

	// ChangeColumn — колонка рода изменения в словаре ВЛАДЕЛЬЦА.
	ChangeColumn string

	// PayloadColumn — колонка состояния предмета (`jsonb`), отдаваемая
	// отображению как есть.
	PayloadColumn string

	// ProjectColumn — колонка проектного якоря. Заполняется тогда и только
	// тогда, когда [Storage.Project] равен [ProjectInColumn].
	ProjectColumn string

	// Project — откуда берётся ПРОЕКТНЫЙ ЯКОРЬ. Ось якорная: по ней принимается
	// решение о показе, поэтому нулевое значение объявлением не является.
	Project ProjectDimension

	// Retention — что владелец обещает про удержание журнала.
	Retention Retention

	// AgeColumn — колонка отметки времени, по которой судится возраст строки.
	//
	// ПАРА к [Storage.Retention], а не независимое поле: заполняется тогда и
	// только тогда, когда объявлено [RetainsFromEarliestRow]. «Журнал чистится»
	// и «по чему судить возраст» — ОДНО решение, и половина его не выражается:
	// без колонки предикат уборки не построить вовсе, а колонка без чистки —
	// объявление, которого не читает никто (класс «принято-и-проигнорировано»,
	// `api-conventions.md`).
	//
	// Обе стороны отвергает [Journal.Validate]; уборщик, который эту колонку
	// читает, — `retention.go`.
	AgeColumn string
}

Storage — координаты журнала и объявленные свойства его формы.

Имена колонок стоят здесь, а не выводятся из «стандартной схемы», потому что стандартной схемы НЕТ: перепись четырёх ресурсных журналов дерева дала ТРИ разные формы — у одного вид предмета зовётся `resource_type`, а род изменения `action`; у остальных — `resource_kind` и `event_type`. Разнобой форм хранения дефектом не является: единой объявлена форма ПОДПИСКИ, а не форма журнала.

type Watermark

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

Watermark — ВОДЯНОЙ ЗНАК: граница устоявшегося в журнале одного потока.

Имя взято у источника намеренно. Приёмка фазы Ф1 назначила предикатом переноса техники именно его: «`watermark` в дереве — не ноль ВНЕ `services/nlb`». Пока техника звалась здесь иначе, предикат отвечал «не перенесена» на перенесённой технике — и снятие мёртвого читателя nlb (kacho#1043) уничтожило бы единственный написанный ответ на этот класс, формально не нарушив ни одного условия.

Что здесь перенесено, а не изобретено

Техника взята у `services/nlb/internal/repo/kacho/pg/lifecycle_feed.go`, где она написана, разобрана и покрыта пробой на трёх подслучаях. Перенос — а не повторное изобретение — был условием фазы: единственный написанный ответ на этот класс жил там, и снятие того файла до переноса уничтожило бы его.

Обобщено ровно одно: имя таблицы и имя колонки номера пришли параметрами. Правило, разбор отката и размен остались дословно.

Правило

Максимальный ВИДИМЫЙ номер M становится границей, когда ни одна транзакция, писавшая в журнал в момент наблюдения M, больше не пишет. Только такая транзакция может ещё выпустить номер ≤ M: писатель, взявший блокировку ПОСЛЕ наблюдения, вызовет счётчик позже и получит номер > M.

Почему блокировка таблицы, а не снимок транзакций

Писатель держит `RowExclusiveLock` на журнале с момента ПЛАНИРОВАНИЯ своей вставки — то есть ещё до того, как счётчик выдал ему номер, — и до конца транзакции. Идентификатор транзакции назначается позже, на самой вставке: существует наблюдаемое состояние «номер уже выдан, идентификатора транзакции ещё нет», в котором горизонт по снимку транзакций СЛЕП. Дополнительно: верхняя граница снимка не является верхней границей работающих транзакций, поэтому «нижняя догнала прежнюю верхнюю» не доказывает ничего.

`virtualtransaction` идентифицирует ТРАНЗАКЦИЮ, а не соединение: тот же процесс в следующей транзакции получит другой идентификатор. Поэтому непрерывно пишущие фоновые работники не могут удерживать горизонт вечно — зафиксированное множество всегда доистекает.

Осознанный размен

Писатель, не завершающий транзакцию, задерживает поток целиком: у него может быть сколь угодно малый невыпущенный номер, и никакое наблюдение этого не опровергнет. Выбор в пользу «не потерять» против «доставить сейчас»; удержание дольше [stallWarnAfter] — в лог. # Наблюдатель ОДИН на дерево, и это держится гейтом

Экземпляров у него столько, сколько журналов; РЕАЛИЗАЦИЯ одна, и второй такой файл — находка гейта `settledwatermarksingularity` ГДЕ БЫ ОН НИ ЛЕЖАЛ, включая сам `pkg/`. Поэтому потребитель вне механизма подписки берёт этот тип, а не пишет своё наблюдение: два наблюдения одного класса расходятся молча.

Тип ЗАЩИЩЁН ОТ СОВМЕСТНОГО ДОСТУПА

У механизма подписки наблюдатель свой на каждый поток, и состязания нет. У чтения, отвечающего на запрос, экземпляр один на процесс, а проходов столько, сколько реплик потребителя поллит эту реплику владельца. Замок неоспариваемый в первом случае и несущий во втором.

func NewWatermark

func NewWatermark(table, positionColumn string, log *slog.Logger) *Watermark

NewWatermark собирает наблюдателя для одной таблицы.

Имя таблицы и колонки попадают в текст запроса, поэтому вызывающий обязан подавать сюда СВОИ имена, а не пришедшие снаружи: наблюдение их не экранирует. У механизма подписки за это отвечает Journal.Validate; у прямого потребителя — то, что оба имени объявлены константами его же пакета.

func (*Watermark) Advance

func (h *Watermark) Advance(ctx context.Context, q Querier) error

advance двигает границу устоявшегося по одному наблюдению.

Два исхода:

  • писателей нет — граница переносится сразу (обычный ненагруженный случай, задержка доставки нулевая);
  • писатели есть — наблюдение запоминается и подтверждается на следующем проходе. Проход инициирует пробуждение того самого писателя (он держит блокировку именно потому, что вставляет строку, а фиксация доставляет его уведомление), поэтому нормальная задержка — время его транзакции. Исключение — ОТКАТ писателя: уведомления не будет, и подтверждение приедет на холостом перепросе.

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

func (*Watermark) Established

func (h *Watermark) Established() bool

Established — подтверждена ли граница хоть однажды.

Признак МОНОТОНЕН: подтвердившись, он не отзывается появлением нового писателя. Поэтому отказ вызывающего, читающий его, есть состояние ХОЛОДНОГО СТАРТА, а не режим работы.

func (*Watermark) Floor

func (h *Watermark) Floor(r Retention) int64

Floor — нижняя возобновимая позиция.

У владельца, удерживающего журнал целиком, её не существует, и служебное сообщение говорит это отдельным признаком, а не нулём в поле позиции. У владельца, который журнал чистит, она выводится из самой ранней удержанной строки: с позиции `earliest-1` возобновление ещё не теряет ничего.

Пустой журнал у чистящего владельца означает, что удержано ноль строк, и возобновиться можно только с текущей границы: всё, что было до неё, уже снято.

Почему ЭКСПОРТИРОВАН, хотя механизм подписки зовёт его изнутри

Пол — половина ответа на «обнаружил ли читатель пропуск»; вторую половину (явный отказ) даёт тот, кто пол спросил. Пока имя было неэкспортировано, спросить его мог ТОЛЬКО механизм подписки, а тот же класс живёт у журнала, читаемого унарным глаголом (`kaname.subject_change_outbox`, `InternalIAMService.PollSubjectChanges`): окно `id > since AND id <= settled` снятой строки не показывает НИЧЕМ — курсор переезжает через неё, и «строк не было» становится неотличимо от «строки убрали».

Экспорт — не послабление инкапсуляции, а исполнение уже объявленного правила: наблюдатель границы в дереве ОДИН и лежит здесь, поэтому потребитель вне механизма подписки берёт этот тип, а не пишет своё наблюдение. Второй расчёт пола разошёлся бы с этим молча.

Полосу удержания называет ВЫЗЫВАЮЩИЙ, и это его решение, а не свойство типа

RetainsEverything означает «нижней границы не существует» и делает отказ «позиция утрачена» непроизводимым; RetainsFromEarliestRow — «журнал чистится, нижняя граница реальна». У механизма подписки значение приходит объявлением журнала и едет клиенту служебным сообщением; у прямого потребителя оно на провод не выходит вовсе и означает ровно то, что читатель обязан ВЫДЕРЖАТЬ.

func (*Watermark) ObserveFloor

func (h *Watermark) ObserveFloor(ctx context.Context, q Querier, r Retention) (int64, error)

ObserveFloor спрашивает нижнюю удержанную строку и отдаёт пол, выведенный ИЗ ЭТОГО ЖЕ наблюдения, а не из общего поля.

Зачем отдельный вопрос, если есть [Watermark.Floor]

Наблюдатель бывает ОБЩИМ. У потока подписки он свой на каждый поток, и там разницы нет; у чтения, отвечающего на запрос, он ОДИН на процесс, а проходов столько, сколько реплик потребителя поллит эту реплику владельца. Тогда Watermark.Floor отдаёт значение, положенное туда ЧУЖИМ проходом, — и положенное им наблюдение может быть СТАРШЕ моего: `observe` пишет `earliest` безусловно, поэтому проход, чей запрос отработал раньше, а запись легла позже, возвращает поле НАЗАД.

Направление вреда здесь несимметрично и потому несущее: пол, занизившийся чужим старым наблюдением, отказа НЕ производит — то есть страница со снятым префиксом уезжает вызывающему как полная. Это ровно тот молчаливый пропуск, ради обнаружения которого пол и заведён.

Общее поле двигается ТОЛЬКО ВПЕРЁД

Нижняя строка непустого журнала монотонна: строки снимаются префиксом, номера растут. Поэтому наблюдение, оказавшееся ниже уже известного, — старое, и принимать его в общее поле нельзя. Свой ответ метод при этом отдаёт из своего наблюдения в любом случае: общее поле — удобство для соседей, а не источник этого ответа.

Ноль означает ПУСТОЙ журнал, а не «раньше всех», и в монотонность не входит: пустота — состояние, а не позиция.

func (*Watermark) RefreshEarliest

func (h *Watermark) RefreshEarliest(ctx context.Context, q Querier) error

RefreshEarliest перечитывает ТОЛЬКО нижнюю удержанную строку.

Зачем отдельный вопрос, а не [Watermark.Advance]

Владелец, который журнал ЧИСТИТ, двигает нижнюю границу под работающим потоком. Читающий поток обязан узнать об этом ДО того, как отдаст следующую партию: иначе снятые строки просто не придут в выборку, курсор переедет через них, и подписчик получит НЕПОЛНОЕ, ничем не отличимое от «изменений не было».

Полное наблюдение отвечать на этот вопрос негодно по цене: оно несёт ещё и просмотр блокировок журнала, а спрашивать границу приходится перед КАЖДОЙ партией. Здесь — один просмотр индекса по первичному ключу.

Верхнюю границу метод НЕ трогает: устоявшееся двигает только полное наблюдение, и подмешивать сюда его половину значило бы завести второй путь к одной величине.

МОМЕНТ вопроса — не раньше выборки

«Узнать ДО выдачи партии» и «узнать до её ВЫБОРКИ» — разные моменты, и второй негоден: между двумя запросами уборка вправе зафиксироваться, и наблюдение, взятое раньше страницы, описывает журнал, которого к моменту выборки уже нет. Разбор — у вызывающего ([Server.drain]); здесь он не пересказывается, чтобы два места об одном предмете не разошлись.

func (*Watermark) Settled

func (h *Watermark) Settled() int64

Settled — граница устоявшегося на момент последнего наблюдения.

Каждый номер ≤ этой границы либо уже видим, либо не появится НИКОГДА (писатель откатился). Окно возобновимого чтения — `(курсор, Settled]`, и никогда «всё, что больше курсора»: позиция, выданная за неустоявшийся номер, теряет его молча и навсегда.

Ноль читать как «журнал пуст» НЕЛЬЗЯ — спрашивай Watermark.Established: у холодного наблюдателя первое наблюдение лишь ЗАПОМИНАЕТ писателей, и до подтверждения ноль означает «позиции ещё нет».

Jump to

Keyboard shortcuts

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