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())) }