Files
healthlog/internal/replay/worker.go
T
av b278501a6e store: при равной полноте точек побеждает пришедшая доставка
- байтовый порядок канонических форм остался тай-брейком только внутри одной
  доставки: на живом корпусе он решал 98,8% спорных координат и системно хранил
  меньшее значение, из-за чего step_count терял род и verify:archive был красным
- правило перестало быть коммутативным осознанно, поэтому порядок свёртки
  приведён к журнальному: проход воркера прекращается на отложенной доставке,
  а свёртка вне порядка журнала пишет WARN
- заведены счётчики PointsHeld и PointsErased — удержание полнотой и
  единственное направление, в котором правило теряет содержание
2026-08-04 11:16:24 +03:00

315 lines
17 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package replay
import (
"context"
"log/slog"
"time"
"git.vakhrushev.me/av/healthlog/internal/fold"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// batchSize — сколько неразобранных доставок берётся одним запросом.
//
// Не ради страниц, а ради памяти: задолженность после миграции, переводящей
// строки в `pending`, равна всему архиву, и материализовать её целиком незачем —
// проход всё равно идёт по одной.
const batchSize = 256
// foldTimeout — сколько отводится свёртке одной доставки.
//
// Свёртка идёт на контексте, отвязанном от остановки, поэтому собственный
// дедлайн обязателен: без него зависшая запись держала бы единственного воркера
// до конца жизни процесса, и очередь перестала бы двигаться вовсе.
const foldTimeout = 2 * time.Minute
// tickInterval — как часто воркер просыпается сам, без сигнала.
//
// Сигнал приносит приём, и для свежей доставки его достаточно. Тик закрывает
// два случая, которых сигнал не закрывает: доставка, оставшаяся в очереди из-за
// занятости базы, иначе ждала бы СЛЕДУЮЩЕЙ доставки (а ночью телефон молчит
// часами), и метка отставания иначе не вычислялась бы вовсе — «работа есть,
// прогресса нет» было бы неотличимо от здорового пустого потока.
const tickInterval = time.Minute
// lagThreshold — с какого ожидания доставка считается задержанной.
//
// Период быстрого прохода синхронизации: если доставка ждала дольше, чем
// интервал между доставками, очередь растёт, а не рассасывается.
const lagThreshold = 5 * time.Minute
// Worker — фоновая свёртка принятых доставок.
//
// Очередь — сама таблица: доставка ждёт свёртки в статусе `pending`, а канал
// несёт только бит «есть работа». Отсюда три свойства, ради которых так и
// сделано: переполнять нечего, падение процесса очереди не теряет, а подбор
// неразобранного при старте не является отдельным кодом — это обычный проход.
type Worker struct {
store *store.Store
player Player
log *slog.Logger
// wake — сигнал «есть работа», ёмкость 1 и неблокирующая отправка. Та же
// форма, что у os/signal.Notify: сигнал ничего не несёт, и потерять лишний
// не только можно, но и нужно.
wake chan struct{}
// startupDone — первый проход завершён.
//
// До него метка отставания молчит: задолженность, накопленная ДО старта,
// ждала не воркера, а его появления, и сотня одинаковых WARN при первом же
// запуске обесценила бы уровень.
startupDone bool
// fold — свёртка одной доставки. Полем, а не прямым вызовом, ради ОДНОГО
// шва: исход `Deferred` наступает только от занятости базы, а удержать её
// по-настоящему стоит двадцать пять секунд (`busy_timeout` 5 с × 5 повторов
// транзакции) — столько же, сколько стоит `task verify:busy`, который в
// гейт намеренно не входит. Барьер журнального порядка при этом и есть
// механизм, на котором держится равенство «пересборка = приём», и оставить
// его без дешёвой проверки значило бы проверять его только тем прогоном,
// который никто не гоняет по расписанию.
fold func(ctx context.Context, deliveryID string) Outcome
}
// NewWorker собирает воркер над рабочей базой.
func NewWorker(st *store.Store, f *fold.Service, log *slog.Logger) *Worker {
w := &Worker{
store: st,
player: Player{Fold: f},
log: log.With("capability", "fold-worker"),
wake: make(chan struct{}, 1),
}
w.fold = w.foldOne
return w
}
// Notify будит воркер. Вызывается приёмом после того, как доставка учтена.
//
// Потеря сигнала отказом не является: доставка от этого не перестаёт числиться
// `pending`, и её подберёт следующий сигнал, тик или старт.
func (w *Worker) Notify() {
select {
case w.wake <- struct{}{}:
default:
}
}
// Run ведёт воркер до отмены контекста.
//
// Отмена проверяется МЕЖДУ доставками: свёртка идёт на отвязанном контексте и
// рваться не должна. Обещания «текущая доставка непременно досворачивается» тут
// нет — бюджет остановки меньше бюджета свёртки; гарантируется другое: после
// выхода не существует доставки, которая числится разобранной, а записана
// наполовину.
//
// Первый проход делается сразу, без ожидания сигнала: он и есть подбор
// неразобранного при старте.
func (w *Worker) Run(ctx context.Context) {
if n, err := w.store.CountPendingDeliveries(ctx); err != nil {
if ctx.Err() == nil {
w.log.ErrorContext(ctx, "pending backlog not counted", "error", err)
}
} else if n > 0 {
// Размер задолженности — ответ на вопрос «что сервис будет делать
// первые минуты после рестарта». Одной строкой и один раз.
w.log.InfoContext(ctx, "pending backlog at start", "deliveries", n)
}
ticker := time.NewTicker(tickInterval)
defer ticker.Stop()
for {
if _, err := w.Pass(ctx); err != nil && ctx.Err() == nil {
// Отказ прохода не убивает цикл: воркер, умерший от временного
// отказа базы, остановил бы свёртку до конца жизни процесса, пока
// приём продолжал бы отвечать 200.
//
// Отмена сюда не попадает: штатная остановка не отказ, а ERROR о
// ней обесценил бы уровень, по которому вмешиваются.
w.log.ErrorContext(ctx, "fold pass failed", "error", err)
}
select {
case <-ctx.Done():
return
case <-w.wake:
case <-ticker.C:
}
}
}
// Pass делает один проход по очереди и возвращает его исход.
//
// Синхронный шов: тесты зовут его напрямую и не ждут по часам. Без него
// проверки «все свёрнуты», «проход конечен», «метка не сработала на первом
// проходе» писались бы опросом базы с таймаутом.
//
// Курсор строго возрастает, и это нужно не ради страниц, а ради завершимости:
// доставка, у которой не удалось записать даже исход разбора, остаётся
// `pending`, и проход без курсора выбирал бы её бесконечно.
func (w *Worker) Pass(ctx context.Context) (Outcome, error) {
var total Outcome
var cursor store.PendingDelivery
var lag lagged
for {
if ctx.Err() != nil {
return total, nil
}
batch, err := w.store.PendingDeliveries(ctx, cursor, batchSize)
if err != nil {
return total, err
}
if len(batch) == 0 {
// Флаг снимается ТОЛЬКО здесь — у прохода, дошедшего до пустой
// выборки. Взведённый на любом выходе (отказ базы, отмена), он
// включал бы метку задержки после прохода, который ничего не
// свернул, и следующий проход выдал бы WARN на всю задолженность —
// ровно тот шквал, против которого метка и подавляется при старте.
//
// Порядок двух строк существен: на задолженности ПЕРВОГО прохода
// метка молчит — та ждала не воркера, а его появления.
w.warnLag(ctx, lag)
w.startupDone = true
return total, nil
}
for _, d := range batch {
if ctx.Err() != nil {
return total, nil
}
lag.add(d)
// Порядок спрашивается ДО свёртки, а говорится ПОСЛЕ и только если
// свёртка состоялась: отложенная доставка витрину не трогала, и
// запись о свёртке вне порядка утверждала бы событие, которого не
// было, — да ещё повторялась бы каждым проходом, пока голова очереди
// занята. Спросить надо всё же до: после свёртки предикат уже видит
// саму эту доставку разобранной.
outOfOrder := w.outOfOrder(ctx, d)
out := w.fold(ctx, d.ID)
total.Add(out)
if outOfOrder && out.Deferred == 0 {
w.warnOutOfOrder(ctx, d)
}
if out.Deferred > 0 {
// Отложенная доставка ДЕРЖИТ очередь: перешагнув её, проход
// свернул бы её преемниц раньше неё, а тай-брейк равной полноты
// разрешается в пользу пришедшей доставки — то есть исход
// слияния стал бы функцией порядка свёртки, отличного от
// порядка журнала, и живая витрина разошлась бы с пересборкой
// молча, в значениях точек.
//
// Очередь от этого не встаёт: отложенным считается ТОЛЬКО
// занятость базы и отмена снаружи (store.Transient), а
// собственный дедлайн свёртки в него намеренно не входит —
// доставка, не уложившаяся в бюджет, получает `failed` и
// очередь освобождает. Занятость же блокирует запись всем
// одинаково: проход, перешагнувший занятую доставку, упёрся бы
// в ту же занятость на следующей.
//
// Метка отставания пишется ЗДЕСЬ, а не только на пустой
// выборке. Иначе барьер отменял бы её ровно в том состоянии,
// ради которого она заведена: занятая голова очереди не
// пропускает проход до пустой выборки никогда, и «очередь стоит
// два часа» стало бы неотличимо от «споткнулась один раз».
//
// Флаг первого прохода при этом НЕ взводится: проход до конца
// очереди не дошёл, и объявлять задолженность разобранной рано.
w.warnLag(ctx, lag)
return total, nil
}
cursor = d
}
}
}
// lagged копит отставание прохода: сколько доставок ждали свёртки и дольше всех
// ждала какая.
//
// Считается на ВЫБОРКЕ, а не по факту успешной свёртки: иначе застрявшая
// доставка молчала бы ровно в том состоянии, ради которого метка и заведена.
type lagged struct {
count int
worst time.Duration
worstID string
}
func (l *lagged) add(d store.PendingDelivery) {
waited := store.Now().Sub(d.ReceivedAt)
if waited < lagThreshold {
return
}
l.count++
if waited > l.worst {
l.worst = waited
l.worstID = d.ID
}
}
// foldOne сворачивает доставку на контексте, ОТВЯЗАННОМ от остановки.
//
// Отмена снаружи не должна превращаться в свойство доставки: свёртка пишет
// исход на переживающем отмену контексте, и оборванная на середине пометила бы
// доставку так, что воркер её больше не подберёт. Собственный дедлайн при этом
// остаётся и означает именно отказ доставки.
func (w *Worker) foldOne(ctx context.Context, deliveryID string) Outcome {
foldCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), foldTimeout)
defer cancel()
// Ошибку не возвращаем: она уже записана свёрткой в лог и в parse_status,
// а отказ одной доставки прохода не прекращает. Паника тоже: её
// перехватывает сама свёртка — там же, где живёт единственный писатель
// исхода разбора.
out, _ := w.player.Play(foldCtx, deliveryID)
return out
}
// warnOutOfOrder называет свёртку, идущую вне порядка журнала.
//
// Остаточное окно закрыть здесь нечем: метка `received_at` фиксируется при
// выпуске ULID, а строка учёта становится видимой только после записи тела
// (измерено 184 мс на 62 МиБ), поэтому при конкурентном приёме доставка с более
// ранней меткой появляется после того, как её преемница уже свёрнута. Барьер
// выше этого не ловит — отставшей доставки в момент прохода не существует.
//
// Молчать об этом нельзя: правило слияния точек стало функцией порядка свёртки,
// и перестановка оставляет в витрине версию не той доставки, что стоит в
// журнале последней. Лечится `healthlog reindex`; без этой строки узнать о
// необходимости неоткуда.
//
// Одна запись на доставку: событие редкое и адресное, шквала здесь не бывает.
// Отказ запроса свёртку НЕ прекращает — наблюдение не может стоить доставки.
func (w *Worker) outOfOrder(ctx context.Context, d store.PendingDelivery) bool {
later, err := w.store.ParsedAfter(ctx, d.ReceivedAt, d.ID)
if err != nil {
if ctx.Err() == nil {
w.log.ErrorContext(ctx, "journal order check failed", "error", err, "delivery_id", d.ID)
}
return false
}
return later
}
func (w *Worker) warnOutOfOrder(ctx context.Context, d store.PendingDelivery) {
w.log.WarnContext(ctx, "delivery folded out of journal order",
"delivery_id", d.ID,
"waited_sec", int64(store.Now().Sub(d.ReceivedAt).Seconds()))
}
// warnLag называет отставание одной строкой на проход.
//
// Одной, а не по строке на доставку: задолженность в сотню тел давала бы сотню
// одинаковых WARN каждую минуту, и уровень, по которому вмешиваются, перестал
// бы что-либо значить.
func (w *Worker) warnLag(ctx context.Context, l lagged) {
if !w.startupDone || l.count == 0 {
return
}
w.log.WarnContext(ctx, "deliveries waited for fold",
"deliveries", l.count,
"worst_delivery_id", l.worstID,
"worst_waited_sec", int64(l.worst.Seconds()))
}