subjectchange

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Sep 13, 2026 License: AGPL-3.0 Imports: 11 Imported by: 0

Documentation

Overview

Package subjectchange — ЧИТАТЕЛЬ журнала смены субъекта, живущий в фундаменте.

Направление: соединение открывает ПОТРЕБИТЕЛЬ

Владелец прав объявлен ЛИСТОМ графа рёбер: его зовут, он не зовёт никого. Прежде он сам дозванивался до потребителя и гасил его кэш толчком — ребро из листа обратно, да ещё и с обязательной ручкой адреса потребителя, из-за которой владелец не поднимался там, где потребителя нет вовсе. Толчок снят (задача #1024); читатель — здесь, и открывает соединение он.

Почему в фундаменте, а не у потребителя

Свойство «смена прав доезжает до кэша решений» держится ОДНОЙ реализацией. Живи читатель у каждого потребителя, свойство держалось бы столькими независимыми копиями, сколько потребителей, — и достаточно одной забывчивой, чтобы отзыв пережил себя у одного из них, причём незаметно: у остальных проба зелёная.

Побочное и не менее важное: живя здесь, читатель проверяется ОДНОЙ пробой вместе с производителем журнала — сквозь обе стороны, а не по половине. Пока он лежал в дереве потребителя, такая проба была невыразима: правило видимости `internal/` не пускает потребителя к производителю и обратно.

Что он делает — ДВА дела, а не одно

Читает журнал курсором по возрастанию позиции и на всякой непустой порции:

  • гасит кэш решений СВОЕГО процесса — это полоса ЗАПРОСА;
  • закрывает открытые потоки НАЗВАННЫХ субъектов — это полоса СОЕДИНЕНИЯ (задача #1022).

Второе не выводится из первого и первым не покрывается. Кэш решений отвечает на СЛЕДУЮЩИЙ запрос, а длинное соединение следующего запроса не делает: сброс кэша выглядит отзывом и потока не касается вовсе — тот самый «контроль, действующий на выдаче, но не на предъявлении».

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

Index

Constants

View Source
const JournalRetention = 10 * time.Minute

JournalRetention — НАИБОЛЬШЕЕ ОТСТАВАНИЕ ЧИТАТЕЛЯ, которое владелец журнала обязан обслужить точно, и потому же — срок, в течение которого строка журнала остаётся возобновимой.

Величина ВЫВЕДЕНА из читателя, а не назначена из удобства

Читатель — реплика края, ведущая курсор в памяти процесса. Отсюда первый и несущий факт: КУРСОР НЕ ПЕРЕЖИВАЕТ ПРОЦЕССА. Свежая реплика праймится на голову журнала (Watcher.Poll, ветвь `!w.primed`) и удержанного хвоста не читает вовсе. Значит курсор отстаёт ТОЛЬКО пока процесс жив и не может опросить владельца, — а не «пока потребителя нет».

Отсюда порог: он обязан покрывать самое долгое молчание ЖИВОГО читателя, которое платформа ещё считает рабочим состоянием. Читатель называет эту границу сам — Config.StaleAfter: за ней он объявляет себя просроченным и на каждом следующем перепросе закрывает ВСЕ открытые потоки ([Watcher.failClosed]).

Чем это держится, а не обещается

New отвергает читателя, объявившего `StaleAfter` больше этой величины. Состояние «порог владельца короче того молчания, которое читатель считает здоровым» становится НЕВЫРАЗИМЫМ, а не маловероятным: посадка, объявившая такое, не поднимается.

Почему не короче

Единственное ШТАТНОЕ событие, при котором живой читатель не может опросить владельца, — перекат самого владельца прав. Десять минут покрывают его с запасом, и запас этот платится строками, а не безопасностью.

Сегодняшняя посадка края объявляет `StaleAfter = max(5 × 2s, 10s) = 10s` (`gateway/cmd/api-gateway`, `revocationStaleAfter`) — на два порядка ниже порога, то есть отказ «позиция утрачена» в штатной работе не производится ничем.

Почему не длиннее

Журнал НЕ является историей выдачи. История — предмет журнала аудита и журнала намерений (`kaname.audit_outbox`, `kaname.fga_outbox`), чьё удержание объявлено политикой намеренно. Держать здесь дольше значило бы завести третье свидетельство одного события — свидетельство без читателя.

ЧЕГО величина не обещает

Она не есть окно отзыва доступа: окно задаёт ПЕРИОД ПЕРЕПРОСА читателя и объявлено политикой (`corelib/authz`.RevocationPolicy). Отставший дальше порога читатель получает не «пропущенный отзыв», а ЯВНЫЙ отказ с возобновимой позицией — и отвечает на него сплошным гашением кэша и закрытием всех потоков ([Watcher.positionLost]). Это дороже точечного закрытия, но не менее безопасно: отзыв применяется ШИРЕ нужного, а не уже.

У порога НЕТ слагаемого на разницу часов

Отметка времени строки ставится умолчанием колонки (`created_at DEFAULT now()`), то есть часами БАЗЫ, и уборка судит теми же часами: момент времени в сигнатуру уборщика не входит, предикат целиком в SQL. Слагаемое, покрывающее разницу, которой нет, было бы запасом без предмета.

View Source
const ReasonPositionLost = "SUBJECT_CHANGE_POSITION_LOST"

ReasonPositionLost — машинный признак полосы «позиция больше не возобновима».

Клиент ключуется на признак, а не разбирает прозу сообщения: тон сообщения — часть контракта, но не его машинная часть.

Variables

This section is empty.

Functions

func PositionLost

func PositionLost(earliestResumable int64) error

PositionLost собирает отказ владельца журнала.

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

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

Types

type Config

type Config struct {
	// Poller — источник изменений субъекта. Обязателен.
	Poller Poller
	// Flush — сброс кэша решений этой реплики. Обязателен.
	Flush func()
	// Interval — период перепроса. Неположительный резолвится в 2s.
	Interval time.Duration

	// Closer — реестр открытых потоков. ОБЯЗАТЕЛЕН.
	//
	// Необязательным он быть не может, и это выяснилось инъекцией: передай в
	// точке сборки ноль — и весь корпус проб потребителя остаётся зелёным, код
	// собирается, а отзыв перестаёт доезжать до потоков совсем. То есть
	// несделанная провязка была бы НЕОТЛИЧИМА от сделанной, и неотличима именно
	// тем механизмом, ради которого задача заведена.
	//
	// Потребитель проекцию строит безусловно, поэтому законного нуля здесь нет.
	Closer StreamCloser

	// StaleAfter — сколько потребитель вправе ДЕРЖАТЬ потоки, не подтвердив
	// чтения отзыва. Обязателен вместе с [Config.Closer].
	//
	// Это ОКНО FAIL-CLOSED, и оно есть следствие решения, а не срока жизни
	// соединения: величина отсчитывается от последнего удачного перепроса и к
	// длительности соединения отношения не имеет. Обязана превосходить
	// [Config.Interval] — иначе один пропущенный перепрос объявляется аварией.
	StaleAfter time.Duration

	// Now — часы. Ноль резолвится в [time.Now]; проба подставляет свои, потому
	// что срок обязан быть свойством решения, а не занятости машины.
	Now func() time.Time

	// Logger — журнал процесса. Обязателен.
	Logger *slog.Logger
}

Config — что приносит композиционный корень потребителя.

type Poller

type Poller interface {
	PollSubjectChanges(ctx context.Context, since int64, limit int32) (changes []SubjectChange, headID int64, err error)
}

Poller — узкий порт над чтением журнала. Реализуется Reader; в пробах — чем угодно, потому что предмет петли не транспорт, а курсор и решение гасить.

Предел порции задаёт ВЫЗЫВАЮЩИЙ, а не адаптер: по нему же вызывающий решает, могла ли порция быть усечена. Оставь предел у адаптера — и решение об усечении принималось бы здесь по числу, объявленному в другом месте, то есть двумя местами об одном предмете.

type PositionLostError

type PositionLostError struct {
	// EarliestResumable — нижняя позиция, с которой возобновление ещё ничего не
	// теряет: «самая ранняя удержанная строка минус один», а у вычищенного
	// целиком журнала — сама граница устоявшегося.
	EarliestResumable int64
}

PositionLostError — разобранный отказ: с какой позиции читателю СЕСТЬ.

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

func AsPositionLost

func AsPositionLost(err error) (*PositionLostError, bool)

AsPositionLost разбирает отказ владельца.

Ключуется на ПРИЗНАК, а не на код

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

Отказ БЕЗ позиции полосой не является

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

func (*PositionLostError) Error

func (e *PositionLostError) Error() string

type Reader

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

Reader — адаптер над порождённым клиентом, исполняющий Poller.

Единственное место пакета, говорящее по gRPC: остальное — курсор и решения «гасить» и «кого закрыть», и они не должны знать транспорта.

func NewReader

func NewReader(cc grpc.ClientConnInterface) *Reader

NewReader навешивает адаптер на УЖЕ ОТКРЫТОЕ соединение к внутреннему слушателю владельца прав. Своего соединения не открывает: адрес владельца потребитель объявляет один раз, и второе объявление того же адреса разошлось бы с первым молча.

func (*Reader) PollSubjectChanges

func (p *Reader) PollSubjectChanges(
	ctx context.Context, since int64, limit int32,
) ([]SubjectChange, int64, error)

PollSubjectChanges читает журнал владельца с позиции since с пределом, который назвал ВЫЗЫВАЮЩИЙ, отдавая изменения и голову журнала.

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

Почему имя субъекта едет дальше, а не остаётся здесь

Прежде адаптер оставлял от строки ОДИН номер и выбрасывал имя субъекта. Для сброса кэша этого хватало: он и так сбрасывался целиком. Для закрытия открытого потока (kacho#1022) — нет: закрыть можно только НАЗВАННОГО, а голый идентификатор не совпадает с ключом реестра ни при каких условиях. То есть отзыв не имел бы действия на длинных соединениях вовсе.

type StreamCloser

type StreamCloser interface {
	// CloseSubject закрывает потоки названного субъекта и возвращает их число.
	CloseSubject(subject string) int
	// CloseAll закрывает ВСЕ открытые потоки и возвращает их число.
	CloseAll() int
}

StreamCloser — реестр открытых длинных соединений потребителя.

Реализуется проекцией потока (`gateway/internal/subscriptionstream`.Handler). Порт объявлен ЗДЕСЬ, а реализация живёт у потребителя: фундамент не знает, что именно потребитель держит открытым, и знать не должен.

type SubjectChange

type SubjectChange struct {
	// ID — номер строки журнала; им двигается курсор.
	ID int64
	// Subject — субъект модели прав («user:usr-x»). Может быть пуст, если
	// владелец журнала его не назвал; такая строка двигает курсор и никого не
	// закрывает.
	Subject string
	// Naming — ПОЧЕМУ [SubjectChange.Subject] пуст.
	//
	// Пустое имя приходит по двум несравнимым причинам, и наблюдение обязано
	// их различать: выдача ГРУППЕ безымянна по устройству продукта (норма,
	// при каждом снятии привязки в предпочтительной форме выдачи), а
	// потерянный производителем тип означает, что отзыв по этой строке не
	// доедет ни до кэша поимённо, ни до открытого потока (дефект). Величина, в
	// которую сложены обе, ненулевая в штатной работе — тревогу на неё не
	// повесить (kacho#1463).
	//
	// Читается ТОЛЬКО при пустом имени. Ноль значения —
	// [authz.SubjectUnnameable]: строка, собранная без разбора, попадает в
	// громкую корзину, а не в тихую.
	Naming authz.SubjectNaming
}

SubjectChange — одна строка `subject_change_outbox` в том объёме, в каком её ЧИТАЕТ потребитель.

Вид события (`op`) сюда намеренно не переносится: потребитель его не читает, а поле, принятое и не прочитанное, обещает вызывающему поведение, которого нет.

type Watcher

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

Watcher читает журнал смены субъекта, гасит кэш решений своего процесса и закрывает открытые потоки названных субъектов.

func New

func New(cfg Config) (*Watcher, error)

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

Отказ на сборке, а не первым отказом в бою: величина, при которой fail-closed не наступает никогда, — ошибка посадки, и обнаруживать её тогда, когда она уже кому-то стоила доступа, поздно.

func (*Watcher) ClosesStreams

func (w *Watcher) ClosesStreams() bool

ClosesStreams — провязан ли реестр открытых потоков.

Существует ради САМООТЧЁТА: он обязан печатать НАБЛЮДЕНИЕ, а не литерал. Литерал продолжает утверждать «закрывает» при отключённом закрывателе — то есть ровно тот класс, ради которого задача заведена.

func (*Watcher) Poll

func (w *Watcher) Poll(ctx context.Context)

Poll исполняет РОВНО ОДИН перепрос: прочитать, затем принять голову за курсор ЛИБО двинуть курсор по прочитанному, погасить кэш и закрыть названные потоки.

Цикл синхронный: когда вызов вернулся, всё, что этот перепрос собирался сделать, сделано. Отсюда два следствия, и оба намеренные.

ПЕРВОЕ — метод публичный. Он и есть та единица работы, которую Watcher.Run повторяет по расписанию; потребителю, ведущему собственный такт, незачем поднимать вторую петлю рядом. Тест-только-шов здесь стоял бы ровно на этом же месте и означал бы то же самое, только не будучи назван.

ВТОРОЕ — им ставится СКВОЗНОЙ вопрос. Проба, спрашивающая «изменение прав доехало ли до кэша решений и до открытого потока», обязана читать наблюдаемое ПОСЛЕ того, как перепрос закончился, а не пока он идёт. Через Watcher.Run это означало бы угадывать момент по таймеру — угадывание верное на свободной машине и неверное на занятой.

func (*Watcher) Run

func (w *Watcher) Run(ctx context.Context)

Run блокирует до отмены контекста. Зовётся в своей горутине.

РЕПЛИКИ: на-реплику — петля гасит кэш СВОЕГО процесса, закрывает СВОИ потоки и держит курсор в памяти. Каждая реплика обязана читать сама: разведи её выбором одной — и кэш невыбранных не погаснет вовсе, а их потоки переживут отзыв. Дубль чтения безвреден не по намерению, а по свойству оператора: чтение журнала ничего не меняет, гашение кэша идемпотентно, а закрытие уже закрытого потока стоит нуля.

func (*Watcher) StaleAfter

func (w *Watcher) StaleAfter() time.Duration

StaleAfter — объявленный срок неподтверждённого чтения отзыва. Для самоотчёта.

Jump to

Keyboard shortcuts

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