Приём отвечает 200 до свёртки, свёртку ведёт фоновый воркер
- Очередью служит сама таблица: доставка ждёт свёртки в статусе `pending`, канал несёт только бит «есть работа». Переполнять нечего, падение процесса очередь не теряет, а подбор `pending` при старте — обычный проход воркера, а не отдельный код. Классификация исхода общая с пересборкой журнала. - Исход разбора начал отражать доставку, а не обстоятельства: отмена и занятость базы статус не меняют (иначе конкуренция за базу выводила бы доставку из очереди навсегда), паника свёртки больше не валит процесс, а учёт доставки идёт через транзакцию с повторами. - Длинный бюджет ответа выдан маршруту приёма, а не всему серверу: `write_timeout` в Go покрывает и чтение тела, и общий подъём снял бы защиту с остальных маршрутов.
This commit is contained in:
@@ -0,0 +1,67 @@
|
||||
package replay
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"testing"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/fold"
|
||||
"git.vakhrushev.me/av/healthlog/internal/hae"
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
// Классификация исхода — та половина, которую пересборка и фоновый воркер
|
||||
// обязаны делить. Проверяется перебором классов, без базы и без архива: второй
|
||||
// классификатор разошёлся бы с первым молча, а по счётчику `partial`
|
||||
// принимается решение о судьбе тела в архиве.
|
||||
func TestClassifyРазводитИсходыПоКлассам(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
cases := []struct {
|
||||
name string
|
||||
err error
|
||||
want Outcome
|
||||
}{
|
||||
{"успех", nil, Outcome{Folded: 1}},
|
||||
{"слой не выведен", hae.ErrLayerUnknown, Outcome{FailedLayer: 1}},
|
||||
{"слой не выведен, обёрнут", fmt.Errorf("свёртка: %w", hae.ErrLayerUnknown), Outcome{FailedLayer: 1}},
|
||||
{"содержимое не разбирается", hae.ErrMalformed, Outcome{FailedMalformed: 1}},
|
||||
{"база занята", store.ErrBusy, Outcome{Deferred: 1}},
|
||||
{"база занята, обёрнута", fmt.Errorf("слияние: %w", store.ErrBusy), Outcome{Deferred: 1}},
|
||||
{"работу прекратили снаружи", context.Canceled, Outcome{Deferred: 1}},
|
||||
// Дедлайн — свойство доставки, а не обстоятельств: она не уложится в
|
||||
// бюджет и в следующий раз, а повтор безнадёжного останавливает очередь.
|
||||
{"свёртка не уложилась в бюджет", context.DeadlineExceeded, Outcome{FailedOther: 1}},
|
||||
{"прочее", errors.New("диск отвалился"), Outcome{FailedOther: 1}},
|
||||
// Паника — дефект нашего кода, а не обстоятельство: доставка выводится
|
||||
// из очереди, тело ждёт пересборки.
|
||||
{"свёртка паниковала", fold.ErrPanicked, Outcome{FailedOther: 1}},
|
||||
}
|
||||
|
||||
for _, c := range cases {
|
||||
t.Run(c.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if got := classify(c.err); got != c.want {
|
||||
t.Errorf("Classify(%v) = %+v, ожидалось %+v", c.err, got, c.want)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Накопление — сложение по классам: у пересборки и у воркера один набор имён
|
||||
// для одних исходов.
|
||||
func TestOutcomeAddСкладываетПоКлассам(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var total Outcome
|
||||
total.Add(Outcome{Folded: 1, Partial: 1})
|
||||
total.Add(Outcome{Folded: 1, Incomparable: 2})
|
||||
total.Add(Outcome{Deferred: 1})
|
||||
|
||||
want := Outcome{Folded: 2, Deferred: 1, Partial: 1, Incomparable: 2}
|
||||
if total != want {
|
||||
t.Errorf("сумма %+v, ожидалась %+v", total, want)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package replay
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/fold"
|
||||
"git.vakhrushev.me/av/healthlog/internal/hae"
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
// Outcome — исход свёртки: одной доставки или их последовательности.
|
||||
//
|
||||
// Классы разведены потому, что читаются по-разному. FailedLayer — штатный исход
|
||||
// (слой не выводится, таких тел в журнале заведомо есть), FailedMalformed —
|
||||
// содержимое не разбирается, Deferred — работа не сделана по обстоятельствам,
|
||||
// и доставка осталась в очереди. Только FailedOther означает, что что-то не так
|
||||
// с самой свёрткой. Один общий счётчик отправлял бы человека искать дефект там,
|
||||
// где его нет.
|
||||
type Outcome struct {
|
||||
Folded int
|
||||
FailedLayer int
|
||||
FailedMalformed int
|
||||
// Deferred — доставка осталась `pending`: отмена или занятость базы. Не
|
||||
// отказ доставки, а несделанная работа; её подберёт следующий проход.
|
||||
Deferred int
|
||||
FailedOther int
|
||||
// Partial — доставок, в теле которых остались непокрытые разбором секции.
|
||||
// Не отклонение, а половина потока; названо потому, что именно эти тела
|
||||
// ретеншену трогать нельзя.
|
||||
Partial int
|
||||
// Incomparable — столкновений с несравнимыми наборами полей. На живом потоке
|
||||
// их не было ни разу, и на этом стоит отказ от объединения полей.
|
||||
Incomparable int
|
||||
}
|
||||
|
||||
// Add накапливает исход одной доставки в общий.
|
||||
func (o *Outcome) Add(other Outcome) {
|
||||
o.Folded += other.Folded
|
||||
o.FailedLayer += other.FailedLayer
|
||||
o.FailedMalformed += other.FailedMalformed
|
||||
o.Deferred += other.Deferred
|
||||
o.FailedOther += other.FailedOther
|
||||
o.Partial += other.Partial
|
||||
o.Incomparable += other.Incomparable
|
||||
}
|
||||
|
||||
// classify раскладывает ошибку свёртки по классам исхода.
|
||||
//
|
||||
// Чистая функция, и это не украшение: она и есть та половина, которую задача
|
||||
// требовала не дублировать между пересборкой и фоновым воркером, — а
|
||||
// проверяется она перебором классов, без базы и без архива.
|
||||
//
|
||||
// Неэкспортируемая намеренно: её результат содержит поля `Partial` и
|
||||
// `Incomparable`, которые дописывает только Play, — вторая публичная дверь
|
||||
// молча занижала бы именно тот счётчик, по которому принимается решение о
|
||||
// судьбе тела в архиве.
|
||||
func classify(err error) Outcome {
|
||||
var out Outcome
|
||||
switch {
|
||||
case err == nil:
|
||||
out.Folded++
|
||||
case store.Transient(err):
|
||||
// Статус доставки свёртка в этих случаях не трогает: она осталась
|
||||
// `pending` и будет свёрнута снова. Правило одно на обоих — то, по
|
||||
// которому свёртка решает не писать исход.
|
||||
out.Deferred++
|
||||
case errors.Is(err, hae.ErrLayerUnknown):
|
||||
out.FailedLayer++
|
||||
case errors.Is(err, hae.ErrMalformed):
|
||||
out.FailedMalformed++
|
||||
default:
|
||||
out.FailedOther++
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// Player сворачивает доставку по идентификатору и классифицирует исход.
|
||||
//
|
||||
// Общий и для пересборки журнала, и для фонового воркера приёма — второй
|
||||
// классификатор разошёлся бы с первым молча, а по одному из его счётчиков
|
||||
// (`Partial`) принимается решение о судьбе тела в архиве.
|
||||
type Player struct {
|
||||
Fold *fold.Service
|
||||
}
|
||||
|
||||
// Play сворачивает одну доставку и возвращает её исход.
|
||||
//
|
||||
// Классифицируется ТОЛЬКО ошибка свёртки: на контекст Play не смотрит, и это
|
||||
// существенно. У двух вызывающих отменённый контекст означает противоположное —
|
||||
// у пересборки в свёртку уходит тот же отменяемый контекст («нас остановили»),
|
||||
// у воркера отвязанный от остановки, с собственным дедлайном («доставка не
|
||||
// уложилась в бюджет»). Решение «работу прекратили снаружи» принимает цикл,
|
||||
// каждый по своему контексту.
|
||||
func (p Player) Play(ctx context.Context, deliveryID string) (Outcome, error) {
|
||||
st, err := p.Fold.Fold(ctx, deliveryID)
|
||||
out := classify(err)
|
||||
|
||||
if err == nil {
|
||||
// Счётчики читаются только у успешной свёртки: при ошибке поля Stats
|
||||
// заполнены частично (Uncovered у отказавшего разбора всегда пуст, хотя
|
||||
// в базу список записан) — и Partial молча занижался бы. А по нему
|
||||
// принимается решение о ретеншене тел.
|
||||
if len(st.Uncovered) > 0 {
|
||||
out.Partial++
|
||||
}
|
||||
out.Incomparable += st.Incomparable
|
||||
}
|
||||
return out, err
|
||||
}
|
||||
+15
-46
@@ -24,7 +24,6 @@ import (
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/archive"
|
||||
"git.vakhrushev.me/av/healthlog/internal/fold"
|
||||
"git.vakhrushev.me/av/healthlog/internal/hae"
|
||||
"git.vakhrushev.me/av/healthlog/internal/ident"
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
@@ -70,24 +69,11 @@ type Report struct {
|
||||
// Orphans — строк учёта, у которых тела в архиве нет. Станет штатным, когда
|
||||
// появится ретеншен архива.
|
||||
Orphans int
|
||||
// Folded — сколько доставок свернулось.
|
||||
Folded int
|
||||
// Отказы разведены по классам, потому что читаются они по-разному.
|
||||
// FailedLayer — слой не выводится: штатный исход, таких доставок в журнале
|
||||
// заведомо есть. FailedMalformed — содержимое не разбирается. FailedOther —
|
||||
// всё прочее (тело не читается, отказ базы); только оно означает, что с
|
||||
// пересборкой что-то не так. Один общий счётчик отправлял бы человека
|
||||
// искать дефект там, где его нет.
|
||||
FailedLayer int
|
||||
FailedMalformed int
|
||||
FailedOther int
|
||||
// Partial — доставок, в теле которых остались непокрытые разбором секции.
|
||||
// Не отклонение, а половина потока; названо потому, что именно эти тела
|
||||
// ретеншену трогать нельзя.
|
||||
Partial int
|
||||
// Incomparable — столкновений с несравнимыми наборами полей. На живом потоке
|
||||
// их не было ни разу, и на этом стоит отказ от объединения полей.
|
||||
Incomparable int
|
||||
|
||||
// Outcome — счётчики свёртки, те же самые, что считает фоновый воркер
|
||||
// приёма. Встроены, а не продублированы именами: два набора имён для одних
|
||||
// исходов разошлись бы при первой же правке классификации.
|
||||
Outcome
|
||||
|
||||
Buckets int64
|
||||
Fingerprint string
|
||||
@@ -148,6 +134,7 @@ func Run(ctx context.Context, o Options) (Report, error) {
|
||||
}
|
||||
}
|
||||
|
||||
player := Player{Fold: o.Fold}
|
||||
for i, d := range journal {
|
||||
if ctx.Err() != nil {
|
||||
rep.Canceled = true
|
||||
@@ -156,36 +143,17 @@ func Run(ctx context.Context, o Options) (Report, error) {
|
||||
if err := o.Target.CreateDelivery(ctx, d); err != nil {
|
||||
return stopOr(rep, err)
|
||||
}
|
||||
st, err := o.Fold.Fold(ctx, d.ID)
|
||||
switch {
|
||||
case err == nil:
|
||||
rep.Folded++
|
||||
case ctx.Err() != nil:
|
||||
// Отмена, застигшая свёртку, — не отказ доставки: считать её отказом
|
||||
// значило бы обвинить разбор в том, чего он не делал, и отправить
|
||||
// человека искать дефект по логу.
|
||||
out, _ := player.Play(ctx, d.ID)
|
||||
// Отмена, застигшая свёртку, — не отказ доставки: считать её отказом
|
||||
// значило бы обвинить разбор в том, чего он не делал, и отправить
|
||||
// человека искать дефект по логу. Решение принимает цикл по СВОЕМУ
|
||||
// контексту — тому же, на котором шла свёртка; классификатор о нём не
|
||||
// знает намеренно, у воркера тот же признак означает другое.
|
||||
if ctx.Err() != nil {
|
||||
rep.Canceled = true
|
||||
return rep, nil
|
||||
case errors.Is(err, hae.ErrLayerUnknown):
|
||||
// Штатный исход, уже записанный свёрткой в лог и в parse_status:
|
||||
// журнал заведомо содержит тела без плотных метрик. Останов на
|
||||
// первом лишил бы пересборки все остальные.
|
||||
rep.FailedLayer++
|
||||
case errors.Is(err, hae.ErrMalformed):
|
||||
rep.FailedMalformed++
|
||||
default:
|
||||
rep.FailedOther++
|
||||
}
|
||||
if err == nil {
|
||||
// Счётчики читаются только у успешной свёртки: при ошибке поля Stats
|
||||
// заполнены частично (Uncovered у отказавшего разбора всегда пуст,
|
||||
// хотя в базу список записан) — и Partial молча занижался бы. А по
|
||||
// нему принимается решение о ретеншене тел.
|
||||
if len(st.Uncovered) > 0 {
|
||||
rep.Partial++
|
||||
}
|
||||
rep.Incomparable += st.Incomparable
|
||||
}
|
||||
rep.Add(out)
|
||||
if o.Progress != nil {
|
||||
o.Progress(i+1, len(journal))
|
||||
}
|
||||
@@ -218,6 +186,7 @@ func Run(ctx context.Context, o Options) (Report, error) {
|
||||
"folded", rep.Folded,
|
||||
"failed_layer", rep.FailedLayer,
|
||||
"failed_malformed", rep.FailedMalformed,
|
||||
"deferred", rep.Deferred,
|
||||
"failed_other", rep.FailedOther,
|
||||
"partial", rep.Partial,
|
||||
"incomparable", rep.Incomparable,
|
||||
|
||||
@@ -0,0 +1,232 @@
|
||||
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()))
|
||||
}
|
||||
@@ -0,0 +1,373 @@
|
||||
package replay_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/archive"
|
||||
"git.vakhrushev.me/av/healthlog/internal/fold"
|
||||
"git.vakhrushev.me/av/healthlog/internal/ident"
|
||||
"git.vakhrushev.me/av/healthlog/internal/replay"
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
// newWorker собирает воркер над свежей базой и архивом, отдавая заодно то, чем
|
||||
// проверяют его следы в логе.
|
||||
func newWorker(t *testing.T, dir string) (*replay.Worker, *store.Store, *archive.Archive, *logSink) {
|
||||
t.Helper()
|
||||
|
||||
arch := openArchive(t, filepath.Join(dir, "raw"))
|
||||
st := openStore(t, filepath.Join(dir, "live.db"))
|
||||
sink := newLogSink()
|
||||
log := slog.New(slog.NewJSONHandler(sink, &slog.HandlerOptions{Level: slog.LevelDebug}))
|
||||
w := replay.NewWorker(st, fold.New(arch, st, 0, log), log)
|
||||
return w, st, arch, sink
|
||||
}
|
||||
|
||||
// logSink собирает записи лога, чтобы проверять их без гонок и без ожиданий по
|
||||
// часам: тест синхронизируется появлением строки, а не сном.
|
||||
type logSink struct {
|
||||
mu sync.Mutex
|
||||
lines []string
|
||||
watch map[string]*watcher
|
||||
}
|
||||
|
||||
// watcher ждёт n-го появления строки.
|
||||
type watcher struct {
|
||||
left int
|
||||
ch chan struct{}
|
||||
}
|
||||
|
||||
func newLogSink() *logSink { return &logSink{watch: map[string]*watcher{}} }
|
||||
|
||||
func (s *logSink) Write(p []byte) (int, error) {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
line := string(p)
|
||||
s.lines = append(s.lines, line)
|
||||
if w, ok := s.watch[msgOf(line)]; ok {
|
||||
w.left--
|
||||
if w.left == 0 {
|
||||
close(w.ch)
|
||||
delete(s.watch, msgOf(line))
|
||||
}
|
||||
}
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
// expect регистрирует ожидание n-го появления строки до того, как она может
|
||||
// появиться: синхронизация идёт событием, а не сном.
|
||||
func (s *logSink) expect(msg string, n int) <-chan struct{} {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
w := &watcher{left: n, ch: make(chan struct{})}
|
||||
s.watch[msg] = w
|
||||
return w.ch
|
||||
}
|
||||
|
||||
func msgOf(line string) string {
|
||||
var rec struct {
|
||||
Msg string `json:"msg"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &rec); err != nil {
|
||||
return ""
|
||||
}
|
||||
return rec.Msg
|
||||
}
|
||||
|
||||
// count считает записи с данным msg.
|
||||
func (s *logSink) count(msg string) int {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
|
||||
n := 0
|
||||
for _, line := range s.lines {
|
||||
if msgOf(line) == msg {
|
||||
n++
|
||||
}
|
||||
}
|
||||
return n
|
||||
}
|
||||
|
||||
func (s *logSink) dump() string {
|
||||
s.mu.Lock()
|
||||
defer s.mu.Unlock()
|
||||
return strings.Join(s.lines, "")
|
||||
}
|
||||
|
||||
// Подбор неразобранного — обычный проход воркера, а не отдельный режим: после
|
||||
// миграции 00005 неразобранными числятся все доставки архива, и подобрать их
|
||||
// сегодня может только пересборка с ручной подменой базы.
|
||||
func TestПроходПодбираетЗадолженность(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, arch, _ := newWorker(t, dir)
|
||||
ctx := context.Background()
|
||||
|
||||
items := journal(t, "minute.json", "hour.json", "raw.json")
|
||||
for _, it := range items {
|
||||
writeBody(t, arch, st, it, fixture(t, it.fixture))
|
||||
}
|
||||
|
||||
out, err := w.Pass(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("проход: %v", err)
|
||||
}
|
||||
if out.Folded != len(items) {
|
||||
t.Fatalf("свёрнуто %d из %d: %+v", out.Folded, len(items), out)
|
||||
}
|
||||
|
||||
n, err := st.CountPendingDeliveries(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("CountPendingDeliveries: %v", err)
|
||||
}
|
||||
if n != 0 {
|
||||
t.Errorf("неразобранными остались %d доставок", n)
|
||||
}
|
||||
|
||||
buckets, err := st.CountBuckets(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("CountBuckets: %v", err)
|
||||
}
|
||||
if buckets == 0 {
|
||||
t.Error("объектов 0: точки не доехали до хранилища")
|
||||
}
|
||||
}
|
||||
|
||||
// Порядок задаётся ЖУРНАЛОМ, а не порядком, в котором доставки попали в учёт.
|
||||
// Проверяется наблюдаемым следствием: доставка без плотных метрик наследует
|
||||
// слой предшествующей ей по `(received_at, id)`.
|
||||
//
|
||||
// Учёт заполняется в обратном хронологии порядке — так выглядит гонка двух
|
||||
// конкурентных приёмов, где поздняя доставка закоммитила строку первой.
|
||||
func TestПроходИдётВПорядкеЖурналаАНеВставки(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, arch, _ := newWorker(t, dir)
|
||||
ctx := context.Background()
|
||||
|
||||
base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
|
||||
minute := item{id: ident.NewID(), at: base.Add(1 * time.Second), automationID: "a", aggregation: "Default", fixture: "minute.json"}
|
||||
sleep := item{id: ident.NewID(), at: base.Add(2 * time.Second), automationID: "a", aggregation: "Default", fixture: "sparse_sleep.json"}
|
||||
|
||||
// Сначала учитывается ПОЗДНЯЯ доставка.
|
||||
writeBody(t, arch, st, sleep, fixture(t, sleep.fixture))
|
||||
writeBody(t, arch, st, minute, fixture(t, minute.fixture))
|
||||
|
||||
out, err := w.Pass(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("проход: %v", err)
|
||||
}
|
||||
if out.FailedLayer != 0 {
|
||||
t.Fatalf("слой не вывелся у %d доставок: проход пошёл в порядке вставки", out.FailedLayer)
|
||||
}
|
||||
|
||||
// Предшественник — минутная доставка, значит эпизоды сна легли в minute.
|
||||
hours, err := st.BucketHours(ctx, "sleep_analysis", "minute")
|
||||
if err != nil {
|
||||
t.Fatalf("часы объектов: %v", err)
|
||||
}
|
||||
if len(hours) == 0 {
|
||||
t.Error("эпизоды сна не унаследовали слой предшествующей доставки")
|
||||
}
|
||||
}
|
||||
|
||||
// Проход конечен и продвигается мимо доставки, которую свернуть не удалось:
|
||||
// курсор двигается вперёд независимо от исхода свёртки. Без этого доставка, у
|
||||
// которой не удалось записать даже исход разбора, выбиралась бы бесконечно.
|
||||
func TestПроходПродвигаетсяМимоНесворачиваемойДоставки(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, arch, _ := newWorker(t, dir)
|
||||
ctx := context.Background()
|
||||
|
||||
items := journal(t, "minute.json", "hour.json")
|
||||
for _, it := range items {
|
||||
writeBody(t, arch, st, it, fixture(t, it.fixture))
|
||||
}
|
||||
// У первой доставки тела больше нет — свернуть её нечем.
|
||||
if err := os.Remove(filepath.Join(arch.Root(), "2026", "08", "01", items[0].id+".json.gz")); err != nil {
|
||||
t.Fatalf("удаление тела: %v", err)
|
||||
}
|
||||
|
||||
done := make(chan replay.Outcome, 1)
|
||||
go func() {
|
||||
out, err := w.Pass(ctx)
|
||||
if err != nil {
|
||||
t.Errorf("проход: %v", err)
|
||||
}
|
||||
done <- out
|
||||
}()
|
||||
|
||||
var out replay.Outcome
|
||||
select {
|
||||
case out = <-done:
|
||||
case <-time.After(30 * time.Second):
|
||||
t.Fatal("проход не завершился: курсор не двигается")
|
||||
}
|
||||
if out.Folded != 1 || out.FailedOther != 1 {
|
||||
t.Fatalf("исход прохода %+v: ожидались одна свёрнутая и одна отказавшая", out)
|
||||
}
|
||||
|
||||
// Отказавшая доставка выбыла из очереди — иначе следующий проход брал бы её
|
||||
// снова и снова.
|
||||
n, err := st.CountPendingDeliveries(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("CountPendingDeliveries: %v", err)
|
||||
}
|
||||
if n != 0 {
|
||||
t.Errorf("неразобранными числятся %d доставок, ожидалось 0", n)
|
||||
}
|
||||
}
|
||||
|
||||
// Метка задержки молчит на задолженности первого прохода и говорит после него:
|
||||
// доставки, накопленные до старта, ждали не воркера, а его появления.
|
||||
func TestМеткаЗадержкиВключаетсяПослеПервогоПрохода(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, arch, sink := newWorker(t, dir)
|
||||
ctx := context.Background()
|
||||
|
||||
old := item{
|
||||
id: ident.NewID(), at: store.Now().Add(-time.Hour),
|
||||
automationID: "a", aggregation: "Minutes", fixture: "minute.json",
|
||||
}
|
||||
writeBody(t, arch, st, old, fixture(t, old.fixture))
|
||||
|
||||
if _, err := w.Pass(ctx); err != nil {
|
||||
t.Fatalf("первый проход: %v", err)
|
||||
}
|
||||
if n := sink.count("deliveries waited for fold"); n != 0 {
|
||||
t.Errorf("на задолженности первого прохода %d предупреждений о задержке:\n%s", n, sink.dump())
|
||||
}
|
||||
|
||||
late := item{
|
||||
id: ident.NewID(), at: store.Now().Add(-time.Hour),
|
||||
automationID: "a", aggregation: "Minutes", fixture: "hour.json",
|
||||
}
|
||||
writeBody(t, arch, st, late, fixture(t, late.fixture))
|
||||
|
||||
if _, err := w.Pass(ctx); err != nil {
|
||||
t.Fatalf("второй проход: %v", err)
|
||||
}
|
||||
if n := sink.count("deliveries waited for fold"); n != 1 {
|
||||
t.Errorf("предупреждений о задержке %d, ожидалось 1:\n%s", n, sink.dump())
|
||||
}
|
||||
}
|
||||
|
||||
// Отмена контекста завершает цикл — без ожиданий по часам: синхронизация идёт
|
||||
// возвратом Run, а не сном.
|
||||
func TestRunЗавершаетсяПоОтмене(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, _, _, sink := newWorker(t, dir)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
w.Run(ctx)
|
||||
}()
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(10 * time.Second):
|
||||
t.Fatalf("Run не вышел по отмене:\n%s", sink.dump())
|
||||
}
|
||||
}
|
||||
|
||||
// Задолженность при старте называется одной строкой: это ответ на вопрос «что
|
||||
// сервис будет делать первые минуты после рестарта».
|
||||
func TestRunНазываетЗадолженностьПриСтарте(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, arch, sink := newWorker(t, dir)
|
||||
|
||||
items := journal(t, "minute.json", "hour.json")
|
||||
for _, it := range items {
|
||||
writeBody(t, arch, st, it, fixture(t, it.fixture))
|
||||
}
|
||||
|
||||
// Ожидание регистрируется ДО запуска: тест синхронизируется появлением
|
||||
// строки, а не сном.
|
||||
said := sink.expect("pending backlog at start", 1)
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
w.Run(ctx)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-said:
|
||||
case <-time.After(30 * time.Second):
|
||||
t.Fatalf("строки о задолженности нет:\n%s", sink.dump())
|
||||
}
|
||||
cancel()
|
||||
<-done
|
||||
}
|
||||
|
||||
// Notify не блокирует и не копит: сигнал ничего не несёт, и лишний теряется
|
||||
// намеренно.
|
||||
func TestNotifyНеБлокирует(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
w, _, _, _ := newWorker(t, t.TempDir())
|
||||
for range 100 {
|
||||
w.Notify()
|
||||
}
|
||||
}
|
||||
|
||||
// Отказ прохода не убивает цикл: воркер, умерший от временного отказа базы,
|
||||
// остановил бы свёртку до конца жизни процесса, пока приём продолжал бы
|
||||
// отвечать 200.
|
||||
func TestRunПереживаетОтказПрохода(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
dir := t.TempDir()
|
||||
w, st, _, sink := newWorker(t, dir)
|
||||
|
||||
// База закрыта — выборка неразобранных отказывает на каждом проходе.
|
||||
if err := st.Close(); err != nil {
|
||||
t.Fatalf("закрытие базы: %v", err)
|
||||
}
|
||||
|
||||
// Второй отказ доказывает, что цикл пережил первый. Разбудить второй проход
|
||||
// без ожидания по часам может только сигнал: тик идёт раз в минуту.
|
||||
twice := sink.expect("fold pass failed", 2)
|
||||
w.Notify()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
w.Run(ctx)
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-twice:
|
||||
case <-time.After(30 * time.Second):
|
||||
t.Fatalf("цикл не пережил отказ прохода:\n%s", sink.dump())
|
||||
}
|
||||
cancel()
|
||||
<-done
|
||||
}
|
||||
Reference in New Issue
Block a user