Files
healthlog/internal/replay/worker.go
T
av 63bffe2865 Приём отвечает 200 до свёртки, свёртку ведёт фоновый воркер
- Очередью служит сама таблица: доставка ждёт свёртки в статусе `pending`,
  канал несёт только бит «есть работа». Переполнять нечего, падение процесса
  очередь не теряет, а подбор `pending` при старте — обычный проход воркера, а
  не отдельный код. Классификация исхода общая с пересборкой журнала.
- Исход разбора начал отражать доставку, а не обстоятельства: отмена и
  занятость базы статус не меняют (иначе конкуренция за базу выводила бы
  доставку из очереди навсегда), паника свёртки больше не валит процесс, а
  учёт доставки идёт через транзакцию с повторами.
- Длинный бюджет ответа выдан маршруту приёма, а не всему серверу:
  `write_timeout` в Go покрывает и чтение тела, и общий подъём снял бы защиту с
  остальных маршрутов.
2026-08-02 11:01:42 +03:00

233 lines
11 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
}
// 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()))
}