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 } // NewWorker собирает воркер над рабочей базой. func NewWorker(st *store.Store, f *fold.Service, log *slog.Logger) *Worker { return &Worker{ store: st, player: Player{Fold: f}, log: log.With("capability", "fold-worker"), wake: make(chan struct{}, 1), } } // 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) total.Add(w.foldOne(ctx, d.ID)) 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 } // 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())) }