diff --git a/CLAUDE.md b/CLAUDE.md index eb1a034..de7178f 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -74,6 +74,9 @@ Module path — `git.vakhrushev.me/av/healthlog`. - `task run` — запуск из исходников, без контейнера - `task build` — статический бинарь linux/amd64 - `task test` / `task lint` — тесты и golangci-lint +- `task verify:archive` — сходимость на живом архиве: весь `./data/raw` через + разбор, повтор обязан дать то же состояние. В гейт не входит намеренно — + минута прогона и данные, которых нет ни на какой другой машине - `task tidy` — `go mod tidy` - `task setup` — установка golangci-lint diff --git a/Taskfile.yml b/Taskfile.yml index 0a018d8..94e784c 100644 --- a/Taskfile.yml +++ b/Taskfile.yml @@ -33,6 +33,14 @@ tasks: cmds: - go test ./... + verify:archive: + desc: 'Сходимость на живом архиве: весь ./data/raw через разбор, повторный прогон обязан дать то же состояние' + cmds: + # Не входит в `task test` и `task gate` намеренно: архив в репозиторий не + # попадает, прогон занимает минуту, и держать его на каждом гейте значит + # платить за проверку, которая возможна только на этой машине. + - go test ./internal/fold -run TestReplay -healthlog.archive={{.ARCHIVE | default (printf "%s/data/raw" .ROOT_DIR)}} -v -count=1 + lint: desc: Запуск golangci-lint cmds: diff --git a/cmd/healthlog/serve.go b/cmd/healthlog/serve.go index cb63e0e..c482949 100644 --- a/cmd/healthlog/serve.go +++ b/cmd/healthlog/serve.go @@ -12,6 +12,7 @@ import ( "git.vakhrushev.me/av/healthlog/internal/archive" "git.vakhrushev.me/av/healthlog/internal/config" + "git.vakhrushev.me/av/healthlog/internal/fold" "git.vakhrushev.me/av/healthlog/internal/httpapi" "git.vakhrushev.me/av/healthlog/internal/ingest" "git.vakhrushev.me/av/healthlog/internal/logging" @@ -51,7 +52,7 @@ func runServe(args []string) error { } handler := httpapi.New(httpapi.Options{ - Ingest: ingest.New(arch, st, log), + Ingest: ingest.New(arch, st, fold.New(arch, st, log), log), Log: log, WriteTokens: cfg.Auth.WriteTokens, MaxBodyMB: cfg.Ingest.MaxBodyMB, diff --git a/docs/architecture.md b/docs/architecture.md index 5888a98..6f62417 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -410,14 +410,24 @@ hour метки выровнены на час heart_rate 00:00:00 (находка 31). Правило: 1. **Плотная метрика** (не меньше десяти точек в доставке) классифицируется - **сама по себе** по выравниванию своих меток. У десяти несуммированных - точек шанс всем лечь на ровную минуту исчезающе мал. + **сама по себе** по выравниванию своих меток — по **самому мелкому** + встретившемуся, а не преобладающему: метка ровно на часе одновременно + является и минутной, и у плотных метрик они перемешаны (`active_energy` — + 1320 минутных и 21 часовая). У десяти несуммированных точек шанс всем лечь + на ровную минуту исчезающе мал. 2. **Редкая метрика** (меньше десяти точек) наследует **преобладающий слой доставки** — самый мелкий среди плотных. У неё выравнивание ничего не доказывает, а Apple многие редкие показатели пишет прямо на границе часа. -3. Плотных метрик в доставке нет вовсе — слой наследуется от предыдущей - доставки той же автоматизации; если её не было, берём заголовок - (`Minutes` → `minute`, `Hours` → `hour`, иначе `raw`). +3. Плотных метрик в доставке нет вовсе — слой наследуется от **предшествующей** + доставки той же автоматизации; если её не было, берём **надёжный** заголовок + (`Minutes` → `minute`, `Hours` → `hour`). Иначе точки не сохраняются вовсе: + молчаливый `raw` создал бы призрачный разрез, который поедет в каталог и в + правило Read API «самый мелкий слой, покрывающий диапазон». + +Слово «предшествующей» в третьем пункте несёт вес: слой обязан быть функцией от +**префикса журнала**. Наследование от последней доставки вообще делает свёртку +зависящей от истории, и пересборка даёт не то состояние, что живой приём — +поймано прогоном архива, 1737 объектов против 1742 (docs/review-journal.md). Классифицировать доставку целиком нельзя: при перенастройке автоматизации приезжают **смешанные доставки**, где часть метрик уже минутная, а часть ещё diff --git a/docs/database.md b/docs/database.md index 850c04d..547e575 100644 --- a/docs/database.md +++ b/docs/database.md @@ -26,7 +26,8 @@ SQLite (`modernc.org/sqlite`, чистый Go), миграции — goose, фа │ parse_status TEXT │ │ sealed INTEGER │ │ points INTEGER │ │ created_at TEXT │ │ headers TEXT │ │ updated_at TEXT │ -└────────────────────────────┘ └──────────────────────────────┘ +│ derived_layer TEXT │ └──────────────────────────────┘ +└────────────────────────────┘ ``` Связь `bucket.first_delivery_id → delivery.id` **внешним ключом не объявлена** @@ -50,9 +51,10 @@ SQLite (`modernc.org/sqlite`, чистый Go), миграции — goose, фа | `parse_status` | `pending` / `parsed` / `failed`. Код ответа приёма от него **не зависит**: сохранили — значит приняли | | `points` | сколько точек дал разбор | | `headers` | все заголовки запроса JSON-объектом, кроме несущих секреты | +| `derived_layer` | слой, выведенный для этой доставки. Нужен не отчётности, а самому выводу: доставка без плотных метрик наследует последний надёжно выведенный слой той же автоматизации, и без хранения этой памяти первая такая доставка после перезапуска осталась бы без слоя | Индексы: `delivery_received_at` (порядок журнала), `delivery_sha256` (учёт -повторов). +повторов), `delivery_automation_layer` (поиск последнего слоя автоматизации). ## `bucket` — часовой объект точек diff --git a/docs/review-journal.md b/docs/review-journal.md index f653554..7aa9b1e 100644 --- a/docs/review-journal.md +++ b/docs/review-journal.md @@ -24,4 +24,26 @@ --- -Пока пусто — конвейер заведён 2026-08-01, задач через него не проходило. +## 2026-08-01 — свёртка не воспроизводилась при пересборке журнала + +- **Где:** `internal/store/delivery.go`, `LastDerivedLayer` +- **Симптом:** прогон живого архива (99 доставок) вторым проходом дал 1742 + объекта вместо 1737, а координат сна 182 вместо 174. Нашёл тест сходимости + на шаге apply — не ревью. +- **Причина:** доставка без плотных метрик наследует слой автоматизации. + Запрос брал последний выведенный слой **вообще**, а не последний до этой + доставки, поэтому при пересборке доставка наследовала слой «из будущего». + Свёртка переставала быть функцией от префикса журнала. +- **Почему не поймали:** формулировка «наследует последний надёжно выведенный + слой той же автоматизации» звучит однозначно и в спеке, и в дизайне — + пропущенное слово «предшествующей» не выглядит пропуском. Проходы `specs` и + `architecture` сверяли код со спекой и понятиями, а инвариант + «`import + replay` даёт то же состояние» ни один из них не проверял на + конкретном правиле: он записан в архитектуре как свойство системы, а не как + критерий для каждого узла, читающего состояние. +- **Что меняем:** в рубрику `healthlog-review-rubric` и в проход `ops` — вопрос + «читает ли узел состояние, которое сам же меняет, и остаётся ли он функцией + от префикса журнала». Дешевле правила: любой запрос к `delivery` из свёртки + обязан иметь границу по `received_at` разбираемой доставки. Тест сходимости + на живом архиве (`internal/fold/replay_test.go`) остаётся постоянным — + именно он это поймал. diff --git a/internal/fold/fold.go b/internal/fold/fold.go new file mode 100644 index 0000000..407102d --- /dev/null +++ b/internal/fold/fold.go @@ -0,0 +1,192 @@ +// Package fold — свёртка доставки в часовые объекты. +// +// Хранилище — свёртка по журналу: состояние пересобирается как +// `import(экспорт) + replay(доставки по received_at)`. Поэтому свёртка +// принимает идентификатор доставки, а тело читает из архива — тем же кодом, +// каким его прочитает пересборка. Разбор «из памяти, раз тело всё равно в +// руках» дал бы второй путь, который разошёлся бы с первым молча. +package fold + +import ( + "context" + "errors" + "fmt" + "io" + "log/slog" + + "git.vakhrushev.me/av/healthlog/internal/archive" + "git.vakhrushev.me/av/healthlog/internal/hae" + "git.vakhrushev.me/av/healthlog/internal/store" +) + +// maxBodyBytes — граница размера тела при чтении из архива. +// +// Наблюдалось 42 МиБ; сотня даёт запас втрое и при этом не даёт битому или +// враждебному архивному файлу выесть память процесса. Граница явная, потому +// что молчаливое «сколько дадут» — это отказ, который проявится только на +// пике потока. +const maxBodyBytes = 100 << 20 + +// Service сворачивает доставки в часовые объекты. +type Service struct { + arch *archive.Archive + store *store.Store + log *slog.Logger +} + +// New собирает свёртку. +func New(arch *archive.Archive, st *store.Store, log *slog.Logger) *Service { + return &Service{arch: arch, store: st, log: log.With("capability", "fold")} +} + +// Stats — итог свёртки одной доставки. +type Stats struct { + Metrics int + Points int + Buckets int + Unchanged int + Overwrites int + SealedHits int + SkippedNoTime int + SkippedMalformed int + Layer string + LayerMismatch bool +} + +// Fold разбирает тело доставки и раскладывает точки по часовым объектам. +// +// Это единственный логирующий чекпоинт свёртки: транспорт и приём исход +// разбора не логируют. Значения точек и имена устройств в лог не попадают — +// данные о здоровье чувствительнее токенов. +func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { + var stats Stats + + d, err := s.store.DeliveryForParse(ctx, deliveryID) + if err != nil { + return stats, err + } + + body, err := s.readBody(d.RawPath) + if err != nil { + s.fail(ctx, deliveryID, err) + return stats, err + } + + fallback, err := s.store.LastDerivedLayer(ctx, d.AutomationID, d.ReceivedAt, d.ID) + if err != nil { + return stats, err + } + + parsed, err := hae.Parse(body, hae.Meta{ + Aggregation: d.Aggregation, + FallbackLayer: hae.Layer(fallback), + }) + if err != nil { + s.fail(ctx, deliveryID, err) + return stats, err + } + + stats.Metrics = parsed.Metrics + stats.SkippedNoTime = parsed.SkippedNoTime + stats.SkippedMalformed = parsed.SkippedMalformed + stats.Layer = string(parsed.Layer) + stats.LayerMismatch = parsed.LayerMismatch + + merge, err := s.store.MergePoints(ctx, toIncoming(parsed.Points), deliveryID) + if err != nil { + s.fail(ctx, deliveryID, err) + return stats, err + } + + stats.Points = merge.Points + stats.Buckets = merge.Buckets + stats.Unchanged = merge.Unchanged + stats.Overwrites = merge.Overwrites + stats.SealedHits = merge.SealedHits + + if err := s.store.FinishParse(ctx, deliveryID, store.ParseDone, int64(stats.Points), stats.Layer); err != nil { + return stats, err + } + + s.logResult(ctx, deliveryID, stats) + return stats, nil +} + +func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) { + attrs := []any{ + "delivery_id", deliveryID, + "metrics", st.Metrics, + "points", st.Points, + "buckets", st.Buckets, + "unchanged", st.Unchanged, + "overwrites", st.Overwrites, + "skipped_no_time", st.SkippedNoTime, + "skipped_malformed", st.SkippedMalformed, + "layer", st.Layer, + } + + switch { + case st.SealedHits > 0: + // Досчёт часа, в который его уже не ждали: единственное наблюдение, по + // которому вообще можно судить о глубине досчёта. + s.log.WarnContext(ctx, "delivery folded, sealed hour changed", + append(attrs, "sealed_hits", st.SealedHits)...) + case st.LayerMismatch: + // Расхождение сверяется только с надёжным заголовком: `Default` не + // означает режима, и сравнение с ним давало бы WARN на каждой доставке. + s.log.WarnContext(ctx, "delivery folded, layer differs from header", attrs...) + default: + s.log.InfoContext(ctx, "delivery folded", attrs...) + } +} + +// fail отмечает доставку неразобранной. Тело остаётся в архиве, и её подберёт +// пересборка — приём при этом не затрагивается: сохранили значит приняли. +func (s *Service) fail(ctx context.Context, deliveryID string, cause error) { + level := slog.LevelError + if errors.Is(cause, hae.ErrLayerUnknown) { + // Слой не определился — это не поломка, а ожидаемый исход для доставки + // без плотных метрик. Тело ждёт пересборки. + level = slog.LevelWarn + } + s.log.Log(ctx, level, "delivery fold failed", "error", cause, "delivery_id", deliveryID) + + if err := s.store.FinishParse(ctx, deliveryID, store.ParseFailed, 0, ""); err != nil { + s.log.ErrorContext(ctx, "delivery parse status not recorded", "error", err, "delivery_id", deliveryID) + } +} + +func (s *Service) readBody(rawPath string) ([]byte, error) { + r, err := s.arch.Open(rawPath) + if err != nil { + return nil, err + } + defer func() { _ = r.Close() }() + + body, err := io.ReadAll(io.LimitReader(r, maxBodyBytes+1)) + if err != nil { + return nil, fmt.Errorf("чтение тела из архива: %w", err) + } + if len(body) > maxBodyBytes { + return nil, fmt.Errorf("тело больше %d байт", maxBodyBytes) + } + return body, nil +} + +func toIncoming(points []hae.Point) []store.IncomingPoint { + out := make([]store.IncomingPoint, 0, len(points)) + for _, p := range points { + out = append(out, store.IncomingPoint{ + Metric: p.Metric, + Layer: string(p.Layer), + Units: p.Units, + Point: store.Point{ + Start: p.Start, + End: p.End, + OffsetSeconds: p.OffsetSeconds, + Raw: p.Raw, + }, + }) + } + return out +} diff --git a/internal/fold/fold_test.go b/internal/fold/fold_test.go new file mode 100644 index 0000000..d125adc --- /dev/null +++ b/internal/fold/fold_test.go @@ -0,0 +1,230 @@ +package fold_test + +import ( + "context" + "log/slog" + "os" + "path/filepath" + "testing" + "time" + + "git.vakhrushev.me/av/healthlog/internal/archive" + "git.vakhrushev.me/av/healthlog/internal/fold" + "git.vakhrushev.me/av/healthlog/internal/store" +) + +func newFold(t *testing.T) (*fold.Service, *archive.Archive, *store.Store) { + t.Helper() + + dir := t.TempDir() + arch, err := archive.New(filepath.Join(dir, "raw")) + if err != nil { + t.Fatalf("архив: %v", err) + } + st, err := store.Open(filepath.Join(dir, "healthlog.db")) + if err != nil { + t.Fatalf("база: %v", err) + } + t.Cleanup(func() { _ = st.Close() }) + + log := slog.New(slog.DiscardHandler) + return fold.New(arch, st, log), arch, st +} + +// deliver кладёт тело в архив и заводит доставку — ровно то, что делает приём. +func deliver(t *testing.T, arch *archive.Archive, st *store.Store, id, aggregation, automationID string, body []byte) { + t.Helper() + + at := store.Now() + rawPath, err := arch.Write(id, at, body) + if err != nil { + t.Fatalf("запись в архив: %v", err) + } + err = st.CreateDelivery(context.Background(), store.Delivery{ + ID: id, + ReceivedAt: at, + AutomationID: automationID, + Aggregation: aggregation, + Bytes: int64(len(body)), + SHA256: "-", + RawPath: rawPath, + ParseStatus: store.ParsePending, + }) + if err != nil { + t.Fatalf("запись доставки: %v", err) + } +} + +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 TestFoldРаскладываетТочкиПоОбъектам(t *testing.T) { + t.Parallel() + + f, arch, st := newFold(t) + ctx := context.Background() + + deliver(t, arch, st, "d1", "Minutes", "auto-1", fixture(t, "minute.json")) + + stats, err := f.Fold(ctx, "d1") + if err != nil { + t.Fatalf("свёртка: %v", err) + } + if stats.Points == 0 || stats.Buckets == 0 { + t.Fatalf("свёртка ничего не разложила: %+v", stats) + } + if stats.Layer != "minute" { + t.Errorf("слой доставки %q, ожидался minute", stats.Layer) + } + + n, err := st.CountBuckets(ctx) + if err != nil { + t.Fatalf("счёт объектов: %v", err) + } + if n != int64(stats.Buckets) { + t.Errorf("объектов в базе %d, свёртка насчитала %d", n, stats.Buckets) + } +} + +// Слой сохраняется, чтобы его унаследовала следующая доставка той же +// автоматизации: без плотных метрик выводить его не из чего, а после +// перезапуска сервиса память об этом иначе теряется. +func TestFoldСлойНаследуетсяВнутриАвтоматизации(t *testing.T) { + t.Parallel() + + f, arch, st := newFold(t) + ctx := context.Background() + + deliver(t, arch, st, "d1", "Minutes", "auto-1", fixture(t, "minute.json")) + if _, err := f.Fold(ctx, "d1"); err != nil { + t.Fatalf("первая свёртка: %v", err) + } + + // Доставка без плотных метрик и с бесполезным заголовком: сама по себе + // слой не выводится. + deliver(t, arch, st, "d2", "Default", "auto-1", fixture(t, "sparse_sleep.json")) + stats, err := f.Fold(ctx, "d2") + if err != nil { + t.Fatalf("вторая свёртка: %v", err) + } + if stats.Layer != "minute" { + t.Errorf("слой %q, ожидался унаследованный minute", stats.Layer) + } +} + +// Наследовать нечего, заголовок ненадёжен: точки не сохраняются, доставка +// помечается неразобранной и ждёт пересборки. Молчаливый выбор raw создал бы +// призрачный разрез. +func TestFoldНеопределимыйСлойНеПишетТочек(t *testing.T) { + t.Parallel() + + f, arch, st := newFold(t) + ctx := context.Background() + + deliver(t, arch, st, "d1", "Default", "auto-1", fixture(t, "sparse_sleep.json")) + + if _, err := f.Fold(ctx, "d1"); err == nil { + t.Fatal("свёртка с неопределимым слоем прошла успешно") + } + + n, err := st.CountBuckets(ctx) + if err != nil { + t.Fatalf("счёт объектов: %v", err) + } + if n != 0 { + t.Errorf("записано %d объектов при неизвестном слое", n) + } +} + +// Повторная свёртка той же доставки не меняет состояния: журнал можно +// проигрывать сколько угодно раз. +func TestFoldИдемпотентен(t *testing.T) { + t.Parallel() + + f, arch, st := newFold(t) + ctx := context.Background() + + deliver(t, arch, st, "d1", "Minutes", "auto-1", fixture(t, "minute.json")) + + first, err := f.Fold(ctx, "d1") + if err != nil { + t.Fatalf("первая свёртка: %v", err) + } + second, err := f.Fold(ctx, "d1") + if err != nil { + t.Fatalf("повторная свёртка: %v", err) + } + + if second.Unchanged != second.Buckets { + t.Errorf("изменилось %d объектов из %d при повторе того же тела", + second.Buckets-second.Unchanged, second.Buckets) + } + if first.Buckets != second.Buckets { + t.Errorf("объектов %d против %d", first.Buckets, second.Buckets) + } +} + +// Порядок воспроизведения доставок не влияет на итог: свёртка по журналу +// обязана давать то же состояние, что приём в реальном времени. +func TestFoldНеЗависитОтПорядкаДоставок(t *testing.T) { + t.Parallel() + + state := func(order []string) (int64, string) { + f, arch, st := newFold(t) + ctx := context.Background() + + bodies := map[string]string{"a": "minute.json", "b": "hour.json", "c": "raw.json"} + aggs := map[string]string{"a": "Minutes", "b": "Hours", "c": "Default"} + for _, id := range order { + deliver(t, arch, st, id, aggs[id], "auto-"+id, fixture(t, bodies[id])) + if _, err := f.Fold(ctx, id); err != nil { + t.Fatalf("свёртка %s: %v", id, err) + } + } + + n, err := st.CountBuckets(ctx) + if err != nil { + t.Fatalf("счёт объектов: %v", err) + } + b, err := st.Bucket(ctx, "active_energy", "minute", mustHour(t, st, "active_energy", "minute")) + if err != nil { + return n, "" + } + return n, b.Hash + } + + nForward, hForward := state([]string{"a", "b", "c"}) + nBackward, hBackward := state([]string{"c", "b", "a"}) + + if nForward != nBackward { + t.Errorf("объектов %d против %d при разном порядке доставок", nForward, nBackward) + } + if hForward != hBackward { + t.Errorf("содержимое разошлось: %s против %s", hForward, hBackward) + } +} + +// mustHour находит час, в котором лежит объект метрики: фикстуры сдвинуты по +// времени, и зашивать конкретную дату в тест значило бы ломать его при каждой +// пересборке набора. +func mustHour(t *testing.T, st *store.Store, metric, layer string) time.Time { + t.Helper() + + hours, err := st.BucketHours(context.Background(), metric, layer) + if err != nil { + t.Fatalf("часы объекта: %v", err) + } + if len(hours) == 0 { + t.Fatalf("объектов метрики %s слоя %s нет", metric, layer) + } + return hours[0] +} diff --git a/internal/fold/replay_test.go b/internal/fold/replay_test.go new file mode 100644 index 0000000..20efd63 --- /dev/null +++ b/internal/fold/replay_test.go @@ -0,0 +1,161 @@ +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() + + var folded, failed 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)) + + if _, err := f.Fold(ctx, id); err != nil { + failed++ + continue + } + folded++ + } + + t.Logf("доставок %d: свёрнуто %d, не свёрнуто %d", len(bodies), folded, failed) + + if folded == 0 { + t.Fatal("ни одна доставка не свернулась") + } + + // Повторный прогон того же журнала не меняет состояния: свёртка + // детерминирована, и пересборка даёт то же, что живой приём. + before, err := st.CountBuckets(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) + } + + // Главное измеренное число: ключ по метке дал бы 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 +} diff --git a/internal/httpapi/httpapi_test.go b/internal/httpapi/httpapi_test.go index 7856bef..8ede0f4 100644 --- a/internal/httpapi/httpapi_test.go +++ b/internal/httpapi/httpapi_test.go @@ -12,6 +12,7 @@ import ( "testing" "git.vakhrushev.me/av/healthlog/internal/archive" + "git.vakhrushev.me/av/healthlog/internal/fold" "git.vakhrushev.me/av/healthlog/internal/httpapi" "git.vakhrushev.me/av/healthlog/internal/ingest" "git.vakhrushev.me/av/healthlog/internal/store" @@ -235,7 +236,7 @@ func newAPI(t *testing.T, writeTokens []string) (http.Handler, *store.Store) { log := slog.New(slog.DiscardHandler) h := httpapi.New(httpapi.Options{ - Ingest: ingest.New(arch, st, log), + Ingest: ingest.New(arch, st, fold.New(arch, st, log), log), Log: log, WriteTokens: writeTokens, MaxBodyMB: 1, diff --git a/internal/ingest/ingest.go b/internal/ingest/ingest.go index b9730ae..f8e23e7 100644 --- a/internal/ingest/ingest.go +++ b/internal/ingest/ingest.go @@ -11,8 +11,10 @@ import ( "errors" "fmt" "log/slog" + "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/store" ) @@ -46,16 +48,25 @@ type Result struct { RawPath string } -// Service принимает пакеты: сохраняет тело в архив и учитывает доставку. +// foldTimeout — сколько отводится свёртке принятой доставки. +// +// Свёртка идёт на контексте, отвязанном от запроса, поэтому собственный +// дедлайн обязателен: без него зависшая запись держала бы горутину до конца +// жизни процесса. +const foldTimeout = 2 * time.Minute + +// Service принимает пакеты: сохраняет тело в архив, учитывает доставку и +// запускает её свёртку. type Service struct { arch *archive.Archive store *store.Store + fold *fold.Service log *slog.Logger } // New собирает use-case приёма. -func New(arch *archive.Archive, st *store.Store, log *slog.Logger) *Service { - return &Service{arch: arch, store: st, log: log.With("capability", "ingest")} +func New(arch *archive.Archive, st *store.Store, f *fold.Service, log *slog.Logger) *Service { + return &Service{arch: arch, store: st, fold: f, log: log.With("capability", "ingest")} } // Accept принимает тело пакета: проверяет форму, кладёт в сырой архив и @@ -117,6 +128,17 @@ func (s *Service) Accept(ctx context.Context, body []byte, meta Meta) (Result, e "aggregation", meta.Aggregation, "period", meta.Period) + // Свёртка идёт после того, как доставка учтена, и на контексте, ОТВЯЗАННОМ + // от запроса: обрыв соединения клиентом или прокси на середине оставил бы + // часть объектов записанной, а доставку — со статусом, по которому её + // никто не подберёт. Исход свёртки на код ответа не влияет — сохранили + // значит приняли. + foldCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), foldTimeout) + defer cancel() + // Ошибку не возвращаем: она уже записана в лог и в parse_status свёрткой, + // а доставка принята. + _, _ = s.fold.Fold(foldCtx, res.DeliveryID) + return res, nil } diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 2177190..be62fc9 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -7,10 +7,12 @@ import ( "errors" "io" "log/slog" + "os" "path/filepath" "testing" "git.vakhrushev.me/av/healthlog/internal/archive" + "git.vakhrushev.me/av/healthlog/internal/fold" "git.vakhrushev.me/av/healthlog/internal/ingest" "git.vakhrushev.me/av/healthlog/internal/store" ) @@ -121,6 +123,71 @@ func TestAcceptRejectsMalformed(t *testing.T) { } } +// Разбор не влияет на исход приёма: сохранили — значит приняли. Непонятое +// содержимое даёт принятую доставку с parse_status=failed, а не отказ. +func TestAcceptНепонятоеСодержимоеПринимается(t *testing.T) { + svc, _, st := newService(t) + ctx := context.Background() + + // Метрика есть, но слой определить нечем: плотных метрик нет, заголовок + // ничего не означает, наследовать не от чего. + body := []byte(`{"data":{"metrics":[{"name":"m","units":"u","data":[` + + `{"date":"2026-07-31 12:00:00 +0300","qty":1}]}]}}`) + + if _, err := svc.Accept(ctx, body, ingest.Meta{Aggregation: "Default"}); err != nil { + t.Fatalf("Accept отверг доставку из-за разбора: %v", err) + } + + d, err := st.LastDelivery(ctx) + if err != nil { + t.Fatalf("LastDelivery: %v", err) + } + if d.ParseStatus != store.ParseFailed { + t.Errorf("parse_status = %q, ожидался %q", d.ParseStatus, store.ParseFailed) + } +} + +// Разобранная доставка отмечается разобранной, и точки доезжают до объектов. +func TestAcceptРазобраннаяДоставкаОтмечена(t *testing.T) { + svc, _, st := newService(t) + ctx := context.Background() + + body, err := os.ReadFile(filepath.Join("..", "hae", "testdata", "minute.json")) + if err != nil { + t.Fatalf("фикстура: %v", err) + } + + if _, err := svc.Accept(ctx, body, ingest.Meta{Aggregation: "Minutes"}); err != nil { + t.Fatalf("Accept: %v", err) + } + + d, err := st.LastDelivery(ctx) + if err != nil { + t.Fatalf("LastDelivery: %v", err) + } + if d.ParseStatus != store.ParseDone { + t.Fatalf("parse_status = %q, ожидался %q", d.ParseStatus, store.ParseDone) + } + if d.Points == 0 { + t.Error("точек 0: разбор не дошёл до учёта") + } + + n, err := st.CountBuckets(ctx) + if err != nil { + t.Fatalf("CountBuckets: %v", err) + } + if n == 0 { + t.Error("объектов 0: точки не доехали до хранилища") + } +} + +// Свёртка идёт на контексте, отвязанном от запроса (context.WithoutCancel в +// Accept), чтобы обрыв соединения не оставил часть объектов записанной. +// Автотестом это не покрыто: отмену надо подать РОВНО между учётом доставки и +// свёрткой, а такого шва снаружи нет, и заводить его ради теста дороже, чем +// проверять глазами. Атомарность самой записи проверена в store +// (TestMergePointsОтменаНеОставляетПоловины). + func newService(t *testing.T) (*ingest.Service, *archive.Archive, *store.Store) { t.Helper() dir := t.TempDir() @@ -137,5 +204,5 @@ func newService(t *testing.T) (*ingest.Service, *archive.Archive, *store.Store) } log := slog.New(slog.DiscardHandler) - return ingest.New(arch, st, log), arch, st + return ingest.New(arch, st, fold.New(arch, st, log), log), arch, st } diff --git a/internal/store/bucket.go b/internal/store/bucket.go index ace9b04..4a56041 100644 --- a/internal/store/bucket.go +++ b/internal/store/bucket.go @@ -457,6 +457,30 @@ func (s *Store) MarkSealed(ctx context.Context, metric, layer string, hour time. return nil } +// BucketHours возвращает часы, за которые есть объекты метрики в слое, по +// возрастанию. Первое, что понадобится каталогу разрезов, и то, чем проверка +// сходимости находит объект, не зная дат в фикстуре. +func (s *Store) BucketHours(ctx context.Context, metric, layer string) ([]time.Time, error) { + const q = ` + SELECT hour_utc FROM bucket + WHERE metric = ? AND layer = ? ORDER BY hour_utc` + + var raw []string + if err := s.db.SelectContext(ctx, &raw, q, metric, layer); err != nil { + return nil, fmt.Errorf("select bucket hours: %w", err) + } + + out := make([]time.Time, 0, len(raw)) + for _, s := range raw { + t, err := ParseTime(s) + if err != nil { + return nil, err + } + out = append(out, t) + } + return out, nil +} + // CountBuckets возвращает число часовых объектов. Нужно проверке сходимости и // healthcheck. func (s *Store) CountBuckets(ctx context.Context) (int64, error) { diff --git a/internal/store/delivery.go b/internal/store/delivery.go index a23f17f..eb78d36 100644 --- a/internal/store/delivery.go +++ b/internal/store/delivery.go @@ -8,10 +8,16 @@ import ( "time" ) -// Статусы разбора доставки. +// Статусы разбора доставки. Код ответа приёма от них не зависит: сохранили — +// значит приняли. const ( // ParsePending — тело сохранено, разбора ещё не было. ParsePending = "pending" + // ParseDone — тело разобрано, точки разложены по объектам. + ParseDone = "parsed" + // ParseFailed — разобрать не удалось. Тело лежит в архиве, доставку + // подберёт пересборка. + ParseFailed = "failed" ) // Delivery — учётная запись одного принятого пакета. @@ -86,6 +92,92 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) { return d, nil } +// FinishParse записывает исход разбора доставки. +// +// Слой сохраняется здесь же, потому что он нужен следующей доставке той же +// автоматизации: без плотных метрик выводить его не из чего, и наследовать +// приходится от прошлого раза. +func (s *Store) FinishParse(ctx context.Context, id, status string, points int64, layer string) error { + const q = ` + UPDATE delivery SET parse_status = ?, points = ?, derived_layer = ? + WHERE id = ?` + + res, err := s.db.ExecContext(ctx, q, status, points, layer, id) + if err != nil { + return fmt.Errorf("update parse status: %w", err) + } + n, err := res.RowsAffected() + if err != nil { + return fmt.Errorf("update parse status: %w", err) + } + if n == 0 { + return ErrNotFound + } + return nil +} + +// LastDerivedLayer возвращает слой, выведенный для этой автоматизации ПЕРЕД +// указанной доставкой. Пустая строка означает, что наследовать нечего. +// +// «Перед» здесь существенно, а не для красоты: слой обязан быть функцией от +// префикса журнала, иначе повторный прогон даёт другое состояние, чем живой +// приём. Запрос без границы по времени брал бы последний слой вообще — и при +// пересборке доставка наследовала бы слой от будущего. Проверено прогоном +// архива: 1737 объектов превращались в 1742. +func (s *Store) LastDerivedLayer(ctx context.Context, automationID string, before time.Time, beforeID string) (string, error) { + if automationID == "" { + return "", nil + } + + const q = ` + SELECT derived_layer FROM delivery + WHERE automation_id = ? AND derived_layer != '' + AND (received_at, id) < (?, ?) + ORDER BY received_at DESC, id DESC LIMIT 1` + + var layer string + err := s.db.GetContext(ctx, &layer, q, automationID, FormatTime(before), beforeID) + if errors.Is(err, sql.ErrNoRows) { + return "", nil + } + if err != nil { + return "", fmt.Errorf("select derived layer: %w", err) + } + return layer, nil +} + +// DeliveryBody — что нужно знать о доставке, чтобы разобрать её тело. +type DeliveryBody struct { + ID string + ReceivedAt time.Time + RawPath string + AutomationID string + Aggregation string +} + +// DeliveryForParse возвращает сведения о доставке, нужные разбору. +func (s *Store) DeliveryForParse(ctx context.Context, id string) (DeliveryBody, error) { + const q = ` + SELECT id, received_at, raw_path, automation_id, aggregation + FROM delivery WHERE id = ?` + + var d DeliveryBody + var receivedAt string + err := s.db.QueryRowxContext(ctx, q, id). + Scan(&d.ID, &receivedAt, &d.RawPath, &d.AutomationID, &d.Aggregation) + if errors.Is(err, sql.ErrNoRows) { + return DeliveryBody{}, ErrNotFound + } + if err != nil { + return DeliveryBody{}, fmt.Errorf("select delivery for parse: %w", err) + } + d.ReceivedAt, err = ParseTime(receivedAt) + if err != nil { + return DeliveryBody{}, err + } + return d, nil +} + // CountDeliveries возвращает число принятых пакетов. Нужно для healthcheck и // быстрой проверки «данные вообще идут». func (s *Store) CountDeliveries(ctx context.Context) (int64, error) { diff --git a/internal/store/migrations/00004_delivery_layer.sql b/internal/store/migrations/00004_delivery_layer.sql new file mode 100644 index 0000000..6a34bed --- /dev/null +++ b/internal/store/migrations/00004_delivery_layer.sql @@ -0,0 +1,17 @@ +-- +goose Up +-- Слой, выведенный для доставки. Нужен не отчётности, а самому выводу слоя: +-- доставка без плотных метрик наследует последний надёжно выведенный слой той +-- же автоматизации, и без хранения этой памяти первая же такая доставка после +-- перезапуска сервиса осталась бы без слоя (измерено: 2 доставки из 89). +-- +-- Пустая строка означает «слой не выводился» — такие доставки в наследование +-- не участвуют. +ALTER TABLE delivery ADD COLUMN derived_layer TEXT NOT NULL DEFAULT ''; + +-- Поиск идёт «последний слой этой автоматизации»: сортировка по времени приёма +-- внутри автоматизации. +CREATE INDEX delivery_automation_layer ON delivery (automation_id, received_at DESC); + +-- +goose Down +DROP INDEX delivery_automation_layer; +ALTER TABLE delivery DROP COLUMN derived_layer; diff --git a/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md b/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md index ee208ce..242e534 100644 --- a/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md +++ b/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md @@ -83,10 +83,20 @@ - **WHEN** ни у одной метрики доставки нет десяти точек - **THEN** слой наследуется от последнего надёжно выведенного слоя той же - автоматизации (`automation-id`) + автоматизации (`automation-id`) среди доставок, **предшествующих** этой - **AND** если наследовать нечего, слой берётся из **надёжного** заголовка (`Minutes` → `minute`, `Hours` → `hour`) +Граница «предшествующих» обязательна: слой обязан быть функцией от префикса +журнала. Наследование от последней доставки вообще делает свёртку зависящей от +истории, и пересборка даёт не то состояние, что живой приём — измерено на +архиве, 1737 объектов против 1742. + +#### Scenario: Пересборка журнала даёт то же состояние + +- **WHEN** те же доставки сворачиваются повторно в том же порядке +- **THEN** число объектов и их содержимое не меняются + #### Scenario: Наследовать нечего и заголовок ненадёжен - **WHEN** плотных метрик нет, предыдущего слоя автоматизации нет, а заголовок diff --git a/openspec/changes/razbor-metrik-v-obekty/tasks.md b/openspec/changes/razbor-metrik-v-obekty/tasks.md index cb43810..e8a5d99 100644 --- a/openspec/changes/razbor-metrik-v-obekty/tasks.md +++ b/openspec/changes/razbor-metrik-v-obekty/tasks.md @@ -44,16 +44,16 @@ ## 5. Сшивка с приёмом -- [ ] 5.1 Свёртка вызывается по идентификатору доставки, тело читается из архива — один код с будущей пересборкой -- [ ] 5.2 Работа после записи в архив — на `context.WithoutCancel` с собственным дедлайном: обрыв соединения не рвёт запись объектов -- [ ] 5.3 Отказ разбора не меняет код ответа: `200`, `parse_status=failed`, запись `ERROR` без значений -- [ ] 5.4 Единственный логирующий чекпоинт на границе: счётчики метрик, точек, объектов, пропусков, перезаписей; идентификатор доставки; ни значений, ни имён устройств -- [ ] 5.5 Тест приёма: битый JSON — 400, непонятое содержимое — 200 с `parse_status=failed` +- [x] 5.1 Свёртка вызывается по идентификатору доставки, тело читается из архива — один код с будущей пересборкой +- [x] 5.2 Работа после записи в архив — на `context.WithoutCancel` с собственным дедлайном: обрыв соединения не рвёт запись объектов +- [x] 5.3 Отказ разбора не меняет код ответа: `200`, `parse_status=failed`, запись `ERROR` без значений +- [x] 5.4 Единственный логирующий чекпоинт на границе: счётчики метрик, точек, объектов, пропусков, перезаписей; идентификатор доставки; ни значений, ни имён устройств +- [x] 5.5 Тест приёма: битый JSON — 400, непонятое содержимое — 200 с `parse_status=failed` ## 6. Сходимость на реальных данных -- [ ] 6.1 Скрипт `tmp/research/verify_buckets.py`: прогоняет архив через разбор и сверяет суммы по часовому слою с проверкой из разведки -- [ ] 6.2 Прогон на всех накопленных доставках: расхождений по накопительным метрикам нет +- [x] 6.1 Скрипт `tmp/research/verify_buckets.py`: прогоняет архив через разбор и сверяет суммы по часовому слою с проверкой из разведки +- [x] 6.2 Прогон на всех накопленных доставках: расхождений по накопительным метрикам нет - [ ] 6.3 `task gate` зелёный; `task restart` поднимает сервис, новая доставка с телефона разбирается ## 7. Приёмочные критерии (рубрика ревью дизайна) @@ -61,16 +61,16 @@ Порождена проходом `healthlog-review-rubric` **до** чтения предложения. Каждый пункт проверяем: понятно, каким тестом его провалить. -- [ ] 7.1 Отсутствие паники на произвольном входе — усечённый JSON, `null` вместо объекта, массив вместо объекта, число вместо строки даты -- [ ] 7.2 Незнакомое поле точки переживает round-trip дословно -- [ ] 7.3 Идемпотентность повторной записи, включая переставленный порядок ключей и точек -- [ ] 7.4 Различие не теряется молча: перезапись и изменение запечатанного часа оставляют след -- [ ] 7.5 Конкурентное слияние того же часа не теряет точки — под `-race`, с проверкой суммы -- [ ] 7.6 Отмена посреди слияния не оставляет половинчатого состояния: объект либо прежний, либо полный -- [ ] 7.7 Граница размера входа явная; вход в сотни мегабайт не кладёт процесс по памяти -- [ ] 7.8 Ошибки различимы по типу, а не по тексту (`errors.Is`/`errors.As`) -- [ ] 7.9 Частично непонятный пакет имеет явную судьбу: что сохранено, что отброшено — видно в счётчиках -- [ ] 7.10 Время нормализовано без потери зоны -- [ ] 7.11 Слой выводится из данных, а не из заголовка; неопределимый слой имеет явную судьбу -- [ ] 7.12 Парсер детерминирован и чист: ни `time.Now`, ни генерации id; два вызова на одном входе равны -- [ ] 7.13 (добавлено после снятия блокера) Точка-интервал не схлопывается по метке: прогон архива даёт 174 координаты сна, а не 170 +- [x] 7.1 Отсутствие паники на произвольном входе — усечённый JSON, `null` вместо объекта, массив вместо объекта, число вместо строки даты +- [x] 7.2 Незнакомое поле точки переживает round-trip дословно +- [x] 7.3 Идемпотентность повторной записи, включая переставленный порядок ключей и точек +- [x] 7.4 Различие не теряется молча: перезапись и изменение запечатанного часа оставляют след +- [x] 7.5 Конкурентное слияние того же часа не теряет точки — под `-race`, с проверкой суммы +- [x] 7.6 Отмена посреди слияния не оставляет половинчатого состояния: объект либо прежний, либо полный +- [x] 7.7 Граница размера входа явная; вход в сотни мегабайт не кладёт процесс по памяти +- [x] 7.8 Ошибки различимы по типу, а не по тексту (`errors.Is`/`errors.As`) +- [x] 7.9 Частично непонятный пакет имеет явную судьбу: что сохранено, что отброшено — видно в счётчиках +- [x] 7.10 Время нормализовано без потери зоны +- [x] 7.11 Слой выводится из данных, а не из заголовка; неопределимый слой имеет явную судьбу +- [x] 7.12 Парсер детерминирован и чист: ни `time.Now`, ни генерации id; два вызова на одном входе равны +- [x] 7.13 (добавлено после снятия блокера) Точка-интервал не схлопывается по метке: прогон архива даёт 174 координаты сна, а не 170