package replay_test import ( "context" "log/slog" "path/filepath" "regexp" "sort" "strconv" "testing" "time" "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" ) // qtyPattern находит числовые значения точек в теле HAE. var qtyPattern = regexp.MustCompile(`("(?:qty|Avg|Min|Max)"\s*:\s*)(\d+(?:\.\d+)?)`) // bumped берёт РЕАЛЬНЫЙ пакет HAE и сдвигает в нём числовые значения, оставляя // координаты нетронутыми. Получается доставка, сталкивающаяся с исходной по // всем координатам и отличающаяся от неё только значениями — то есть ровно та // форма столкновения, ради которой изменение и делается (равные наборы полей, // разные значения; на живом корпусе это 98,8% спорных координат). // // Синтетическое тело писать нельзя: разбор формата HAE проверяется на реальных // пакетах — документация формата тонка и местами расходится с потоком. Здесь // сохраняется и то и другое: пакет настоящий, меняются только числа. func bumped(t *testing.T, name string, delta int) []byte { t.Helper() body := fixture(t, name) out := qtyPattern.ReplaceAllFunc(body, func(m []byte) []byte { g := qtyPattern.FindSubmatch(m) v, err := strconv.ParseFloat(string(g[2]), 64) if err != nil { return m } return append(append([]byte{}, g[1]...), strconv.FormatFloat(v+float64(delta), 'f', -1, 64)...) }) if string(out) == string(body) { t.Fatalf("фикстура %s не содержит числовых значений — столкновение не построить", name) } return out } // collidingJournal собирает журнал, где каждая следующая доставка приносит те // же координаты с другими значениями. func collidingJournal(t *testing.T, n int) []item { t.Helper() ids := make([]string, 0, n) for range n { ids = append(ids, ident.NewID()) } sort.Strings(ids) base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC) out := make([]item, 0, n) for i, id := range ids { out = append(out, item{ id: id, at: base.Add(time.Duration(i+1) * time.Second), automationID: "auto-1", aggregation: "Minutes", fixture: "minute.json", }) } return out } // Главный оракул задачи: тай-брейк «побеждает пришедшая» сделал содержимое // витрины функцией ПОРЯДКА свёртки, и витрина остаётся свёрткой журнала только // потому, что порядок свёртки приведён к порядку журнала. // // Проверяется на входе, воспроизводящем настоящую гонку приёма: строки учёта // вставляются в ОБРАТНОМ хронологии порядке — так выглядит конкурентный приём, // где поздняя доставка закоммитила строку первой. Воркер обязан всё равно // свернуть их в порядке `(received_at, id)`, и отпечаток обязан совпасть с // пересборкой. // // Проверка снабжена отрицательным контролем: те же тела, свёрнутые в обратном // порядке, обязаны дать ДРУГОЙ отпечаток. Без него тест зеленел бы и на // правиле, которое к порядку безразлично, — то есть не проверял бы ничего. func TestЖиваяСвёрткаРавнаПересборкеНаСтолкновениях(t *testing.T) { t.Parallel() ctx := context.Background() dir := t.TempDir() arch := openArchive(t, filepath.Join(dir, "raw")) src := openStore(t, filepath.Join(dir, "live.db")) sink := newLogSink() log := slog.New(slog.NewJSONHandler(sink, nil)) items := collidingJournal(t, 4) bodies := make([][]byte, len(items)) for i := range items { bodies[i] = bumped(t, items[i].fixture, i) } // Учёт заполняется от последней доставки к первой. for i := len(items) - 1; i >= 0; i-- { writeBody(t, arch, src, items[i], bodies[i]) } w := replay.NewWorker(src, fold.New(arch, src, 0, log), log) 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) } // Порядок журнала соблюдён, значит жаловаться не на что. if n := sink.count("delivery folded out of journal order"); n != 0 { t.Errorf("записей о свёртке вне порядка журнала %d, ожидалось 0:\n%s", n, sink.dump()) } rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db")) if got, want := rep.Fingerprint, fingerprint(t, src); got != want { t.Errorf("живая свёртка и пересборка разошлись:\n приём %s\n пересборка %s", want, got) } // Отрицательный контроль: тот же журнал, свёрнутый задом наперёд. rev := openArchive(t, filepath.Join(dir, "rev-raw")) revStore := openStore(t, filepath.Join(dir, "rev.db")) f := fold.New(rev, revStore, 0, slog.New(slog.DiscardHandler)) for i := len(items) - 1; i >= 0; i-- { writeBody(t, rev, revStore, items[i], bodies[i]) if _, err := f.Fold(ctx, items[i].id); err != nil { t.Fatalf("свёртка %s: %v", items[i].id, err) } } if fingerprint(t, revStore) == rep.Fingerprint { t.Error("обратный порядок свёртки дал тот же отпечаток: тай-брейк не зависит от порядка, то есть правило не сработало") } } // Остаточное окно не молчит: доставка, ставшая видимой после того, как её // преемница по журналу уже свёрнута, обязана оставить `WARN`. // // Закрыть окно здесь нечем — метка `received_at` фиксируется при выпуске ULID, // а строка учёта появляется только после записи тела. Наблюдение и есть всё, // что можно сделать без изменения приёма; лечится расхождение пересборкой. func TestСвёрткаВнеПорядкаЖурналаНеМолчит(t *testing.T) { t.Parallel() ctx := context.Background() dir := t.TempDir() arch := openArchive(t, filepath.Join(dir, "raw")) st := openStore(t, filepath.Join(dir, "live.db")) sink := newLogSink() log := slog.New(slog.NewJSONHandler(sink, nil)) items := collidingJournal(t, 2) // Поздняя доставка доехала и свернулась первой — её строка учёта успела // закоммититься, пока ранняя писала тело. writeBody(t, arch, st, items[1], bumped(t, items[1].fixture, 1)) f := fold.New(arch, st, 0, slog.New(slog.DiscardHandler)) if _, err := f.Fold(ctx, items[1].id); err != nil { t.Fatalf("свёртка поздней доставки: %v", err) } // Теперь появляется ранняя. writeBody(t, arch, st, items[0], bumped(t, items[0].fixture, 0)) w := replay.NewWorker(st, fold.New(arch, st, 0, log), log) if _, err := w.Pass(ctx); err != nil { t.Fatalf("проход воркера: %v", err) } if n := sink.count("delivery folded out of journal order"); n != 1 { t.Fatalf("записей о свёртке вне порядка журнала %d, ожидалась 1:\n%s", n, sink.dump()) } // Значений точек в записи быть не должно — данные о здоровье чувствительнее // токенов. Проверяется по разобранной записи, а не по сырому буферу: метка // времени содержит цифры и совпала бы с любым коротким числом. assertNoPointValues(t, sink, "delivery folded out of journal order") } // Доставка в статусе `failed` преемницей для этого наблюдения не считается: она // в витрину ничего не записала, перестановка относительно неё содержимого не // разводит, а исход этот штатный — невыводимый слой. Учитывай его предикат, // `WARN` стал бы шумом, на который перестают смотреть. func TestОтказавшаяПреемницаНеСчитаетсяРасхождением(t *testing.T) { t.Parallel() ctx := context.Background() dir := t.TempDir() arch := openArchive(t, filepath.Join(dir, "raw")) st := openStore(t, filepath.Join(dir, "live.db")) sink := newLogSink() log := slog.New(slog.NewJSONHandler(sink, nil)) items := collidingJournal(t, 2) // Поздняя доставка приезжает без плотных метрик и без предшественника — // слой выводить не из чего, исход `failed`. late := items[1] late.aggregation = "Default" late.automationID = "auto-lonely" late.fixture = "sparse_sleep.json" writeBody(t, arch, st, late, fixture(t, late.fixture)) f := fold.New(arch, st, 0, slog.New(slog.DiscardHandler)) if _, err := f.Fold(ctx, late.id); err == nil { t.Fatal("поздняя доставка свернулась: слой не должен был вывестись") } if status, err := st.DeliveryStatus(ctx, late.id); err != nil || status != store.ParseFailed { t.Fatalf("статус поздней доставки %q (%v), ожидался %q", status, err, store.ParseFailed) } writeBody(t, arch, st, items[0], bumped(t, items[0].fixture, 0)) w := replay.NewWorker(st, fold.New(arch, st, 0, log), log) if _, err := w.Pass(ctx); err != nil { t.Fatalf("проход воркера: %v", err) } if n := sink.count("delivery folded out of journal order"); n != 0 { t.Errorf("записей о свёртке вне порядка журнала %d, ожидалось 0:\n%s", n, sink.dump()) } } // assertNoPointValues убеждается, что запись лога не несёт значений точек. // // Разбирает запись и выбрасывает служебные поля, а не ищет подстроку в сыром // буфере: метка времени содержит цифры, и проверка по буферу была бы флаки по // построению (docs/conventions/testing.md, запись 2026-08-02). func assertNoPointValues(t *testing.T, sink *logSink, msg string) { t.Helper() for _, line := range sink.linesOf(msg) { for _, k := range []string{"qty", "Avg", "Min", "Max", "value", "points"} { if _, ok := line[k]; ok { t.Errorf("запись %q несёт поле %q: значения точек в лог не попадают", msg, k) } } } }