reindex: пересборка витрины проигрыванием журнала

- `healthlog reindex` собирает витрину из журнала (тела архива + учёт
  доставок) в ОТДЕЛЬНЫЙ файл базы, строго по `(received_at, id)`; рабочую
  базу читает без наката миграций и не трогает вовсе. Подмену делает
  человек при остановленном сервисе: переименование поверх открытого
  дескриптора портит базу молча.
- Журналом считается архив, а не таблица доставок: тело без учётной записи
  заводится заново (метка из ULID, размер и хеш по распакованному телу),
  запись без тела переносится, но не сворачивается. Оракул сходимости
  встроен — два отпечатка и «объектов было/стало»; пустой журнал успехом не
  считается.
- Прогон живого архива переехал на новый пакет: второго проигрывателя
  журнала в проекте не осталось, а его утверждение о ключе сна перестало
  быть константой, протухающей с каждой доставкой.
This commit is contained in:
av
2026-08-02 09:07:46 +03:00
parent 84bcbbea5c
commit 5ae0c5ff81
36 changed files with 4452 additions and 258 deletions
+79
View File
@@ -10,11 +10,18 @@ import (
"compress/gzip"
"fmt"
"io"
"io/fs"
"os"
"path/filepath"
"sort"
"strings"
"time"
)
// bodyExt — расширение файла тела. Всё, что ему не соответствует, телом не
// является: остаток `*.json.gz.tmp` от прерванной записи, чужой файл, каталог.
const bodyExt = ".json.gz"
// Archive — каталог сырых тел, разложенных по дате приёма.
type Archive struct {
root string
@@ -28,6 +35,23 @@ func New(root string) (*Archive, error) {
return &Archive{root: root}, nil
}
// Existing открывает уже существующий архив и каталога не создаёт.
//
// Отличие от New не косметическое: приёму каталог создать надо, а пересборке —
// нельзя. Созданный на лету пустой каталог превращает запуск не из той
// директории в успешный прогон по пустому журналу, а пустая витрина совпадает
// по отпечатку с пустой витриной, то есть выглядит идеальной сходимостью.
func Existing(root string) (*Archive, error) {
fi, err := os.Stat(root)
if err != nil {
return nil, fmt.Errorf("open archive dir %q: %w", root, err)
}
if !fi.IsDir() {
return nil, fmt.Errorf("archive dir %q: не каталог", root)
}
return &Archive{root: root}, nil
}
// Root возвращает корневой каталог архива.
func (a *Archive) Root() string { return a.root }
@@ -60,6 +84,61 @@ func (a *Archive) Write(id string, at time.Time, body []byte) (string, error) {
return rel, nil
}
// Listing — что нашлось в архиве.
type Listing struct {
// Bodies — относительные пути тел, отсортированные лексикографически.
// Порядок здесь только для воспроизводимости перечисления: журнал
// упорядочивает не он, а время приёма.
Bodies []string
// Skipped — файлы, телом не являющиеся. Считаются, а не выбрасываются:
// молчаливый пропуск означал бы «тело есть, а в отчёте его нет».
Skipped []string
}
// List перечисляет тела архива.
//
// Ошибку чтения каталога отдаёт наружу, а не превращает в пустой список, и
// каталога не создаёт. Различие принципиально для пересборки: нечитаемый или
// отсутствующий каталог означает «неизвестно, есть ли тела», а не «тел нет», —
// а пустой журнал даёт пустую витрину, чей отпечаток совпадает с отпечатком
// любой другой пустой витрины, то есть выглядит идеальной сходимостью.
func (a *Archive) List() (Listing, error) {
var out Listing
err := filepath.WalkDir(a.root, func(path string, d fs.DirEntry, err error) error {
if err != nil {
// Отказ чтения каталога прекращает обход целиком: пропустить его
// значило бы молча потерять сутки журнала.
return fmt.Errorf("walk archive %q: %w", path, err)
}
if d.IsDir() {
return nil
}
rel, err := filepath.Rel(a.root, path)
if err != nil {
return fmt.Errorf("relative path %q: %w", path, err)
}
if !strings.HasSuffix(d.Name(), bodyExt) {
out.Skipped = append(out.Skipped, rel)
return nil
}
out.Bodies = append(out.Bodies, rel)
return nil
})
if err != nil {
return Listing{}, err //nolint:wrapcheck // ошибка уже обёрнута внутри обхода
}
sort.Strings(out.Bodies)
sort.Strings(out.Skipped)
return out, nil
}
// BodyID возвращает идентификатор доставки по относительному пути тела.
func BodyID(rel string) string {
return strings.TrimSuffix(filepath.Base(rel), bodyExt)
}
// Open открывает сохранённое тело для чтения (распакованным). Нужен для
// пересборки витрины из архива.
func (a *Archive) Open(rel string) (io.ReadCloser, error) {
+75
View File
@@ -4,6 +4,7 @@ import (
"io"
"os"
"path/filepath"
"reflect"
"strings"
"testing"
"time"
@@ -90,3 +91,77 @@ func newArchive(t *testing.T) *archive.Archive {
}
return a
}
// Перечисление тел — вход пересборки. Всё, что телом не является, обязано быть
// посчитано, а не выброшено молча: тело есть, а в отчёте его нет.
func TestListРазводитТелаИПрочиеФайлы(t *testing.T) {
t.Parallel()
root := t.TempDir()
a, err := archive.New(root)
if err != nil {
t.Fatalf("New: %v", err)
}
at := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
for _, id := range []string{"b", "a"} {
if _, err := a.Write(id, at, []byte(`{"data":{}}`)); err != nil {
t.Fatalf("Write: %v", err)
}
}
day := filepath.Join(root, "2026", "08", "01")
for _, name := range []string{"a.json.gz.tmp", "readme.txt"} {
if err := os.WriteFile(filepath.Join(day, name), []byte("x"), 0o600); err != nil {
t.Fatalf("подготовка файла: %v", err)
}
}
got, err := a.List()
if err != nil {
t.Fatalf("List: %v", err)
}
want := []string{
filepath.Join("2026", "08", "01", "a.json.gz"),
filepath.Join("2026", "08", "01", "b.json.gz"),
}
if !reflect.DeepEqual(got.Bodies, want) {
t.Errorf("тела %v, ожидались %v", got.Bodies, want)
}
if len(got.Skipped) != 2 {
t.Errorf("пропущено %v, ожидалось два файла", got.Skipped)
}
if id := archive.BodyID(got.Bodies[0]); id != "a" {
t.Errorf("BodyID = %q, ожидался %q", id, "a")
}
}
// Нечитаемый каталог означает «неизвестно, есть ли тела», а не «тел нет»:
// пустой журнал даёт пустую витрину, которая по отпечатку совпадает с любой
// другой пустой витриной и выглядит идеальной сходимостью.
func TestОтсутствующийКаталогНеЯвляетсяПустымАрхивом(t *testing.T) {
t.Parallel()
missing := filepath.Join(t.TempDir(), "нет-такого")
if _, err := archive.Existing(missing); err == nil {
t.Error("Existing создал или принял отсутствующий каталог")
}
if _, err := os.Stat(missing); err == nil {
t.Error("Existing создал каталог — пересборке этого делать нельзя")
}
// А приёму каталог создать надо: это и есть разница между конструкторами.
if _, err := archive.New(missing); err != nil {
t.Errorf("New: %v", err)
}
a, err := archive.Existing(missing)
if err != nil {
t.Fatalf("Existing после New: %v", err)
}
list, err := a.List()
if err != nil {
t.Fatalf("List: %v", err)
}
if len(list.Bodies) != 0 {
t.Errorf("в пустом архиве нашлись тела: %v", list.Bodies)
}
}
+11
View File
@@ -312,6 +312,17 @@ func (s *Service) fail(ctx context.Context, deliveryID string, cause error, unco
}
}
// ReadBody читает тело из архива с той же границей размера, что и свёртка.
//
// Экспортировано ради пересборки: ей нужно прочесть тело, у которого ещё нет
// учётной записи, чтобы посчитать размер и хеш. Своей копией чтения это делать
// нельзя — граница обязана быть общей, иначе тело, принятое приёмом со `200`,
// начнёт вечно отказывать на каждой пересборке, и договорённость «предел тот
// же» ничем не проверяется.
func (s *Service) ReadBody(rawPath string) ([]byte, error) {
return s.readBody(rawPath)
}
func (s *Service) readBody(rawPath string) ([]byte, error) {
r, err := s.arch.Open(rawPath)
if err != nil {
-207
View File
@@ -1,207 +0,0 @@
package fold_test
import (
"bytes"
"compress/gzip"
"context"
"flag"
"io"
"os"
"path/filepath"
"sort"
"strings"
"testing"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// archiveDir включает прогон сходимости на живом архиве.
//
// Флагом, а не переменной окружения и не путём по умолчанию: архив в
// репозиторий не попадает (данные о здоровье), прогон занимает минуту и не
// должен висеть на каждом `task gate`. Запускается командой
// `task verify:archive`.
var archiveDir = flag.String("healthlog.archive", "",
"каталог сырого архива для прогона сходимости (по умолчанию прогон пропускается)")
// Сходимость на живом архиве: тот же корпус, на котором выводились правила
// разбора, обязан пройти через код без потерь и без расхождений — и повторный
// прогон журнала обязан дать то же состояние.
func TestReplayЖивогоАрхива(t *testing.T) {
if *archiveDir == "" {
t.Skip("прогон живого архива выключен: задайте -healthlog.archive")
}
root := *archiveDir
bodies := collectBodies(t, root)
if len(bodies) == 0 {
t.Skipf("живого архива нет в %s — прогон пропущен", root)
}
f, arch, st := newFold(t)
ctx := context.Background()
partial := 0
sections := map[string]int{}
var folded, failed, incomparable int
for _, path := range bodies {
body, err := os.ReadFile(path)
if err != nil {
t.Fatalf("чтение %s: %v", path, err)
}
// Тела в архиве сжаты; распаковываем и кладём через тот же архив, чтобы
// путь чтения был ровно тот, каким пойдёт пересборка.
//
// Заголовки доставки в архиве не лежат — они были заголовками запроса.
// Поэтому автоматизация у всех одна: так проверяется в том числе
// наследование слоя по цепочке доставок.
id := strings.TrimSuffix(filepath.Base(path), ".json.gz")
deliver(t, arch, st, id, "", "auto", gunzip(t, body))
res, err := f.Fold(ctx, id)
if err != nil {
failed++
continue
}
folded++
incomparable += res.Incomparable
if len(res.Uncovered) > 0 {
partial++
for _, s := range res.Uncovered {
sections[s]++
}
}
}
// Несравнимые наборы полей — посылка, на которой стоит отказ от объединения
// полей: их не было ни разу на всём корпусе. Число печатается, а не
// проверяется: появление такого набора — событие для разбора, а не отказ
// сходимости.
t.Logf("доставок %d: свёрнуто %d, не свёрнуто %d, несравнимых наборов %d",
len(bodies), folded, failed, incomparable)
t.Logf("частично разобрано %d, непокрытые секции: %v", partial, sections)
if folded == 0 {
t.Fatal("ни одна доставка не свернулась")
}
// Половина живого потока не несёт metrics вовсе (находка 50): такие
// доставки обязаны быть отличимы от разобранных целиком, иначе ретеншен
// срежет тела, которые для stateOfMind единственный источник.
if partial == 0 {
t.Error("ни одной частично разобранной доставки — перечисление непокрытых секций не работает")
}
// Повторный прогон того же журнала не меняет состояния: свёртка
// детерминирована, и пересборка даёт то же, что живой приём.
//
// Сравнивается ОТПЕЧАТОК содержимого, а не число объектов: на координате
// всегда лежит ровно одна точка, и правило разрешения столкновений выбирает,
// какая это будет точка, а не сколько их. Счёт объектов совпал бы и при
// заведомо сломанном правиле.
before, err := st.CountBuckets(ctx)
if err != nil {
t.Fatalf("счёт объектов: %v", err)
}
fingerprintBefore, err := st.Fingerprint(ctx)
if err != nil {
t.Fatalf("отпечаток: %v", err)
}
var refolded int
for _, path := range bodies {
if _, err := f.Fold(ctx, strings.TrimSuffix(filepath.Base(path), ".json.gz")); err != nil {
continue
}
refolded++
}
if refolded != folded {
t.Fatalf("повторно свёрнуто %d доставок из %d — проверка идемпотентности вхолостую",
refolded, folded)
}
after, err := st.CountBuckets(ctx)
if err != nil {
t.Fatalf("счёт объектов: %v", err)
}
if before != after {
t.Errorf("повторный прогон журнала изменил число объектов: %d → %d", before, after)
}
fingerprintAfter, err := st.Fingerprint(ctx)
if err != nil {
t.Fatalf("отпечаток: %v", err)
}
if fingerprintBefore != fingerprintAfter {
t.Errorf("повторный прогон журнала изменил содержимое объектов:\n %s\n %s",
fingerprintBefore, fingerprintAfter)
}
// Отпечаток печатается всегда: это единственный способ сравнить состояние с
// тем, что давала прежняя редакция правила слияния. Эталон в репозитории не
// живёт — он производен от архива, которого нет ни на одной другой машине.
// Значений точек отпечаток не раскрывает: содержимое входит в него хешем.
t.Logf("объектов %d, отпечаток содержимого %s", after, fingerprintAfter)
// Главное измеренное число: ключ по метке дал бы 170 координат сна, ключ по
// интервалу — 174 (docs/local-research.md, находка 47). Если координата
// когда-нибудь схлопнется обратно до метки, здесь станет 170.
if got := countPoints(t, st, "sleep_analysis"); got != 174 {
t.Errorf("координат sleep_analysis %d, измерено 174: ключ схлопнул записи", got)
}
}
// countPoints считает точки метрики во всех слоях. Каталог разрезов — отдельная
// задача, поэтому здесь перебор по известным слоям, а не запрос к нему.
func countPoints(t *testing.T, st *store.Store, metric string) int {
t.Helper()
ctx := context.Background()
total := 0
for _, layer := range []string{"sample", "raw", "minute", "hour", "day"} {
hours, err := st.BucketHours(ctx, metric, layer)
if err != nil {
t.Fatalf("часы объектов: %v", err)
}
for _, h := range hours {
b, err := st.Bucket(ctx, metric, layer, h)
if err != nil {
t.Fatalf("чтение объекта: %v", err)
}
total += len(b.Points)
}
}
return total
}
func collectBodies(t *testing.T, root string) []string {
t.Helper()
var out []string
err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
if err != nil {
return nil //nolint:nilerr // архива может не быть — это не отказ теста
}
if !info.IsDir() && filepath.Ext(path) == ".gz" {
out = append(out, path)
}
return nil
})
if err != nil {
return nil
}
sort.Strings(out)
return out
}
func gunzip(t *testing.T, body []byte) []byte {
t.Helper()
gz, err := gzip.NewReader(bytes.NewReader(body))
if err != nil {
t.Fatalf("распаковка: %v", err)
}
defer func() { _ = gz.Close() }()
out, err := io.ReadAll(gz)
if err != nil {
t.Fatalf("чтение: %v", err)
}
return out
}
+20
View File
@@ -9,6 +9,7 @@ import (
"errors"
"fmt"
"strings"
"time"
"github.com/oklog/ulid/v2"
)
@@ -22,6 +23,25 @@ func NewID() string {
return strings.ToLower(ulid.Make().String())
}
// TimeOf возвращает время создания идентификатора: UTC, секундная точность —
// ровно та форма, в которой время хранится (см. store.Now).
//
// Нужен пересборке витрины: у тела, лежащего в архиве без учётной записи,
// другого источника метки приёма нет. Дата каталога архива не годится — она
// задаёт сутки, а порядок проигрывания нужен внутри суток.
//
// `.UTC()` здесь обязателен и не для красоты: ulid.Time собирает время через
// time.Unix, то есть в локальной зоне машины, а Truncate работает с абсолютной
// длительностью и зону не нормализует. Без приведения метка подобранного тела
// сравнивалась бы с меткой из БД по-разному на разных машинах.
func TimeOf(s string) (time.Time, error) {
id, err := ulid.ParseStrict(strings.ToUpper(strings.TrimSpace(s)))
if err != nil {
return time.Time{}, fmt.Errorf("%w: %q", ErrInvalid, s)
}
return ulid.Time(id.Time()).UTC().Truncate(time.Second), nil
}
// Parse валидирует внешний идентификатор и приводит его к каноническому виду.
// Вызывается на входных границах (HTTP, CLI) до запроса к БД: сравнение строк
// в SQLite побайтовое, а base32 ULID при декодировании нечувствителен к
+57
View File
@@ -4,6 +4,7 @@ import (
"errors"
"strings"
"testing"
"time"
"git.vakhrushev.me/av/healthlog/internal/ident"
)
@@ -56,3 +57,59 @@ func TestParseRejectsGarbage(t *testing.T) {
}
}
}
// Время из ULID нужно телу, лежащему в архиве без учётной записи: другого
// источника метки приёма у него нет.
func TestTimeOfВUTCИСекундах(t *testing.T) {
t.Parallel()
id := ident.NewID()
at, err := ident.TimeOf(id)
if err != nil {
t.Fatalf("TimeOf(%q): %v", id, err)
}
// UTC, а не локальная зона: ulid.Time собирает время через time.Unix, то
// есть в зоне машины, и метка подобранного тела сравнивалась бы с меткой из
// БД по-разному на разных машинах.
if at.Location() != time.UTC {
t.Errorf("зона %v, ожидался UTC", at.Location())
}
// Секундная точность — та же, что у store.Now: метка уезжает в колонку
// фиксированной ширины, где лексикографический порядок равен хронологии.
if at.Nanosecond() != 0 {
t.Errorf("метка %v несёт доли секунды", at)
}
if d := time.Since(at); d < 0 || d > time.Minute {
t.Errorf("метка %v далека от настоящего времени (%v)", at, d)
}
}
// Монотонность ULID — то, на чём стоит порядок проигрывания журнала для
// подобранных тел.
func TestTimeOfСохраняетПорядок(t *testing.T) {
t.Parallel()
first, err := ident.TimeOf(ident.NewID())
if err != nil {
t.Fatalf("TimeOf: %v", err)
}
time.Sleep(1100 * time.Millisecond)
second, err := ident.TimeOf(ident.NewID())
if err != nil {
t.Fatalf("TimeOf: %v", err)
}
if !second.After(first) {
t.Errorf("порядок меток не сохранился: %v, затем %v", first, second)
}
}
func TestTimeOfОтвергаетНеИдентификатор(t *testing.T) {
t.Parallel()
for _, s := range []string{"", "не-ulid", "01kyzbb5zy4cbkbc0agb6ad07", "readme.txt"} {
if _, err := ident.TimeOf(s); !errors.Is(err, ident.ErrInvalid) {
t.Errorf("TimeOf(%q) дал %v, ожидался ErrInvalid", s, err)
}
}
}
+8 -1
View File
@@ -90,7 +90,14 @@ func (s *Service) Accept(ctx context.Context, body []byte, meta Meta) (Result, e
Bytes: int64(len(body)),
SHA256: hex.EncodeToString(sum[:]),
}
receivedAt := store.Now()
// Метка приёма выводится ИЗ идентификатора, а не берётся вторым обращением
// к часам. Источник обязан быть один: у тела, лежащего в архиве без учётной
// записи, метку восстанавливают из ULID, и два разных источника разошлись бы
// на границе секунды — а от порядка журнала зависит наследование слоя.
receivedAt, err := ident.TimeOf(res.DeliveryID)
if err != nil {
receivedAt = store.Now()
}
rawPath, err := s.arch.Write(res.DeliveryID, receivedAt, body)
if err != nil {
+23 -5
View File
@@ -2,6 +2,7 @@
package logging
import (
"io"
"log/slog"
"os"
"strings"
@@ -10,21 +11,38 @@ import (
// New возвращает slog-логгер с указанным уровнем и форматом ("json"|"text").
func New(level, format string) *slog.Logger {
return NewTo(os.Stdout, level, format)
}
// NewErr — тот же логгер, что New, но в stderr.
//
// Нужен командам CLI: stdout у них занят отчётом человеку, и лог в том же
// потоке сделал бы отчёт неразбираемым — а именно его человек перенаправляет в
// файл и читает глазами.
func NewErr(level, format string) *slog.Logger {
return NewTo(os.Stderr, level, format)
}
// NewTo — общая форма: приёмник вывода параметром, как у log.New и
// slog.NewTextHandler. New и NewErr — тонкие обёртки над ней; отдельные имена
// существуют ради читаемости места вызова, а не ради разного поведения.
func NewTo(w io.Writer, level, format string) *slog.Logger {
opts := &slog.HandlerOptions{Level: parseLevel(level), ReplaceAttr: utcTime}
var handler slog.Handler
if strings.EqualFold(format, "text") {
handler = slog.NewTextHandler(os.Stdout, opts)
handler = slog.NewTextHandler(w, opts)
} else {
handler = slog.NewJSONHandler(os.Stdout, opts)
handler = slog.NewJSONHandler(w, opts)
}
return slog.New(handler)
}
// NewStderr — JSON-логгер в stderr (UTC) для фатальных ошибок старта, когда
// основной логгер ещё не собран (конфиг не прочитан).
// NewStderr — JSON-логгер в stderr для фатальных ошибок старта, когда основной
// логгер ещё не собран (конфиг не прочитан). Уровень и формат брать неоткуда,
// поэтому они фиксированы — этим он и отличается от NewErr.
func NewStderr() *slog.Logger {
return slog.New(slog.NewJSONHandler(os.Stderr, &slog.HandlerOptions{ReplaceAttr: utcTime}))
return NewTo(os.Stderr, "info", "json")
}
// utcTime приводит метку времени записи к UTC: однозначный порядок событий и
+65
View File
@@ -0,0 +1,65 @@
package logging
import (
"bytes"
"encoding/json"
"log/slog"
"strings"
"testing"
)
// Уровень — это адресат, а не громкость: DEBUG предназначен разработчику при
// отладке и в штатной работе наружу не выходит.
func TestУровеньОтсекаетНижние(t *testing.T) {
t.Parallel()
var buf bytes.Buffer
log := NewTo(&buf, "warn", "json")
log.Debug("отладка")
log.Info("событие")
log.Warn("может стать проблемой")
out := buf.String()
if strings.Contains(out, "отладка") || strings.Contains(out, "событие") {
t.Errorf("уровень warn пропустил записи ниже себя: %s", out)
}
if !strings.Contains(out, "может стать проблемой") {
t.Errorf("запись уровня warn потерялась: %s", out)
}
}
// Время в логах — UTC: иначе порядок событий между машинами не сравнить.
func TestВремяВUTC(t *testing.T) {
t.Parallel()
var buf bytes.Buffer
NewTo(&buf, "info", "json").Info("событие")
var rec map[string]any
if err := json.Unmarshal(buf.Bytes(), &rec); err != nil {
t.Fatalf("запись не JSON: %v (%s)", err, buf.String())
}
ts, _ := rec[slog.TimeKey].(string)
if !strings.HasSuffix(ts, "Z") {
t.Errorf("метка времени %q не в UTC", ts)
}
}
// Команды CLI пишут отчёт в stdout, поэтому их лог обязан идти в stderr —
// иначе отчёт не разобрать ни глазами, ни перенаправлением.
func TestNewErrПишетВStderrАNewВStdout(t *testing.T) {
t.Parallel()
// Прямой проверки потока здесь нет намеренно: подмена os.Stdout была бы
// мутацией глобала. Проверяем то, что от этих конструкторов зависит на
// самом деле, — что они собирают разные назначения и оба живые.
if New("info", "json") == nil || NewErr("info", "text") == nil {
t.Fatal("конструктор вернул nil")
}
var buf bytes.Buffer
NewTo(&buf, "info", "text").Info("событие", "ключ", "значение")
if !strings.Contains(buf.String(), "ключ=значение") {
t.Errorf("текстовый формат не применён: %s", buf.String())
}
}
+201
View File
@@ -0,0 +1,201 @@
package replay_test
import (
"bytes"
"compress/gzip"
"context"
"flag"
"io"
"os"
"path/filepath"
"sort"
"strings"
"testing"
"git.vakhrushev.me/av/healthlog/internal/ident"
"git.vakhrushev.me/av/healthlog/internal/store"
)
// archiveDir включает прогон сходимости на живом архиве.
//
// Флагом, а не переменной окружения и не путём по умолчанию: архив в
// репозиторий не попадает (данные о здоровье), прогон занимает минуту и не
// должен висеть на каждом `task gate`. Запускается командой
// `task verify:archive`.
var archiveDir = flag.String("healthlog.archive", "",
"каталог сырого архива для прогона сходимости (по умолчанию прогон пропускается)")
// Сходимость на живом архиве: тот же корпус, на котором выводились правила
// разбора, обязан пройти через код без потерь и без расхождений — и повторное
// проигрывание журнала обязано дать то же состояние.
//
// Прогон идёт через ту же операцию, которой пересобирает витрину
// `healthlog reindex`. Собственный обход и собственный порядок здесь были
// раньше и были ошибкой: они образовывали второй проигрыватель журнала, чьи
// правила разошлись с настоящим, — а зеленел бы при этом он.
func TestReplayЖивогоАрхива(t *testing.T) {
if *archiveDir == "" {
t.Skip("прогон живого архива выключен: задайте -healthlog.archive")
}
bodies := collectBodies(t, *archiveDir)
if len(bodies) == 0 {
t.Skipf("живого архива нет в %s — прогон пропущен", *archiveDir)
}
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
// Тела в архиве сжаты; распаковываем и кладём через тот же архив, чтобы
// путь чтения был ровно тот, каким пойдёт пересборка.
//
// Заголовки доставки в архиве не лежат — они были заголовками запроса.
// Поэтому автоматизация у всех одна: так проверяется в том числе
// наследование слоя по цепочке доставок. Метка приёма берётся из ULID,
// то есть хронология журнала настоящая.
for _, path := range bodies {
body, err := os.ReadFile(path)
if err != nil {
t.Fatalf("чтение %s: %v", path, err)
}
id := strings.TrimSuffix(filepath.Base(path), ".json.gz")
at, err := ident.TimeOf(id)
if err != nil {
t.Fatalf("время из ULID %s: %v", id, err)
}
writeBody(t, arch, src, item{id: id, at: at, automationID: "auto"}, gunzip(t, body))
}
first, dst := run(t, ctx, arch, src, filepath.Join(dir, "first.db"))
// Несравнимые наборы полей — посылка, на которой стоит отказ от объединения
// полей: их не было ни разу на всём корпусе. Число печатается, а не
// проверяется: появление такого набора — событие для разбора, а не отказ
// сходимости.
t.Logf("тел %d: свёрнуто %d; отказов: слой %d, содержимое %d, прочее %d; несравнимых наборов %d",
first.Bodies, first.Folded, first.FailedLayer, first.FailedMalformed,
first.FailedOther, first.Incomparable)
t.Logf("частично разобрано %d, подобрано без учёта %d, записей без тела %d",
first.Partial, first.Adopted, first.Orphans)
if first.Folded == 0 {
t.Fatal("ни одна доставка не свернулась")
}
// Половина живого потока не несёт metrics вовсе (находка 50): такие
// доставки обязаны быть отличимы от разобранных целиком, иначе ретеншен
// срежет тела, которые для stateOfMind единственный источник.
if first.Partial == 0 {
t.Error("ни одной частично разобранной доставки — перечисление непокрытых секций не работает")
}
// Повторное проигрывание того же журнала даёт то же состояние: свёртка
// детерминирована, и пересборка даёт то же, что живой приём.
//
// Сравнивается ОТПЕЧАТОК содержимого, а не число объектов: на координате
// всегда лежит ровно одна точка, и правило разрешения столкновений выбирает,
// какая это будет точка, а не сколько их. Счёт объектов совпал бы и при
// заведомо сломанном правиле.
second, _ := run(t, ctx, arch, src, filepath.Join(dir, "second.db"))
if second.Buckets != first.Buckets {
t.Errorf("повторное проигрывание изменило число объектов: %d → %d", first.Buckets, second.Buckets)
}
if second.Fingerprint != first.Fingerprint {
t.Errorf("повторное проигрывание изменило содержимое объектов:\n %s\n %s",
first.Fingerprint, second.Fingerprint)
}
// Отпечаток печатается всегда: это единственный способ сравнить состояние с
// тем, что давала прежняя редакция правила слияния. Эталон в репозитории не
// живёт — он производен от архива, которого нет ни на одной другой машине.
// Значений точек отпечаток не раскрывает: содержимое входит в него хешем.
t.Logf("объектов %d, отпечаток содержимого %s", first.Buckets, first.Fingerprint)
// Главное свойство ключа: у записей сна он ИНТЕРВАЛ, а не метка — под одним
// `date` лежит до трёх записей (docs/local-research.md, находка 47).
//
// Проверяется само свойство, а не измеренное когда-то число. Прежняя
// редакция сравнивала с константой 174, снятой на 94 доставках, и покраснела
// молча, когда архив дорос до 116: константа, производная от корпуса,
// протухает с каждой новой доставкой, а прогон живого архива в гейт не
// входит, так что краснота никому не видна.
coords, labels := countSleepKeys(t, dst)
if coords == 0 {
t.Fatal("записей сна в витрине нет — проверять нечего")
}
if coords <= labels {
t.Errorf("координат сна %d при %d различных метках: ключ схлопнул записи до метки",
coords, labels)
}
t.Logf("координат сна %d, различных меток %d", coords, labels)
}
// countSleepKeys возвращает число различных координат записей сна и число
// различных меток начала. Разница между ними и есть то, что теряет ключ по
// метке.
//
// Каталог разрезов — отдельная задача, поэтому здесь перебор по известным
// слоям, а не запрос к нему.
func countSleepKeys(t *testing.T, st *store.Store) (coords, labels int) {
t.Helper()
type key struct{ start, end int64 }
ctx := context.Background()
seenCoord := map[key]struct{}{}
seenLabel := map[int64]struct{}{}
for _, layer := range []string{"sample", "raw", "minute", "hour", "day"} {
hours, err := st.BucketHours(ctx, "sleep_analysis", layer)
if err != nil {
t.Fatalf("часы объектов: %v", err)
}
for _, h := range hours {
b, err := st.Bucket(ctx, "sleep_analysis", layer, h)
if err != nil {
t.Fatalf("чтение объекта: %v", err)
}
for _, p := range b.Points {
seenCoord[key{p.Start.UnixNano(), p.End.UnixNano()}] = struct{}{}
seenLabel[p.Start.UnixNano()] = struct{}{}
}
}
}
return len(seenCoord), len(seenLabel)
}
func collectBodies(t *testing.T, root string) []string {
t.Helper()
var out []string
err := filepath.Walk(root, func(path string, info os.FileInfo, err error) error {
if err != nil {
return nil //nolint:nilerr // архива может не быть — это не отказ теста
}
if !info.IsDir() && filepath.Ext(path) == ".gz" {
out = append(out, path)
}
return nil
})
if err != nil {
return nil
}
sort.Strings(out)
return out
}
func gunzip(t *testing.T, body []byte) []byte {
t.Helper()
gz, err := gzip.NewReader(bytes.NewReader(body))
if err != nil {
t.Fatalf("распаковка: %v", err)
}
defer func() { _ = gz.Close() }()
out, err := io.ReadAll(gz)
if err != nil {
t.Fatalf("чтение: %v", err)
}
return out
}
+412
View File
@@ -0,0 +1,412 @@
// Package replay — проигрывание журнала доставок в витрину.
//
// Состояние healthlog есть свёртка по журналу:
// `import(снапшот экспорта) + replay(доставки по received_at)`. Здесь живёт
// вторая половина формулы; пересборка из архива — её вырожденный случай с
// пустым снапшотом, а не отдельная операция. Когда появится импорт родного
// экспорта Apple, он добавит стадию снапшота ПЕРЕД проигрыванием и переиспользует
// эту же операцию.
//
// Собственного разбора и собственного слияния пакет не имеет: он зовёт ту же
// свёртку, что и приём, по идентификатору доставки. Второй путь разбора
// разошёлся бы с первым молча — и уже расходился: прежний прогон живого архива
// сортировал тела по путям, а не по времени приёма.
package replay
import (
"context"
"crypto/sha256"
"encoding/hex"
"errors"
"fmt"
"log/slog"
"sort"
"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"
)
// Options — что нужно проигрыванию.
type Options struct {
// Archive — журнал тел. Именно он источник состава: перечислять только
// строки учёта значило бы пересобирать витрину из витрины.
Archive *archive.Archive
// Source — рабочая база, откуда берётся учёт доставок. Открывается на
// чтение: заголовков доставки в архиве нет, а восстановить их неоткуда.
// Пустой источник допустим — тогда весь журнал состоит из подобранных тел.
Source *store.Store
// Target — база назначения, куда собирается витрина.
Target *store.Store
// Fold — свёртка над Target. Собирается вызывающим с тем же пределом
// размера тела, что и у приёма; тем же пределом пересборка читает тело при
// подборе — своей копии границы у неё нет намеренно.
Fold *fold.Service
// Progress зовётся по мере продвижения. Нужен потому, что прогон на полном
// архиве молчит минутами, и зависший неотличим от идущего.
Progress func(done, total int)
Log *slog.Logger
}
// Report — итог проигрывания. Значений точек не несёт: содержимое входит в
// отчёт только отпечатком.
type Report struct {
// Bodies — тел в архиве, признанных телами.
Bodies int
// SkippedFiles — файлов, телом не являющихся (чужое расширение, остаток
// прерванной записи, имя не разбирается как ULID).
SkippedFiles int
// Duplicates — тел, чей идентификатор уже встретился в другом каталоге.
// Проигрывается первое; второе считается, а не роняет прогон.
Duplicates int
// Adopted — тел, у которых учётной записи не было: приём успел записать
// тело и не успел строку.
Adopted int
// AdoptFailed — тел без учёта, которые не удалось прочитать, чтобы завести
// запись. Тело остаётся в архиве.
AdoptFailed int
// Orphans — строк учёта, у которых тела в архиве нет. Станет штатным, когда
// появится ретеншен архива.
Orphans int
// Folded — сколько доставок свернулось.
Folded int
// Отказы разведены по классам, потому что читаются они по-разному.
// FailedLayer — слой не выводится: штатный исход, таких доставок в журнале
// заведомо есть. FailedMalformed — содержимое не разбирается. FailedOther —
// всё прочее (тело не читается, отказ базы); только оно означает, что с
// пересборкой что-то не так. Один общий счётчик отправлял бы человека
// искать дефект там, где его нет.
FailedLayer int
FailedMalformed int
FailedOther int
// Partial — доставок, в теле которых остались непокрытые разбором секции.
// Не отклонение, а половина потока; названо потому, что именно эти тела
// ретеншену трогать нельзя.
Partial int
// Incomparable — столкновений с несравнимыми наборами полей. На живом потоке
// их не было ни разу, и на этом стоит отказ от объединения полей.
Incomparable int
Buckets int64
Fingerprint string
// Canceled — проигрывание прервано отменой, а не дошло до конца.
Canceled bool
}
// Run проигрывает журнал в базу назначения.
//
// Порядок строго `(received_at, id)`: слой доставки без плотных метрик
// наследуется от ПРЕДШЕСТВУЮЩЕЙ доставки той же автоматизации, то есть является
// функцией префикса журнала. Обход каталога совпадает с хронологией только по
// датам каталогов и внутри суток не упорядочивает ничего.
//
// Проигрывание последовательное. Распараллеливать его нельзя по той же причине:
// слой зависит от префикса, а запись часового объекта — это чтение, слияние и
// запись обратно.
func Run(ctx context.Context, o Options) (Report, error) {
var rep Report
log := o.Log
if log == nil {
log = slog.New(slog.DiscardHandler)
}
log = log.With("capability", "replay")
// База назначения обязана быть пустой: проигрывание поверх накопленного
// оставило бы в витрине результат прежнего разбора — точки из объекта не
// удаляются никогда, — то есть не сделало бы того, ради чего пересборка и
// существует.
n, err := o.Target.CountDeliveries(ctx)
if err != nil {
return stopOr(rep, err)
}
if n > 0 {
return rep, fmt.Errorf("база назначения не пуста: %d доставок", n)
}
journal, orphans, rep, err := o.collect(ctx, rep, log)
if err != nil {
return stopOr(rep, err)
}
// Строка старта: прогон идёт минутами, и убитый на середине не оставлял бы
// о себе в логе ни следа — только чужие с виду записи свёртки.
log.InfoContext(ctx, "journal replay started",
"archive_dir", o.Archive.Root(), "bodies", len(journal), "orphans", len(orphans))
// Учётные записи, тела которых в архиве нет, переносятся ПЕРВЫМИ и не
// сворачиваются. Не перенести их значило бы потерять факты журнала —
// заголовки, хеш, автоматизацию — при первой же подмене базы: тела уже нет,
// и восстановить их будет нечем. В наследовании слоя они не участвуют:
// выведенного слоя у них нет.
for _, d := range orphans {
if err := o.Target.CreateDelivery(ctx, d); err != nil {
return stopOr(rep, err)
}
}
for i, d := range journal {
if ctx.Err() != nil {
rep.Canceled = true
return rep, nil
}
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:
// Отмена, застигшая свёртку, — не отказ доставки: считать её отказом
// значило бы обвинить разбор в том, чего он не делал, и отправить
// человека искать дефект по логу.
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
}
if o.Progress != nil {
o.Progress(i+1, len(journal))
}
}
// Отмена могла прийти на последней доставке: без этой проверки запросы ниже
// вернули бы context.Canceled как обычную ошибку, и команда завершилась бы
// одной строкой «fatal», не напечатав частичного отчёта.
if ctx.Err() != nil {
rep.Canceled = true
return rep, nil
}
rep.Buckets, err = o.Target.CountBuckets(ctx)
if err != nil {
return stopOr(rep, err)
}
rep.Fingerprint, err = o.Target.Fingerprint(ctx)
if err != nil {
return stopOr(rep, err)
}
log.InfoContext(ctx, "journal replayed",
"bodies", rep.Bodies,
"skipped_files", rep.SkippedFiles,
"adopted", rep.Adopted,
"adopt_failed", rep.AdoptFailed,
"orphans", rep.Orphans,
"duplicates", rep.Duplicates,
"folded", rep.Folded,
"failed_layer", rep.FailedLayer,
"failed_malformed", rep.FailedMalformed,
"failed_other", rep.FailedOther,
"partial", rep.Partial,
"incomparable", rep.Incomparable,
"buckets", rep.Buckets)
return rep, nil
}
// collect собирает состав журнала и упорядочивает его.
//
// Второй возврат — учётные записи, тел которых в архиве нет. Они переносятся в
// базу назначения, но не сворачиваются.
func (o Options) collect(ctx context.Context, rep Report, log *slog.Logger) ([]store.Delivery, []store.Delivery, Report, error) {
listing, err := o.Archive.List()
if err != nil {
// Нечитаемый каталог означает «неизвестно, есть ли тела», а не «тел
// нет»: пустой журнал дал бы пустую витрину, чей отпечаток совпадает с
// отпечатком любой другой пустой витрины.
return nil, nil, rep, err
}
rep.SkippedFiles = len(listing.Skipped)
for _, rel := range listing.Skipped {
// Путь в архиве — дата и ULID, содержимого тела в нём нет. Без этой
// строки счётчик пропусков в отчёте не на что раскрыть: имя файла не
// узнать иначе как обходом архива руками.
log.DebugContext(ctx, "archive file is not a body", "file", rel)
}
var known []store.Delivery
if o.Source != nil {
known, err = o.Source.ListDeliveries(ctx)
if err != nil {
return nil, nil, rep, err
}
}
byID := make(map[string]store.Delivery, len(known))
for _, d := range known {
byID[d.ID] = d
}
journal := make([]store.Delivery, 0, len(listing.Bodies))
seen := make(map[string]struct{}, len(listing.Bodies))
for _, rel := range listing.Bodies {
// Имя обязано быть КАНОНИЧЕСКИМ идентификатором, а не приводиться к
// нему: ident.Parse — функция входной границы, она обрезает пробелы и
// поднимает регистр, и `\u00a0<ulid>.json.gz` дал бы ту же координату
// журнала, что настоящее тело. Подложенный файл сортируется раньше и
// вытеснил бы настоящее — а невидимые пробелы в этих данных уже
// встречались.
id, err := ident.Parse(archive.BodyID(rel))
if err != nil || archive.BodyID(rel) != id {
rep.SkippedFiles++
log.DebugContext(ctx, "archive file is not a body", "file", rel)
continue
}
if _, dup := seen[id]; dup {
// Одно имя в двух каталогах: копия, восстановленная руками, или
// тело, переложенное не туда. Вторая запись журнала с тем же
// идентификатором сорвала бы весь прогон отказом по первичному
// ключу — то есть один посторонний файл лишал бы пересборки всё
// остальное.
rep.Duplicates++
log.WarnContext(ctx, "archive body id repeats", "file", rel, "delivery_id", id)
continue
}
rep.Bodies++
seen[id] = struct{}{}
if d, ok := byID[id]; ok {
journal = append(journal, prepare(d))
continue
}
d, err := o.adopt(id, rel)
if err != nil {
// Тело есть, а завести по нему запись не вышло: без этой строки
// счётчик в отчёте не на что раскрыть, а отчёт при этом отсылает
// человека «разобраться по логу».
rep.AdoptFailed++
log.WarnContext(ctx, "archive body not adopted", "error", err, "file", rel, "delivery_id", id)
continue
}
rep.Adopted++
journal = append(journal, d)
}
// Учётная запись, тела которой в архиве нет. Переносится, но не
// сворачивается: не перенести значило бы стереть первой же подменой базы
// единственное свидетельство, что доставка была, — тела-то уже нет.
orphans := make([]store.Delivery, 0)
for _, d := range known {
if _, ok := seen[d.ID]; !ok {
rep.Orphans++
// Не `pending`: тела нет и не будет, а «этим разбором ещё не
// смотрели» обещало бы данные, которых не появится, и ретеншен,
// который pending не трогает никогда, берёг бы такие строки вечно.
o := prepare(d)
o.ParseStatus = store.ParseFailed
orphans = append(orphans, o)
log.DebugContext(ctx, "delivery body missing", "delivery_id", d.ID, "file", d.RawPath)
}
}
// Тотальный ключ: `received_at` хранится с секундной точностью, и доставки
// одной секунды без второго ключа шли бы в неопределённом порядке — два
// прогона одного журнала могли бы разойтись.
sort.Slice(journal, func(i, j int) bool {
a, b := journal[i], journal[j]
if !a.ReceivedAt.Equal(b.ReceivedAt) {
return a.ReceivedAt.Before(b.ReceivedAt)
}
return a.ID < b.ID
})
return journal, orphans, rep, nil
}
// stopOr отличает отмену от настоящего отказа: первая не ошибка операции, а
// требование прекратить работу, и застать она может любой шаг.
//
// Единая точка на весь пакет: разбросанные по шагам проверки означали бы, что
// каждый новый шаг обязан вспомнить правило руками, а забытый превращает
// Ctrl-C в «ошибку окружения» без частичного отчёта.
func stopOr(rep Report, err error) (Report, error) {
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
rep.Canceled = true
return rep, nil
}
return rep, err
}
// prepare оставляет от учётной записи ФАКТЫ ЖУРНАЛА и сбрасывает производные от
// разбора поля.
//
// `parse_status`, `points`, `derived_layer` и `uncovered_sections` — результат
// ПРЕДЫДУЩЕЙ свёртки, а не то, что приехало вместе с доставкой. Перенести их
// значило бы сделать пересобранную витрину функцией прошлого прогона: доставка,
// чей повторный разбор отказал (штатный исход, когда слой не выводится),
// сохранила бы слой прежнего разбора — свёртка не затирает его намеренно, — и
// следующая доставка той же автоматизации унаследовала бы его молча. Оба прогона
// при этом самосогласованы, поэтому проверка «повторная пересборка ничего не
// меняет» такого не ловит.
//
// Незаполненные здесь колонки получают значения по умолчанию схемы: пустой слой
// и пустой список непокрытых секций.
func prepare(d store.Delivery) store.Delivery {
return store.Delivery{
ID: d.ID,
ReceivedAt: d.ReceivedAt,
AutomationName: d.AutomationName,
AutomationID: d.AutomationID,
Aggregation: d.Aggregation,
Period: d.Period,
SessionID: d.SessionID,
Bytes: d.Bytes,
SHA256: d.SHA256,
RawPath: d.RawPath,
Headers: d.Headers,
ParseStatus: store.ParsePending,
}
}
// adopt заводит учётную запись для тела, лежащего в архиве без неё.
//
// Такое тело — не экзотика: приём кладёт тело на диск раньше строки в базе
// (обратный порядок дал бы учтённую доставку без данных), и отказ на вставке
// оставляет тело без учёта. Обещание подобрать его записано в пакете приёма.
//
// Восстанавливается ровно то, что выводится из самого тела и его имени.
// Заголовков в архиве нет вовсе, поэтому вывод слоя у такой доставки честно
// деградирует: наследовать не от чего и подтверждать нечем.
func (o Options) adopt(id, rel string) (store.Delivery, error) {
at, err := ident.TimeOf(id)
if err != nil {
return store.Delivery{}, err
}
body, err := o.Fold.ReadBody(rel)
if err != nil {
return store.Delivery{}, err
}
sum := sha256.Sum256(body)
return store.Delivery{
ID: id,
ReceivedAt: at,
// Размер и хеш — по РАСПАКОВАННОМУ телу, как их считает приём. Иначе в
// тех же колонках появились бы значения другой природы, и индекс по
// хешу начал бы врать на границе подобранных тел.
Bytes: int64(len(body)),
SHA256: hex.EncodeToString(sum[:]),
RawPath: rel,
ParseStatus: store.ParsePending,
}, nil
}
+627
View File
@@ -0,0 +1,627 @@
package replay_test
import (
"context"
"log/slog"
"os"
"path/filepath"
"sort"
"strings"
"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"
)
// item — одна доставка тестового журнала.
type item struct {
id string
at time.Time
automationID string
aggregation string
fixture string
}
func fixture(t *testing.T, name string) []byte {
t.Helper()
body, err := os.ReadFile(filepath.Join("..", "hae", "testdata", name))
if err != nil {
t.Fatalf("фикстура %s: %v", name, err)
}
return body
}
func openStore(t *testing.T, path string) *store.Store {
t.Helper()
st, err := store.Open(path)
if err != nil {
t.Fatalf("база %s: %v", path, err)
}
t.Cleanup(func() { _ = st.Close() })
return st
}
func openArchive(t *testing.T, root string) *archive.Archive {
t.Helper()
arch, err := archive.New(root)
if err != nil {
t.Fatalf("архив: %v", err)
}
return arch
}
// live воспроизводит живой приём: тело в архив, строка учёта, свёртка — в том
// порядке и тем кодом, каким это делает `internal/ingest`.
func live(t *testing.T, arch *archive.Archive, st *store.Store, items []item) {
t.Helper()
f := fold.New(arch, st, 0, slog.New(slog.DiscardHandler))
for _, it := range items {
writeBody(t, arch, st, it, fixture(t, it.fixture))
_, _ = f.Fold(context.Background(), it.id)
}
}
func writeBody(t *testing.T, arch *archive.Archive, st *store.Store, it item, body []byte) {
t.Helper()
rawPath, err := arch.Write(it.id, it.at, body)
if err != nil {
t.Fatalf("запись в архив: %v", err)
}
if st == nil {
return
}
err = st.CreateDelivery(context.Background(), store.Delivery{
ID: it.id,
ReceivedAt: it.at,
AutomationID: it.automationID,
Aggregation: it.aggregation,
Bytes: int64(len(body)),
SHA256: "-",
RawPath: rawPath,
Headers: `{"x-test":["1"]}`,
ParseStatus: store.ParsePending,
})
if err != nil {
t.Fatalf("запись доставки: %v", err)
}
}
// run проигрывает журнал в свежую базу и возвращает отчёт вместе с ней.
func run(t *testing.T, ctx context.Context, arch *archive.Archive, src *store.Store, out string) (replay.Report, *store.Store) {
t.Helper()
dst := openStore(t, out)
rep, err := replay.Run(ctx, replay.Options{
Archive: arch,
Source: src,
Target: dst,
Fold: fold.New(arch, dst, 0, slog.New(slog.DiscardHandler)),
})
if err != nil {
t.Fatalf("проигрывание: %v", err)
}
return rep, dst
}
func fingerprint(t *testing.T, st *store.Store) string {
t.Helper()
fp, err := st.Fingerprint(context.Background())
if err != nil {
t.Fatalf("отпечаток: %v", err)
}
return fp
}
// journal собирает журнал из фикстур с монотонными идентификаторами.
func journal(t *testing.T, fixtures ...string) []item {
t.Helper()
base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
out := make([]item, 0, len(fixtures))
for i, f := range fixtures {
out = append(out, item{
id: ident.NewID(),
at: base.Add(time.Duration(i+1) * time.Second),
automationID: "auto-1",
aggregation: "Default",
fixture: f,
})
}
return out
}
// Главная проверка задачи: пересборка с нуля даёт то же состояние, что
// накопленный приём, а повторный прогон ничего не меняет.
func TestПересборкаСовпадаетСПриёмом(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json", "hour.json", "raw.json", "mixed.json", "sparse_sleep.json")
live(t, arch, src, items)
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Bodies != len(items) {
t.Fatalf("тел в журнале %d, ожидалось %d", rep.Bodies, len(items))
}
if rep.Folded == 0 {
t.Fatal("ни одна доставка не свернулась")
}
if rep.Adopted != 0 || rep.Orphans != 0 || rep.SkippedFiles != 0 {
t.Errorf("журнал не должен был дать подобранных/сирот/пропусков: %+v", rep)
}
want := fingerprint(t, src)
if rep.Fingerprint != want {
t.Errorf("отпечаток пересобранной витрины не совпал с накопленной:\n приём %s\n пересборка %s",
want, rep.Fingerprint)
}
// Повторный прогон в ещё одну базу обязан дать то же самое: победитель
// координаты — функция множества кандидатов, а не порядка прихода.
again, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild2.db"))
if again.Fingerprint != rep.Fingerprint {
t.Errorf("повторная пересборка изменила состояние:\n %s\n %s", rep.Fingerprint, again.Fingerprint)
}
}
// Порядок проигрывания задаётся журналом, а не раскладкой файлов: доставка без
// плотных метрик наследует слой ПРЕДШЕСТВУЮЩЕЙ доставки той же автоматизации.
//
// Журнал устроен так, что порядок имён файлов ОБРАТЕН хронологии. Проигрывание
// по каталогу поставило бы доставку без плотных метрик первой — наследовать ей
// было бы не от чего, слой не вывелся бы, и точки не сохранились бы вовсе.
func TestПорядокЗадаётсяЖурналомАНеКаталогом(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
ids := []string{ident.NewID(), ident.NewID(), ident.NewID()}
sort.Strings(ids)
base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
// Хронология обратна лексикографике имён: самый ранний файл — самая поздняя
// доставка.
items := []item{
{id: ids[2], at: base.Add(1 * time.Second), automationID: "a", aggregation: "Default", fixture: "hour.json"},
{id: ids[1], at: base.Add(2 * time.Second), automationID: "a", aggregation: "Default", fixture: "minute.json"},
{id: ids[0], at: base.Add(3 * time.Second), automationID: "a", aggregation: "Default", fixture: "sparse_sleep.json"},
}
for _, it := range items {
writeBody(t, arch, src, it, fixture(t, it.fixture))
}
rep, dst := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.FailedLayer+rep.FailedMalformed+rep.FailedOther != 0 {
t.Fatalf("отказов %d: доставка без плотных метрик не нашла предшественника — порядок взят из каталога",
rep.FailedLayer+rep.FailedMalformed+rep.FailedOther)
}
// Предшественник — минутная доставка, значит эпизоды сна легли в minute.
hours, err := dst.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()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
orphanBody := fixture(t, "minute.json")
orphan := item{id: ident.NewID(), at: time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)}
// Строки учёта нет — только тело.
writeBody(t, arch, nil, orphan, orphanBody)
rep, dst := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Adopted != 1 {
t.Fatalf("подобрано %d тел, ожидалось 1: %+v", rep.Adopted, rep)
}
if rep.Folded != 1 {
t.Fatalf("свёрнуто %d, ожидалась 1", rep.Folded)
}
if rep.Buckets == 0 {
t.Error("точки подобранного тела не доехали до витрины")
}
// Размер и хеш считаются по РАСПАКОВАННОМУ телу — как их считает приём.
got, err := dst.DeliveryForParse(ctx, orphan.id)
if err != nil {
t.Fatalf("учёт подобранного тела: %v", err)
}
wantAt, err := ident.TimeOf(orphan.id)
if err != nil {
t.Fatalf("время из ULID: %v", err)
}
if !got.ReceivedAt.Equal(wantAt) {
t.Errorf("метка приёма %v, ожидалась из ULID %v", got.ReceivedAt, wantAt)
}
all, err := dst.ListDeliveries(ctx)
if err != nil {
t.Fatalf("учёт: %v", err)
}
if len(all) != 1 || all[0].Bytes != int64(len(orphanBody)) {
t.Errorf("размер подобранного тела %v, ожидался по распакованному %d", all, len(orphanBody))
}
}
// Учётная запись без тела станет штатной, когда появится ретеншен архива:
// тела срезаются до даты проверенного экспорта, а строки живут дольше.
func TestУчётБезТелаНеРоняетПрогон(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json")
live(t, arch, src, items)
// Строка есть, тело исчезло.
if err := os.Remove(filepath.Join(arch.Root(), "2026", "08", "01", items[0].id+".json.gz")); err != nil {
t.Fatalf("удаление тела: %v", err)
}
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Orphans != 1 {
t.Errorf("записей без тела %d, ожидалась 1: %+v", rep.Orphans, rep)
}
if rep.Bodies != 0 {
t.Errorf("тел %d, ожидался 0", rep.Bodies)
}
}
// Файл, телом не являющийся, считается отдельно: молчаливый пропуск означал бы
// «тело есть, а в отчёте его нет».
func TestФайлНеТелоСчитаетсяОтдельно(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json")
live(t, arch, src, items)
day := filepath.Join(arch.Root(), "2026", "08", "01")
// Остаток прерванной записи и файл с именем, которое не идентификатор.
for _, name := range []string{items[0].id + ".json.gz.tmp", "readme.txt", "не-ulid.json.gz"} {
if err := os.WriteFile(filepath.Join(day, name), []byte("x"), 0o600); err != nil {
t.Fatalf("подготовка файла: %v", err)
}
}
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.SkippedFiles != 3 {
t.Errorf("пропущено файлов %d, ожидалось 3: %+v", rep.SkippedFiles, rep)
}
if rep.Bodies != 1 || rep.Folded != 1 {
t.Errorf("тел %d, свёрнуто %d, ожидалось 1 и 1", rep.Bodies, rep.Folded)
}
}
// Слой прошлого разбора не должен доживать до наследования: иначе витрина
// оказывается функцией предыдущего прогона, а не журнала.
func TestСлойПрошлогоРазбораНеНаследуется(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json", "hour.json")
live(t, arch, src, items)
clean, _ := run(t, ctx, arch, src, filepath.Join(dir, "clean.db"))
// Портим производное поле: как если бы прежний разбор вывел другой слой.
for _, it := range items {
err := src.FinishParse(ctx, it.id, store.ParseOutcome{
Status: store.ParseDone,
Layer: "raw",
})
if err != nil {
t.Fatalf("порча derived_layer: %v", err)
}
}
dirty, _ := run(t, ctx, arch, src, filepath.Join(dir, "dirty.db"))
if dirty.Fingerprint != clean.Fingerprint {
t.Errorf("слой прошлого разбора повлиял на пересборку:\n чистая %s\n с порчей %s",
clean.Fingerprint, dirty.Fingerprint)
}
}
// Отмена прекращает проигрывание: это требование прекратить работу, а не
// свойство доставки.
func TestОтменаПрекращаетПроигрывание(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
items := journal(t, "minute.json", "hour.json", "raw.json")
live(t, arch, src, items)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
dst := openStore(t, filepath.Join(dir, "rebuild.db"))
rep, err := replay.Run(ctx, replay.Options{
Archive: arch,
Source: src,
Target: dst,
Fold: fold.New(arch, dst, 0, slog.New(slog.DiscardHandler)),
// Отмена приходит посреди журнала — так же, как её принесёт сигнал.
Progress: func(done, _ int) {
if done == 1 {
cancel()
}
},
})
if err != nil {
t.Fatalf("проигрывание: %v", err)
}
if !rep.Canceled {
t.Error("отмена не отмечена в отчёте")
}
if rep.Folded != 1 {
t.Errorf("свёрнуто %d доставок, ожидалась 1 до отмены", rep.Folded)
}
// Отмена до начала работы — тоже отмена, а не отказ.
stopped, stop := context.WithCancel(context.Background())
stop()
early, err := replay.Run(stopped, replay.Options{
Archive: arch,
Source: src,
Target: openStore(t, filepath.Join(dir, "rebuild2.db")),
Fold: fold.New(arch, dst, 0, slog.New(slog.DiscardHandler)),
})
if err != nil {
t.Fatalf("проигрывание при отменённом контексте: %v", err)
}
if !early.Canceled || early.Folded != 0 {
t.Errorf("отмена до старта дала %+v", early)
}
}
// Нечитаемый или отсутствующий каталог архива — отказ, а не пустой журнал:
// пустая витрина совпадает по отпечатку с пустой витриной и выглядит идеальной
// сходимостью.
func TestОтсутствующийАрхивЭтоОтказ(t *testing.T) {
t.Parallel()
dir := t.TempDir()
if _, err := archive.Existing(filepath.Join(dir, "нет-такого")); err == nil {
t.Fatal("отсутствующий каталог архива не дал отказа")
}
arch := openArchive(t, filepath.Join(dir, "raw"))
if err := os.RemoveAll(arch.Root()); err != nil {
t.Fatalf("удаление каталога: %v", err)
}
dst := openStore(t, filepath.Join(dir, "rebuild.db"))
_, err := replay.Run(context.Background(), replay.Options{
Archive: arch,
Target: dst,
Fold: fold.New(arch, dst, 0, slog.New(slog.DiscardHandler)),
})
if err == nil {
t.Error("исчезнувший каталог архива дал пустой журнал вместо отказа")
}
}
// Битое тело не срывает прогон: остальные доставки обязаны проиграться.
func TestБитоеТелоНеСрываетПрогон(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json", "hour.json")
live(t, arch, src, items)
broken := filepath.Join(arch.Root(), "2026", "08", "01", items[0].id+".json.gz")
if err := os.WriteFile(broken, []byte("не gzip"), 0o600); err != nil {
t.Fatalf("порча тела: %v", err)
}
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
// Битый gzip — «прочее», а не невыведенный слой: классы разведены именно
// затем, чтобы человек не искал дефект там, где его нет.
if rep.FailedOther != 1 {
t.Errorf("прочих отказов %d, ожидался 1: %+v", rep.FailedOther, rep)
}
if rep.Folded != 1 {
t.Errorf("свёрнуто %d, ожидалась 1 — прогон сорвался на битом теле", rep.Folded)
}
}
// Одно имя тела в двух каталогах — копия, восстановленная руками, или тело,
// переложенное не туда. Вторая запись журнала с тем же идентификатором сорвала
// бы весь прогон отказом по первичному ключу.
func TestПовторИдентификатораНеСрываетПрогон(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json")
live(t, arch, src, items)
// Тот же файл, но под другой датой.
body, err := os.ReadFile(filepath.Join(arch.Root(), "2026", "08", "01", items[0].id+".json.gz"))
if err != nil {
t.Fatalf("чтение тела: %v", err)
}
other := filepath.Join(arch.Root(), "2026", "07", "31")
if err := os.MkdirAll(other, 0o755); err != nil {
t.Fatalf("каталог: %v", err)
}
if err := os.WriteFile(filepath.Join(other, items[0].id+".json.gz"), body, 0o600); err != nil {
t.Fatalf("копия тела: %v", err)
}
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Duplicates != 1 {
t.Errorf("повторов %d, ожидался 1: %+v", rep.Duplicates, rep)
}
if rep.Folded != 1 {
t.Errorf("свёрнуто %d, ожидалась 1 — повтор сорвал прогон", rep.Folded)
}
}
// Доставки одной секунды упорядочиваются идентификатором: `received_at` хранится
// с секундной точностью, и без второго ключа два прогона одного журнала могли бы
// разойтись.
func TestДоставкиОднойСекундыУпорядоченыИдентификатором(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
ids := []string{ident.NewID(), ident.NewID()}
sort.Strings(ids)
at := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
// Обе доставки — одна секунда. Первая по идентификатору минутная, вторая без
// плотных метрик: если тай-брейк исчезнет, вторая может пойти первой и
// остаться без предшественника.
writeBody(t, arch, src, item{id: ids[0], at: at, automationID: "a"}, fixture(t, "minute.json"))
writeBody(t, arch, src, item{id: ids[1], at: at, automationID: "a"}, fixture(t, "sparse_sleep.json"))
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.FailedLayer != 0 {
t.Fatalf("слой не вывелся у %d доставок: порядок внутри секунды не задан", rep.FailedLayer)
}
// Повтор в другую базу обязан дать тот же отпечаток.
again, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild2.db"))
if again.Fingerprint != rep.Fingerprint {
t.Errorf("порядок внутри секунды не детерминирован:\n %s\n %s", rep.Fingerprint, again.Fingerprint)
}
}
// Учётная запись, тела которой нет, переносится в базу назначения: не перенести
// значило бы стереть первой же подменой единственное свидетельство, что
// доставка была, — тела уже нет, восстановить нечем.
func TestУчётБезТелаПереноситсяВБазуНазначения(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json", "hour.json")
live(t, arch, src, items)
if err := os.Remove(filepath.Join(arch.Root(), "2026", "08", "01", items[0].id+".json.gz")); err != nil {
t.Fatalf("удаление тела: %v", err)
}
rep, dst := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Orphans != 1 {
t.Fatalf("записей без тела %d, ожидалась 1", rep.Orphans)
}
got, err := dst.ListDeliveries(ctx)
if err != nil {
t.Fatalf("учёт: %v", err)
}
if len(got) != len(items) {
t.Fatalf("строк учёта %d, ожидалось %d — запись без тела потеряна", len(got), len(items))
}
// Заголовки переносятся дословно: в архиве их нет вовсе.
for _, d := range got {
if d.Headers != `{"x-test":["1"]}` {
t.Errorf("заголовки доставки %s не дошли дословно: %q", d.ID, d.Headers)
}
}
// Статус — `failed`, а не `pending`. Различие несущее: `pending` означает
// «этим разбором ещё не смотрели» и обещает данные, которых не появится —
// тела уже нет, — а ретеншен, который pending не трогает никогда, берёг бы
// такие строки вечно.
b, err := dst.DeliveryStatus(ctx, items[0].id)
if err != nil {
t.Fatalf("статус записи без тела: %v", err)
}
if b != store.ParseFailed {
t.Errorf("статус записи без тела %q, ожидался %q", b, store.ParseFailed)
}
}
// Имя тела обязано быть КАНОНИЧЕСКИМ идентификатором. Разбор с приведением
// (обрезка пробелов, регистр) дал бы одну координату журнала двум файлам, и
// подложенный вытеснил бы настоящий — невидимые пробелы в этих данных уже
// встречались.
func TestИмяТелаОбязаноБытьКаноническим(t *testing.T) {
t.Parallel()
dir := t.TempDir()
arch := openArchive(t, filepath.Join(dir, "raw"))
src := openStore(t, filepath.Join(dir, "live.db"))
ctx := context.Background()
items := journal(t, "minute.json")
live(t, arch, src, items)
day := filepath.Join(arch.Root(), "2026", "08", "01")
body, err := os.ReadFile(filepath.Join(day, items[0].id+".json.gz"))
if err != nil {
t.Fatalf("чтение тела: %v", err)
}
// Те же 26 знаков, но с невидимым префиксом и в верхнем регистре: обе формы
// ident.Parse приводит к тому же идентификатору.
for _, name := range []string{" " + items[0].id + ".json.gz", strings.ToUpper(items[0].id) + ".json.gz"} {
if err := os.WriteFile(filepath.Join(day, name), body, 0o600); err != nil {
t.Fatalf("подложенный файл: %v", err)
}
}
rep, _ := run(t, ctx, arch, src, filepath.Join(dir, "rebuild.db"))
if rep.Bodies != 1 {
t.Errorf("тел %d, ожидалось 1: подложенное имя принято за тело (%+v)", rep.Bodies, rep)
}
if rep.SkippedFiles != 2 {
t.Errorf("пропущено %d, ожидалось 2", rep.SkippedFiles)
}
if rep.Duplicates != 0 {
t.Errorf("повторов %d: настоящее тело вытеснено подложенным", rep.Duplicates)
}
}
+56
View File
@@ -105,6 +105,62 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
return d, nil
}
// ListDeliveries возвращает учёт доставок в порядке журнала — `(received_at,
// id)`, тем же, в котором их проигрывает пересборка.
//
// Отдаются только **факты журнала**: то, что пришло вместе с доставкой.
// Производные от разбора поля (`parse_status`, `points`, `derived_layer`,
// `uncovered_sections`) сюда не попадают намеренно — перенос их в пересобранную
// базу сделал бы витрину функцией предыдущего прогона. Особенно `derived_layer`:
// доставка, чей повторный разбор отказал, отдала бы в наследование слой
// прежнего разбора, и следующая доставка той же автоматизации унаследовала бы
// его молча.
func (s *Store) ListDeliveries(ctx context.Context) ([]Delivery, error) {
const q = `
SELECT id, received_at, automation_name, automation_id, aggregation,
period, session_id, bytes, sha256, raw_path, headers
FROM delivery ORDER BY received_at, id`
rows, err := s.db.QueryContext(ctx, q)
if err != nil {
return nil, fmt.Errorf("select deliveries: %w", err)
}
defer func() { _ = rows.Close() }()
var out []Delivery
for rows.Next() {
var d Delivery
var receivedAt string
if err := rows.Scan(&d.ID, &receivedAt, &d.AutomationName, &d.AutomationID,
&d.Aggregation, &d.Period, &d.SessionID, &d.Bytes, &d.SHA256,
&d.RawPath, &d.Headers); err != nil {
return nil, fmt.Errorf("scan delivery: %w", err)
}
d.ReceivedAt, err = ParseTime(receivedAt)
if err != nil {
return nil, err
}
out = append(out, d)
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("select deliveries: %w", err)
}
return out, nil
}
// DeliveryStatus возвращает статус разбора доставки.
func (s *Store) DeliveryStatus(ctx context.Context, id string) (string, error) {
var status string
err := s.db.GetContext(ctx, &status, `SELECT parse_status FROM delivery WHERE id = ?`, id)
if errors.Is(err, sql.ErrNoRows) {
return "", ErrNotFound
}
if err != nil {
return "", fmt.Errorf("select parse status: %w", err)
}
return status, nil
}
// FinishParse записывает исход разбора доставки.
//
// Слой сохраняется здесь же, потому что он нужен следующей доставке той же
+110
View File
@@ -0,0 +1,110 @@
package store_test
import (
"context"
"database/sql"
"errors"
"path/filepath"
"testing"
"time"
_ "modernc.org/sqlite" // чистый Go-драйвер SQLite, без cgo
"git.vakhrushev.me/av/healthlog/internal/store"
)
// Пересборка читает рабочую базу, пока в неё может писать сервис. Обычное
// открытие накатывает миграции безусловно, а миграции меняют и данные — та, что
// ввела частичный разбор, переписала parse_status у всех строк. Значит утилите
// нужен путь чтения, который базу не трогает.
func TestOpenForReadНеПишетВБазу(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "healthlog.db")
seed(t, path)
ro, err := store.OpenForRead(path)
if err != nil {
t.Fatalf("OpenForRead: %v", err)
}
defer func() { _ = ro.Close() }()
ctx := context.Background()
got, err := ro.ListDeliveries(ctx)
if err != nil {
t.Fatalf("ListDeliveries: %v", err)
}
if len(got) != 2 {
t.Fatalf("доставок %d, ожидалось 2", len(got))
}
// Порядок журнала — (received_at, id), тот же, в котором проигрывает
// пересборка.
if got[0].ID != "01hzzzzzzzzzzzzzzzzzzzzzz1" || got[1].ID != "01hzzzzzzzzzzzzzzzzzzzzzz0" {
t.Errorf("порядок %q, %q — не по времени приёма", got[0].ID, got[1].ID)
}
// Заголовки переносятся дословно: восстановить их неоткуда, в архиве их нет.
if got[0].Headers != `{"x-test":["1"]}` {
t.Errorf("заголовки %q не дошли дословно", got[0].Headers)
}
if err := ro.CreateDelivery(ctx, store.Delivery{
ID: "01hzzzzzzzzzzzzzzzzzzzzzz2", ReceivedAt: store.Now(),
Bytes: 1, SHA256: "-", RawPath: "x", ParseStatus: store.ParsePending,
}); err == nil {
t.Error("запись в базу, открытую на чтение, удалась")
}
}
// Расхождение версии схемы — отказ, а не повод мигрировать: иначе свежий бинарь
// молча меняет схему под работающим старым сервисом.
func TestOpenForReadОтвергаетЧужуюВерсиюСхемы(t *testing.T) {
t.Parallel()
path := filepath.Join(t.TempDir(), "healthlog.db")
seed(t, path)
// Откатываем учёт миграций мимо store: имитируем базу, к которой бинарь
// новее.
db, err := sql.Open("sqlite", "file:"+path)
if err != nil {
t.Fatalf("sql.Open: %v", err)
}
_, err = db.Exec(`DELETE FROM goose_db_version
WHERE version_id = (SELECT max(version_id) FROM goose_db_version)`)
if err != nil {
t.Fatalf("откат версии: %v", err)
}
_ = db.Close()
if _, err := store.OpenForRead(path); !errors.Is(err, store.ErrSchemaMismatch) {
t.Errorf("OpenForRead дал %v, ожидался ErrSchemaMismatch", err)
}
}
func seed(t *testing.T, path string) {
t.Helper()
st, err := store.Open(path)
if err != nil {
t.Fatalf("Open: %v", err)
}
defer func() { _ = st.Close() }()
base := time.Date(2026, 8, 1, 12, 0, 0, 0, time.UTC)
// Второй идентификатор меньше первого, а приехал он раньше: так проверяется,
// что порядок берётся из времени приёма, а не из имени.
rows := []store.Delivery{
{ID: "01hzzzzzzzzzzzzzzzzzzzzzz1", ReceivedAt: base},
{ID: "01hzzzzzzzzzzzzzzzzzzzzzz0", ReceivedAt: base.Add(time.Second)},
}
for _, d := range rows {
d.Bytes = 1
d.SHA256 = "-"
d.RawPath = d.ID + ".json.gz"
d.ParseStatus = store.ParsePending
d.Headers = `{"x-test":["1"]}`
if err := st.CreateDelivery(context.Background(), d); err != nil {
t.Fatalf("CreateDelivery: %v", err)
}
}
}
+79
View File
@@ -5,9 +5,12 @@ package store
import (
"context"
"embed"
"errors"
"fmt"
"io/fs"
"net/url"
"strconv"
"strings"
"time"
"github.com/jmoiron/sqlx"
@@ -38,6 +41,70 @@ func Open(dbPath string) (*Store, error) {
return &Store{db: db}, nil
}
// ErrSchemaMismatch — версия схемы базы не та, которую знает бинарь.
var ErrSchemaMismatch = errors.New("версия схемы базы не совпадает с версией бинаря")
// OpenForRead открывает базу только для чтения и **без наката миграций**.
//
// Обычный Open мигрирует безусловно, а миграции здесь меняют не только схему, но
// и данные: та, что ввела частичный разбор, переписала parse_status у всех строк.
// Значит утилита, которой достаточно прочитать учёт, обычным открытием нарушала
// бы обещание «рабочую базу не трогаем», — и хуже: свежий бинарь мигрировал бы
// схему под работающим старым сервисом, который держит запросы к прежней.
//
// Расхождение версий — отказ с указанием обеих, а не повод мигрировать.
func OpenForRead(dbPath string) (*Store, error) {
db, err := sqlx.Connect("sqlite", readOnlyDSN(dbPath))
if err != nil {
return nil, fmt.Errorf("open sqlite %q read-only: %w", dbPath, err)
}
want, err := latestMigration()
if err != nil {
_ = db.Close()
return nil, err
}
var got int64
if err := db.Get(&got, `SELECT max(version_id) FROM goose_db_version`); err != nil {
_ = db.Close()
return nil, fmt.Errorf("read schema version: %w", err)
}
if got != want {
_ = db.Close()
return nil, fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, got, want)
}
return &Store{db: db}, nil
}
// latestMigration — номер последней миграции, вшитой в бинарь.
func latestMigration() (int64, error) {
entries, err := fs.ReadDir(migrationsFS, "migrations")
if err != nil {
return 0, fmt.Errorf("read migrations dir: %w", err)
}
var top int64
for _, e := range entries {
name := e.Name()
// Неразобранное имя — отказ, а не пропуск: страж «версия схемы не та»,
// молча не заметивший миграцию, перестаёт страховать, не сказав об этом.
idx := strings.IndexByte(name, '_')
if idx <= 0 {
return 0, fmt.Errorf("имя миграции %q не вида NNNNN_*.sql", name)
}
v, err := strconv.ParseInt(name[:idx], 10, 64)
if err != nil {
return 0, fmt.Errorf("имя миграции %q не вида NNNNN_*.sql", name)
}
if v > top {
top = v
}
}
if top == 0 {
return 0, errors.New("миграций не найдено")
}
return top, nil
}
// Close закрывает соединение с БД.
func (s *Store) Close() error {
if err := s.db.Close(); err != nil {
@@ -63,6 +130,18 @@ func dsn(path string) string {
return "file:" + path + "?" + q.Encode()
}
// readOnlyDSN — подключение только для чтения.
//
// `journal_mode` здесь не задаётся: сменить его на read-only соединении нельзя,
// а читать базу в режиме WAL это не мешает. `_txlock=immediate` тоже не нужен —
// он лечит повышение блокировки с чтения на запись, которого здесь не бывает.
func readOnlyDSN(path string) string {
q := url.Values{}
q.Add("mode", "ro")
q.Add("_pragma", "busy_timeout(5000)")
return "file:" + path + "?" + q.Encode()
}
// migrate накатывает миграции.
//
// Через Provider, а не через пакетные функции: goose.SetBaseFS и