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 } // linesOf отдаёт записи с данным msg РАЗОБРАННЫМИ. // // Разобранными, а не строками: проверка «в логе нет значения точки» обязана // смотреть поля записи, а не сырой буфер — служебное поле `time` содержит доли // секунды, и подстрока вроде «5.1» находится в нём примерно в одном проценте // прогонов (docs/review.md, 2026-08-02). func (s *logSink) linesOf(msg string) []map[string]any { s.mu.Lock() defer s.mu.Unlock() var out []map[string]any for _, line := range s.lines { if msgOf(line) != msg { continue } var rec map[string]any if err := json.Unmarshal([]byte(line), &rec); err != nil { continue } delete(rec, "time") out = append(out, rec) } return out } 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 }