From cd7a4c1493b0ef4fc985c04d2a90fbdbc84fc441 Mon Sep 17 00:00:00 2001 From: Anton Vakhrushev Date: Sat, 1 Aug 2026 19:02:05 +0300 Subject: [PATCH] =?UTF-8?q?=D0=BE=D1=82=D1=80=D0=B0=D0=B1=D0=BE=D1=82?= =?UTF-8?q?=D0=B0=D0=BD=D1=8B=20=D0=BD=D0=B0=D1=85=D0=BE=D0=B4=D0=BA=D0=B8?= =?UTF-8?q?=20=D1=80=D0=B5=D0=B2=D1=8C=D1=8E=20=D0=BA=D0=BE=D0=B4=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Триаж свёл 62 сырые находки девяти проходов к 33 причинам: 3 блокера, 4 «сейчас», 2 развилки. Все закрыты регрессионными тестами. - схема точки сна определяется по самой точке, а не по индексу в исходном массиве: одна пропущенная точка меняла эпизод и сводку местами - доставка сворачивается одной транзакцией: частичное состояние было недетерминированным (восемь прогонов — семь состояний) - граница размера на распакованном теле: 400 КиБ gzip разворачивались в 400 МиБ мимо max_body_mb - столкновение — расхождение канонических форм, а не байтов; WARN с координатами объекта; payload без HTML-экранирования - выравнивание по местной метке: получасовые зоны уводили часовую выгрузку в minute - доставка из одних суточных сводок больше не отвергается целиком - единицы не переписываются молча; счётчик считает сохранённые точки - все выходы Fold логируются, исход пишется на переживающем отмену контексте Четыре развилки вынесены блокерами в беклог. --- cmd/healthlog/serve.go | 2 +- docs/architecture.md | 20 +- docs/backlog/README.md | 4 + docs/backlog/edinicy-metriki-v-razreze.md | 53 ++++ docs/backlog/nerazobrannye-sekcii-dostavki.md | 58 +++++ docs/backlog/otvet-i-svyortka.md | 80 ++++++ docs/backlog/pravilo-sliyaniya-tochek.md | 65 +++++ internal/canon/canon.go | 22 ++ internal/fold/fold.go | 141 +++++++++-- internal/fold/fold_test.go | 2 +- internal/fold/log_test.go | 115 +++++++++ internal/hae/hae.go | 136 +++++----- internal/hae/layer.go | 18 +- internal/hae/regress_test.go | 147 +++++++++++ internal/httpapi/httpapi_test.go | 42 +++- internal/httpapi/ingest.go | 26 +- internal/ingest/ingest_test.go | 2 +- internal/store/bucket.go | 235 ++++++++++++------ internal/store/delivery.go | 10 +- internal/store/regress_test.go | 232 +++++++++++++++++ .../specs/parsing/spec.md | 44 ++++ .../specs/storage/spec.md | 40 ++- .../changes/razbor-metrik-v-obekty/tasks.md | 21 +- 23 files changed, 1340 insertions(+), 175 deletions(-) create mode 100644 docs/backlog/edinicy-metriki-v-razreze.md create mode 100644 docs/backlog/nerazobrannye-sekcii-dostavki.md create mode 100644 docs/backlog/otvet-i-svyortka.md create mode 100644 docs/backlog/pravilo-sliyaniya-tochek.md create mode 100644 internal/fold/log_test.go create mode 100644 internal/hae/regress_test.go create mode 100644 internal/store/regress_test.go diff --git a/cmd/healthlog/serve.go b/cmd/healthlog/serve.go index c482949..b27cc59 100644 --- a/cmd/healthlog/serve.go +++ b/cmd/healthlog/serve.go @@ -52,7 +52,7 @@ func runServe(args []string) error { } handler := httpapi.New(httpapi.Options{ - Ingest: ingest.New(arch, st, fold.New(arch, st, log), log), + Ingest: ingest.New(arch, st, fold.New(arch, st, int64(cfg.Ingest.MaxBodyMB)<<20, log), log), Log: log, WriteTokens: cfg.Auth.WriteTokens, MaxBodyMB: cfg.Ingest.MaxBodyMB, diff --git a/docs/architecture.md b/docs/architecture.md index 6f62417..31178e3 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -230,6 +230,18 @@ HRV); у накопительных — только `date`. Поэтому то безопасности, исход разбора виден в логе, в `delivery.parse_status` и в `/stats`, а доразобрать их можно командой `reindex`. +- **413** — тело больше допустимого. Граница стоит на **распакованном** + потоке, а не только на сжатом: `MaxBytesReader` поверх `r.Body` ограничивает + то, что приехало по сети, а в память попадает то, что из этого развернулось. + Измерено: 400 КиБ сжатого тела давали 400 МиБ и гигабайт выделений при + лимите в мегабайт. Потолок степени сжатия gzip около 1030:1, так что при + штатных 64 МиБ речь о десятках гигабайт на запрос, и параллельные + складываются. Цена отказа здесь наивысшая в проекте: приём — единственное + место, где поток вообще существует, и доставка, не попавшая в архив, не + попадает в журнал. Та же граница действует при чтении тела из архива — + иначе тело между двумя границами принималось бы с `200`, а потом вечно + валилось бы при каждой пересборке. + Причина такого разделения: неизвестно, шлёт ли HAE отклонённый пакет повторно при периоде «Since Last Sync». Если не шлёт, строгий приём означал бы дыру в истории. Многоуровневая синхронизация страхует тот же риск с другой @@ -340,11 +352,11 @@ HAE. Значит для него доставки не хвост журнал ``` delivery(id, received_at, automation_name, automation_id, aggregation, period, session_id, bytes, sha256, raw_path, parse_status, points, - headers) + headers, derived_layer) -bucket(metric, layer, hour_utc, hash, points_count, first_ts, last_ts, - units, payload BLOB, first_delivery_id, updated_at, sealed) - PK (metric, layer, hour_utc) +bucket(metric, layer, hour_utc, units, payload BLOB, content_hash, points, + first_ts, last_ts, first_delivery_id, sealed, created_at, updated_at) + PK (metric, layer, hour_utc) WITHOUT ROWID workout(id PK, name, start_utc, end_utc, tz_offset, duration_sec, payload JSON, delivery_id, updated_at) diff --git a/docs/backlog/README.md b/docs/backlog/README.md index 03787e0..f6ac3d1 100644 --- a/docs/backlog/README.md +++ b/docs/backlog/README.md @@ -12,6 +12,10 @@ одному — прерывать поток ради каждого дороже, чем накопить. ## блокеры +- [Разнести ответ приёма и свёртку доставки](otvet-i-svyortka.md) — синхронная свёртка не помещается в write_timeout: широкие проходы получают обрыв вместо 200 +- [Правило выбора победителя при столкновении точек](pravilo-sliyaniya-tochek.md) — полнота считается числом ключей, поэтому мусорные поля бьют измерение +- [Единицы метрики: часть координаты или свойство объекта](edinicy-metriki-v-razreze.md) — смена единиц в настройках HAE делит час на точки в разных единицах +- [Судьба доставки, у которой разобрана не вся секция data](nerazobrannye-sekcii-dostavki.md) — тело с одним stateOfMind помечается parsed, а ретеншен снесёт его как разобранное ## высокий - [Разбор метрик в часовые объекты](razbor-metrik-v-obekty.md) — Доставки копятся непрозрачными телами — точек в хранилище нет вовсе, всё остальное упирается в это diff --git a/docs/backlog/edinicy-metriki-v-razreze.md b/docs/backlog/edinicy-metriki-v-razreze.md new file mode 100644 index 0000000..3ce40c9 --- /dev/null +++ b/docs/backlog/edinicy-metriki-v-razreze.md @@ -0,0 +1,53 @@ +# Единицы метрики: часть координаты или свойство объекта + +**Приоритет:** блокеры + +Вынуто ревью кода задачи `razbor-metrik-v-obekty` (профиль `deep`, найдено +тремя проходами независимо). + +## Что решить + +Единицы измерения живут колонкой объекта — одна на все точки часа. Что делать, +когда в тот же час приезжают точки в **других** единицах. + +## Почему это не мелочь + +Внутри точки единиц нет: проверено на 89 доставках, поле `units` не +встретилось ни разу, оно живёт только на уровне метрики. Значит у точки, +сохранённой раньше, не остаётся **ничего**, по чему её единицы восстановимы — +кроме сырого архива, пока он жив. + +Смена реальна и не требует злого умысла: переключатель единиц в настройках +HAE, смена локали телефона, переименование единицы в новой версии приложения. + +## Варианты и цена + +**(1) Единицы — часть координаты объекта** (`метрика + слой + единицы + час`). +Цена: миграция, ключ шире, каталог разрезов обязан показывать единицы. Зато +точки никогда не подписаны чужим — разные единицы просто разные ряды. + +**(2) Первое непустое побеждает, расхождение — `WARN` и счётчик** (сделано +сейчас как безопасное умолчание). +Цена: нулевая, уже работает. Но час, начавшийся в километрах, останется +километровым навсегда, даже если телефон окончательно переехал на мили. + +**(3) Хранить обе величины у объекта** (`units` и `units_seen`). +Цена: малая, но это откладывание решения: читателю всё равно придётся +выбирать. + +## Рекомендация + +**(1)**, если единицы вообще могут меняться на живом потоке — а они могут. +Сегодняшний (2) безопасен, но оставляет систематическую ложь в подписи. + +## Что стоит без решения + +Ничего не теряется: сохранённые единицы больше **не перезаписываются** молча +(было — пришедшие всегда побеждали), расхождение считается в +`MergeStats.UnitsConflicts` и даёт `WARN` в свёртке. + +Отдельная тонкость, которую надо закрыть тем же решением: единицы не входят в +хеш содержимого, поэтому доставка с теми же точками и исправленными единицами +уходит по ветке «ничего не изменилось» и колонку не трогает. То есть единицы +сегодня обновляются тогда и только тогда, когда изменилось содержимое, — +правило, которого никто не формулировал. diff --git a/docs/backlog/nerazobrannye-sekcii-dostavki.md b/docs/backlog/nerazobrannye-sekcii-dostavki.md new file mode 100644 index 0000000..aff11da --- /dev/null +++ b/docs/backlog/nerazobrannye-sekcii-dostavki.md @@ -0,0 +1,58 @@ +# Судьба доставки, у которой разобрана не вся секция data + +**Приоритет:** блокеры + +Вынуто ревью кода задачи `razbor-metrik-v-obekty` (профиль `deep`, проход +негативного пространства). **Решить до реализации ретеншена.** + +## Что решить + +Разбор читает только `data.metrics`. Доставка, состоящая из `workouts`, +`stateOfMind`, `symptoms` или `ecg`, помечается `parse_status=parsed` с нулём +точек — неотличимо от доставки с пустой секцией метрик. + +## Почему это блокер, а не задача + +Само по себе это некритично: секции пока не разбираются сознательно, тела +лежат в архиве. Опасность в **сцепке с ретеншеном**. + +Ретеншен по замыслу срезает архив до следующего проверенного экспорта. Если он +будет ориентироваться на `parse_status`, он снесёт тела, которые числятся +разобранными, — а для `stateOfMind` это необратимо: **в экспорте Apple его +нет** (находка 46), доставки HAE для него единственный источник. + +То есть цена ошибки здесь не «придётся пересобрать», а «истории состояния +разума больше не существует». + +## Варианты и цена + +**(1) Разбор возвращает список верхнеуровневых ключей `data`, которые он не +покрыл; статус `partial`.** +Цена: малая — один `json.Decoder` по верхнему уровню, без чтения содержимого. +Ретеншен и `reindex` получают честный признак, а в логе появляется момент, +когда поток принёс новую секцию (ровно то, ради чего заведена задача +`proverka-novyh-sekcij`). + +**(2) Считать `parsed` только доставку, у которой разобрано всё; остальные — +`pending`.** +Цена: малая, но `pending` перестаёт означать «ещё не смотрели», и подбор +зависших доставок теряет свой признак. + +**(3) Ничего не менять, но ретеншену запретить смотреть на `parse_status` — +резать только по дате экспорта.** +Цена: нулевая сейчас, но ретеншен становится тупым и не защищает от «тело +разобрано неверно, а мы его уже срезали». + +## Рекомендация + +**(1).** Список непокрытых ключей — дешёвая честность, и он же закрывает +задачу «не пропустить момент, когда поедет новая секция». + +## Что стоит без решения + +Ничего: сегодня ретеншена нет, тела не удаляются. Задача — не дать сцепке +сложиться позже. + +Связано: [retenshen-syrogo-arhiva](retenshen-syrogo-arhiva.md) — решить **до** +неё; [proverka-novyh-sekcij](proverka-novyh-sekcij.md) — тот же признак закрыл +бы и её. diff --git a/docs/backlog/otvet-i-svyortka.md b/docs/backlog/otvet-i-svyortka.md new file mode 100644 index 0000000..27341dc --- /dev/null +++ b/docs/backlog/otvet-i-svyortka.md @@ -0,0 +1,80 @@ +# Разнести ответ приёма и свёртку доставки + +**Приоритет:** блокеры + +Вынуто ревью кода задачи `razbor-metrik-v-obekty` (профиль `deep`, находка №4 +триажа, severity major). + +## Что решить + +Свёртка выполняется **синхронно внутри обработчика запроса**, поэтому время +ответа равно времени свёртки. Вопрос: разносить ли их, и какой ценой. + +## Оракул: измерено + +`WriteTimeout` в Go ставится в `readRequest` — **до** чтения тела и до вызова +обработчика (`net/http/server.go:993-997`, прочитано в исходниках). Значит +30 секунд по умолчанию это бюджет на всё сразу: дочитать до 64 МиБ по +мобильной сети, сделать `fsync` архива, вставить доставку и свернуть. + +Воспроизведено минимальной программой: сервер с `WriteTimeout=200ms`, +обработчик спит 500 мс. + +``` +handler: WriteHeader(200), body Write err= +client: elapsed=501ms err=EOF +``` + +Сервер считает, что отдал `200` — ошибки записи не видно, ответ ушёл в буфер и +сбрасывается позже. Клиент получил обрыв. Код обработчика этого не видит, а +`accessLog` честно запишет `status_code=200`: единственный сегодняшний канал +наблюдаемости в этом сценарии врёт. + +Стоимость свёртки измерена **до** перехода на одну транзакцию на доставку: + +| тело | объектов | свёртка | +|---|---|---| +| 80 КиБ | 1001 | 815 мс | +| 323 КиБ | 4001 | 3.07 с | +| 1302 КиБ | 16001 | 11.07 с | + +Одна транзакция на доставку убрала около 0.7 мс на объект (прогон живого +архива ускорился с 64 до 52 секунд), но порядок величины остался: широкая +доставка по-прежнему измеряется секундами. + +Бьёт это по **широким проходам** — `Today`, `Previous 7 Days`, ручной +экспорт, — то есть ровно по тем, ради которых заведён инвариант «дыры +закрываются сами». + +## Варианты и цена + +**(а) Отвечать `200` сразу после архивации и учёта; свёртка — воркером в +порядке журнала, с подбором `pending` при старте.** +Цена: средняя — воркер, очередь, подбор при старте. Бонусом закрываются ещё +две дыры: параллельные доставки одной автоматизации перестают гонять +наследование слоя (сейчас вторая может не найти слоя первой и уйти в +`failed`), и доставка, застрявшая в `pending` из-за сбоя записи, наконец +кем-то подбирается. + +**(б) Поднять `write_timeout` до согласованного с `foldTimeout`.** +Цена: малая. Но худший случай (64 МиБ) всё равно минуты, и молчание +`accessLog` остаётся. + +**(в) Оставить как есть, задокументировав потолок размера доставки.** +Цена: нулевая. Широкие проходы продолжают рваться. + +## Рекомендация + +**(а).** Единственный вариант, который решает причину, а не симптом, и попутно +снимает две смежные находки. Он же приближает `reindex`: подбор `pending` при +старте — его половина. + +## Что стоит без решения + +Ничего: свёртка работает, просто рискует не уложиться в таймаут на самых +широких доставках. Данные при этом не теряются — тело ложится в архив **до** +свёртки. + +Связано: [reindex-iz-arhiva](reindex-iz-arhiva.md), +[stats-nablyudaemost](stats-nablyudaemost.md) — метка «ответ не уложился в +таймаут» должна попасть туда. diff --git a/docs/backlog/pravilo-sliyaniya-tochek.md b/docs/backlog/pravilo-sliyaniya-tochek.md new file mode 100644 index 0000000..a1e52bc --- /dev/null +++ b/docs/backlog/pravilo-sliyaniya-tochek.md @@ -0,0 +1,65 @@ +# Правило выбора победителя при столкновении точек + +**Приоритет:** блокеры + +Вынуто ревью кода задачи `razbor-metrik-v-obekty` (профиль `deep`, находка №7 +триажа, severity major). Трогает записанный инвариант «при столкновении +выигрывает более полная точка». + +## Что решить + +Как выбирать победителя, когда по одним координатам приехали разные +содержимые. Сегодняшнее правило измеримо неверно в двух местах. + +## Оракул: измерено + +**Полнота — это счётчик ключей**, чьё значение не `null`, не `""` и не пусто. +Поэтому `{}`, `[]`, `0` и `false` считаются содержательными: + +``` +точка {"qty":0,"a":0,"b":0,"c":{},"d":[]} полнота 5 +точка {"date":"…","qty":123.4} полнота 2 +``` + +Вторая точка — настоящее измерение — проигрывает первой и стирается +безвозвратно. Восстановить можно только из сырого архива, пока он жив. + +**Тай-брейк при равной полноте** — лексикографический порядок канонических +форм, то есть исход зависит от самого значения. Для накопительной метрики, +которую досчитывают задним числом, это систематическая победа **меньшего** +числа: `qty:1` бьёт `qty:100`. Недосчёт, неотличимый от нормы. + +Частота на живом потоке **не измерена** — стоит прогнать по архиву до решения. + +## Варианты и цена + +**(1) Полнота по множеству ключей: побеждает надмножество; несравнимые +множества — объединять поля, а не выбирать точку целиком.** +Цена: средняя — правка спеки, `resolve` и тестов. Самое честное прочтение +инварианта «ничего не теряем молча»: при несравнимых наборах не теряется +ничего вообще. + +**(2) Оставить счёт ключей, но не считать содержательными `{}`, `[]`, `0`, +`false`, `""`.** +Цена: малая. Правило остаётся играбельным (точку с лишними непустыми полями +никто не мешает прислать), и произвол тай-брейка не решён. + +**(3) При равной полноте побеждает точка более поздней доставки по +`received_at`.** +Цена: малая, но спека должна признать зависимость от журнала, и нужен +детерминированный порядок **внутри** одной доставки — а именно там измерены +столкновения (31 координата в 21 доставке из 89). + +## Рекомендация + +**(1).** Объединение полей при несравнимых наборах — единственный вариант, при +котором столкновение вообще перестаёт быть выбором «кого потерять». + +## Что стоит без решения + +Правило работает и теперь **наблюдаемо**: столкновение считается расхождением +канонических форм (а не байтов, как было), даёт `WARN` с координатами объекта +и счётчик `overwrites`. Если правило начнёт терять — это станет видно в логе, +а не через месяцы при сверке с экспортом Apple. + +Связано: `openspec/specs/storage` → «Разрешение столкновений по полноте». diff --git a/internal/canon/canon.go b/internal/canon/canon.go index 04cd5ed..0d2b216 100644 --- a/internal/canon/canon.go +++ b/internal/canon/canon.go @@ -85,6 +85,28 @@ func HashAll(raws [][]byte) (string, error) { return hex.EncodeToString(h.Sum(nil)), nil } +// Equal говорит, одинаковы ли значения с точностью до канонической формы: +// порядка ключей и дребезга последнего разряда. +// +// Именно это, а не побайтовое равенство, является отношением «то же самое» в +// хранилище. Байты нестабильны: порядок ключей в JSON от HAE меняется между +// доставками, а числа расходятся последним разрядом double при одинаковом +// измерении. +func Equal(a, b []byte) bool { + if bytes.Equal(a, b) { + return true + } + fa, err := Form(a) + if err != nil { + return false + } + fb, err := Form(b) + if err != nil { + return false + } + return bytes.Equal(fa, fb) +} + // Completeness — мера полноты точки: сколько значащих полей она несёт. // // Нужна правилу разрешения столкновений: по одним координатам приезжают точки diff --git a/internal/fold/fold.go b/internal/fold/fold.go index 407102d..b629a56 100644 --- a/internal/fold/fold.go +++ b/internal/fold/fold.go @@ -13,44 +13,73 @@ import ( "fmt" "io" "log/slog" + "strings" + "time" "git.vakhrushev.me/av/healthlog/internal/archive" "git.vakhrushev.me/av/healthlog/internal/hae" "git.vakhrushev.me/av/healthlog/internal/store" ) -// maxBodyBytes — граница размера тела при чтении из архива. +// defaultMaxBodyBytes — граница размера тела при чтении из архива, когда +// вызывающий свою не задал. // -// Наблюдалось 42 МиБ; сотня даёт запас втрое и при этом не даёт битому или -// враждебному архивному файлу выесть память процесса. Граница явная, потому -// что молчаливое «сколько дадут» — это отказ, который проявится только на -// пике потока. -const maxBodyBytes = 100 << 20 +// Наблюдалось 42 МиБ. Граница явная, потому что молчаливое «сколько дадут» — +// это отказ, который проявится только на пике потока: битый или враждебный +// файл архива выест память процесса. +const defaultMaxBodyBytes = 64 << 20 // Service сворачивает доставки в часовые объекты. type Service struct { arch *archive.Archive store *store.Store log *slog.Logger + + // maxBody — та же граница, что у приёма, и это существенно. Тело между + // границами приёма и свёртки было бы принято со статусом 200, легло бы в + // архив и потом вечно валилось бы при каждой пересборке. + maxBody int64 } -// New собирает свёртку. -func New(arch *archive.Archive, st *store.Store, log *slog.Logger) *Service { - return &Service{arch: arch, store: st, log: log.With("capability", "fold")} +// New собирает свёртку. maxBody — граница размера РАСПАКОВАННОГО тела; ноль +// означает умолчание. +func New(arch *archive.Archive, st *store.Store, maxBody int64, log *slog.Logger) *Service { + if maxBody <= 0 { + maxBody = defaultMaxBodyBytes + } + return &Service{ + arch: arch, + store: st, + maxBody: maxBody, + log: log.With("capability", "fold"), + } } +// finishTimeout — сколько отводится записи исхода свёртки. +// +// Исход пишется на контексте, ПЕРЕЖИВАЮЩЕМ отмену исходного: иначе при +// срабатывании дедлайна свёртки запись статуса гарантированно провалится, и +// доставка навсегда останется `pending` — притом что часть объектов уже +// записана. То есть ровно в том случае, ради которого дедлайн и заведён, +// учёт разошёлся бы с содержимым витрины. +const finishTimeout = 10 * time.Second + // Stats — итог свёртки одной доставки. type Stats struct { Metrics int Points int + Stored int Buckets int Unchanged int Overwrites int SealedHits int + UnitsConflicts int SkippedNoTime int SkippedMalformed int + SkippedBadEnd int Layer string LayerMismatch bool + Collisions []store.Collision } // Fold разбирает тело доставки и раскладывает точки по часовым объектам. @@ -63,6 +92,10 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { d, err := s.store.DeliveryForParse(ctx, deliveryID) if err != nil { + // Ни один выход с ошибкой не молчит: вызывающий эту ошибку сознательно + // отбрасывает, и молчащий путь означал бы доставку вообще без событий + // в журнале — её не найти ни по какому запросу. + s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID) return stats, err } @@ -74,6 +107,7 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { fallback, err := s.store.LastDerivedLayer(ctx, d.AutomationID, d.ReceivedAt, d.ID) if err != nil { + s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID) return stats, err } @@ -87,8 +121,10 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { } stats.Metrics = parsed.Metrics + stats.Points = len(parsed.Points) stats.SkippedNoTime = parsed.SkippedNoTime stats.SkippedMalformed = parsed.SkippedMalformed + stats.SkippedBadEnd = parsed.SkippedBadEnd stats.Layer = string(parsed.Layer) stats.LayerMismatch = parsed.LayerMismatch @@ -98,13 +134,16 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { return stats, err } - stats.Points = merge.Points + stats.Stored = merge.Stored stats.Buckets = merge.Buckets stats.Unchanged = merge.Unchanged stats.Overwrites = merge.Overwrites stats.SealedHits = merge.SealedHits + stats.UnitsConflicts = merge.UnitsConflicts + stats.Collisions = merge.Collisions - if err := s.store.FinishParse(ctx, deliveryID, store.ParseDone, int64(stats.Points), stats.Layer); err != nil { + if err := s.finish(ctx, deliveryID, store.ParseDone, int64(stats.Points), stats.Layer); err != nil { + s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID) return stats, err } @@ -112,25 +151,59 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) { return stats, nil } +// logResult — единственный логирующий чекпоинт свёртки. +// +// Все признаки идут АТРИБУТАМИ всегда, а уровень выбирается отдельно. Раньше +// признаки жили только в тексте сообщения, и `switch` их терял: доставка, +// одновременно задевшая запечатанный час и разошедшаяся с заголовком, не +// оставляла следа о слое вовсе — а это единственный индикатор того, что слой +// выводится неправильно. +// +// Значений точек и имён устройств здесь нет и быть не может: данные о здоровье +// чувствительнее токенов. Координаты столкновений — метрика, слой, час — не +// значения. func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) { + skipped := st.SkippedNoTime + st.SkippedMalformed + st.SkippedBadEnd + attrs := []any{ "delivery_id", deliveryID, "metrics", st.Metrics, "points", st.Points, + "stored", st.Stored, "buckets", st.Buckets, "unchanged", st.Unchanged, "overwrites", st.Overwrites, + "units_conflicts", st.UnitsConflicts, + "sealed_hits", st.SealedHits, + "skipped", skipped, "skipped_no_time", st.SkippedNoTime, "skipped_malformed", st.SkippedMalformed, + "skipped_bad_end", st.SkippedBadEnd, "layer", st.Layer, + "layer_mismatch", st.LayerMismatch, + } + if len(st.Collisions) > 0 { + attrs = append(attrs, "collisions", formatCollisions(st.Collisions)) } + // Доставка, у которой отброшены ВСЕ точки, — это сломавшийся формат, а не + // штатная работа. Без этого условия смена формата метки выглядела бы как + // здоровый поток: 200, parsed, INFO, points=0. + allSkipped := st.Points == 0 && skipped > 0 + switch { + case st.Overwrites > 0: + // Единственное наблюдение, по которому проверяется правило слияния. + // В INFO оно тонуло: поток идёт раз в пять минут. + s.log.WarnContext(ctx, "delivery folded, points overwritten", attrs...) + case st.UnitsConflicts > 0: + s.log.WarnContext(ctx, "delivery folded, units differ from stored", attrs...) case st.SealedHits > 0: // Досчёт часа, в который его уже не ждали: единственное наблюдение, по // которому вообще можно судить о глубине досчёта. - s.log.WarnContext(ctx, "delivery folded, sealed hour changed", - append(attrs, "sealed_hits", st.SealedHits)...) + s.log.WarnContext(ctx, "delivery folded, sealed hour changed", attrs...) + case allSkipped: + s.log.WarnContext(ctx, "delivery folded, all points skipped", attrs...) case st.LayerMismatch: // Расхождение сверяется только с надёжным заголовком: `Default` не // означает режима, и сравнение с ним давало бы WARN на каждой доставке. @@ -140,18 +213,50 @@ func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) { } } +// formatCollisions превращает координаты столкновений в строку для лога. +func formatCollisions(cs []store.Collision) string { + parts := make([]string, 0, len(cs)) + for _, c := range cs { + parts = append(parts, c.Metric+"/"+c.Layer+"@"+store.FormatTime(c.HourUTC)) + } + return strings.Join(parts, " ") +} + +// keepLayer — значение слоя, означающее «оставить как было». +const keepLayer = "" + +// finish записывает исход разбора на контексте, переживающем отмену исходного. +func (s *Service) finish(ctx context.Context, deliveryID, status string, points int64, layer string) error { + ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), finishTimeout) + defer cancel() + + if err := s.store.FinishParse(ctx, deliveryID, status, points, layer); err != nil { + return fmt.Errorf("запись исхода разбора: %w", err) + } + return nil +} + // fail отмечает доставку неразобранной. Тело остаётся в архиве, и её подберёт // пересборка — приём при этом не затрагивается: сохранили значит приняли. func (s *Service) fail(ctx context.Context, deliveryID string, cause error) { level := slog.LevelError - if errors.Is(cause, hae.ErrLayerUnknown) { + switch { + case errors.Is(cause, hae.ErrLayerUnknown): // Слой не определился — это не поломка, а ожидаемый исход для доставки // без плотных метрик. Тело ждёт пересборки. level = slog.LevelWarn + case errors.Is(cause, hae.ErrMalformed): + // Непонятое содержимое от отправителя — норма жизни, разбирать нечего. + // На границе приёма такой же отказ уходит в DEBUG; два разных уровня у + // одного класса ошибки давали бы постоянный ERROR-шум. + 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 { + // Слой НЕ затирается: доставка могла свернуться успешно раньше, и пустая + // строка здесь оборвала бы цепочку наследования, то есть изменила бы + // результат пересборки журнала. + if err := s.finish(ctx, deliveryID, store.ParseFailed, 0, keepLayer); err != nil { s.log.ErrorContext(ctx, "delivery parse status not recorded", "error", err, "delivery_id", deliveryID) } } @@ -163,12 +268,12 @@ func (s *Service) readBody(rawPath string) ([]byte, error) { } defer func() { _ = r.Close() }() - body, err := io.ReadAll(io.LimitReader(r, maxBodyBytes+1)) + body, err := io.ReadAll(io.LimitReader(r, s.maxBody+1)) if err != nil { return nil, fmt.Errorf("чтение тела из архива: %w", err) } - if len(body) > maxBodyBytes { - return nil, fmt.Errorf("тело больше %d байт", maxBodyBytes) + if int64(len(body)) > s.maxBody { + return nil, fmt.Errorf("тело больше %d байт", s.maxBody) } return body, nil } diff --git a/internal/fold/fold_test.go b/internal/fold/fold_test.go index d125adc..5bf36d1 100644 --- a/internal/fold/fold_test.go +++ b/internal/fold/fold_test.go @@ -28,7 +28,7 @@ func newFold(t *testing.T) (*fold.Service, *archive.Archive, *store.Store) { t.Cleanup(func() { _ = st.Close() }) log := slog.New(slog.DiscardHandler) - return fold.New(arch, st, log), arch, st + return fold.New(arch, st, 0, log), arch, st } // deliver кладёт тело в архив и заводит доставку — ровно то, что делает приём. diff --git a/internal/fold/log_test.go b/internal/fold/log_test.go new file mode 100644 index 0000000..b296eec --- /dev/null +++ b/internal/fold/log_test.go @@ -0,0 +1,115 @@ +package fold_test + +import ( + "bytes" + "context" + "encoding/json" + "log/slog" + "path/filepath" + "strings" + "testing" + + "git.vakhrushev.me/av/healthlog/internal/archive" + "git.vakhrushev.me/av/healthlog/internal/fold" + "git.vakhrushev.me/av/healthlog/internal/store" +) + +// Данные о здоровье чувствительнее токенов, и требование «значения точек не в +// логах» до сих пор не проверялось ничем: все тестовые логгеры выбрасывали +// записи. Тест ловит записи и смотрит на них. +func TestFoldНеПишетЗначенийВЛог(t *testing.T) { + t.Parallel() + + var buf bytes.Buffer + log := slog.New(slog.NewJSONHandler(&buf, &slog.HandlerOptions{Level: slog.LevelInfo})) + + 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() }) + + f := fold.New(arch, st, 0, log) + + // Значения и имя устройства выбраны так, чтобы их нельзя было спутать ни с + // чем: если они окажутся в логе, это будет видно. + body := []byte(`{"data":{"metrics":[{"name":"step_count","units":"count","data":[` + + `{"date":"2025-06-05 10:00:00 +0300","qty":424242.7,"source":"Секретные Часы Антона"},` + + `{"date":"2025-06-05 10:01:00 +0300","qty":313131.9,"source":"Секретные Часы Антона"},` + + `{"date":"2025-06-05 10:02:00 +0300","qty":151515.1,"source":"Секретные Часы Антона"}` + + `]}]}}`) + + deliver(t, arch, st, "d1", "Minutes", "auto-1", body) + if _, err := f.Fold(context.Background(), "d1"); err != nil { + t.Fatalf("свёртка: %v", err) + } + + logged := buf.String() + if logged == "" { + t.Fatal("свёртка не записала ни одной строки — чекпоинт молчит") + } + + for _, secret := range []string{"424242", "313131", "151515", "Секретные Часы"} { + if strings.Contains(logged, secret) { + t.Errorf("в логе оказалось %q:\n%s", secret, logged) + } + } + + // И одновременно — чекпоинт обязан нести счётчики, иначе он бесполезен. + var rec map[string]any + for line := range strings.SplitSeq(strings.TrimSpace(logged), "\n") { + if err := json.Unmarshal([]byte(line), &rec); err != nil { + t.Fatalf("строка лога не JSON: %v", err) + } + if rec["msg"] == "delivery folded" { + break + } + } + for _, attr := range []string{"delivery_id", "points", "stored", "buckets", "layer", "skipped"} { + if _, ok := rec[attr]; !ok { + t.Errorf("в записи нет атрибута %q: %v", attr, rec) + } + } +} + +// Доставка, у которой отброшены ВСЕ точки, — это сломавшийся формат, а не +// штатная работа. Без WARN смена формата метки выглядела бы как здоровый +// поток: 200, parsed, INFO, points=0. +func TestFoldВсеТочкиОтброшеныДаётWarn(t *testing.T) { + t.Parallel() + + var buf bytes.Buffer + log := slog.New(slog.NewJSONHandler(&buf, &slog.HandlerOptions{Level: slog.LevelInfo})) + + 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() }) + + f := fold.New(arch, st, 0, log) + + // Метки в формате, которого разбор не знает: HAE сменил формат. + body := []byte(`{"data":{"metrics":[{"name":"step_count","units":"count","data":[` + + `{"date":"2025-06-05T10:00:00Z","qty":1},{"date":"2025-06-05T10:01:00Z","qty":2}` + + `]}]}}`) + + deliver(t, arch, st, "d1", "Minutes", "auto-1", body) + if _, err := f.Fold(context.Background(), "d1"); err != nil { + t.Fatalf("свёртка: %v", err) + } + + if !strings.Contains(buf.String(), `"level":"WARN"`) { + t.Errorf("доставка без единой сохранённой точки записана не как WARN:\n%s", buf.String()) + } +} diff --git a/internal/hae/hae.go b/internal/hae/hae.go index a46c989..f86400c 100644 --- a/internal/hae/hae.go +++ b/internal/hae/hae.go @@ -85,6 +85,12 @@ type Point struct { // сдвигаются, невалидный UTF-8 → U+FFFD), и потеря не видна тестам на // фикстурах: они сравнивают разобранное с разобранным. Raw json.RawMessage + + // local — метка в исходной зоне. Наружу не отдаётся: хранится всегда UTC + // плюс офсет. Нужна выводу слоя — HAE строит сетку по МЕСТНОМУ времени, и + // выравнивание, посчитанное по UTC, объявляет часовую выгрузку минутной в + // зонах с получасовым смещением (+0530, +0545, +0930). + local time.Time } // Result — итог разбора доставки. Частичные исходы живут в счётчиках, а не в @@ -99,6 +105,10 @@ type Result struct { SkippedNoTime int // SkippedMalformed — точки, не разобравшиеся как объект JSON. SkippedMalformed int + // SkippedBadEnd — точки, у которых есть `end`, но он не разбирается. + // Вырождать такую точку в мгновенную нельзя: она схлопнулась бы с соседней + // по координате. + SkippedBadEnd int // Layer — слой, выведенный для доставки в целом (тот, что наследуют редкие // метрики). Пустой, если плотных метрик не было и наследовать было нечего. @@ -158,9 +168,15 @@ func Parse(body []byte, meta Meta) (res Result, err error) { groups := make([]group, 0, len(metrics)) for _, m := range metrics { g := decodeGroup(m, &res) - if g.spillover != nil { - groups = append(groups, *g.spillover) - g.spillover = nil + if len(g.summaries) > 0 { + groups = append(groups, group{ + metric: sleepSummaryMetric, + units: g.units, + points: g.summaries, + layer: LayerDay, + fixed: true, + }) + g.summaries = nil } if len(g.points) > 0 { groups = append(groups, g) @@ -208,9 +224,11 @@ type group struct { alignment Layer dense bool - // spillover — вторая группа, отделившаяся от этой при разделении схем под - // одним именем метрики. - spillover *group + // summaries — точки, схема которых опознана прямо при разборе и слой + // которым назначен, а не выведен (суточная сводка сна). Держатся отдельно, + // потому что в определении слоя доставки не участвуют: сводок бывает больше + // порога плотности, и их полуночные метки назначили бы всей доставке `hour`. + summaries []Point } // envelope — форма тела, ровно настолько подробная, насколько нужно разбору. @@ -254,6 +272,12 @@ func decodeMetrics(body []byte) ([]metricEnvelope, error) { // decodeGroup разбирает точки одной метрики. Точка без разбираемой метки // пропускается со счётчиком — ронять из-за неё остальную доставку незачем. +// +// Схема точки определяется ЗДЕСЬ ЖЕ, в том же проходе, где точка разобрана. +// Второй проход по исходному массиву был бы неверен: список разобранных точек +// уже отфильтрован пропусками, и любой пропуск сдвигал бы соответствие — эпизод +// сна уезжал бы под имя суточной сводки, а сводка под имя эпизода. Индексной +// корреляции между двумя списками здесь не существует по построению. func decodeGroup(m metricEnvelope, res *Result) group { g := group{metric: m.Name, units: m.Units, points: make([]Point, 0, len(m.Data))} @@ -272,30 +296,42 @@ func decodeGroup(m metricEnvelope, res *Result) group { end := start if head.End != "" { - if e, err := time.Parse(timeLayout, head.End); err == nil { - end = e + e, err := time.Parse(timeLayout, head.End) + if err != nil { + // Интервал, конец которого не читается, вырождать в точку + // нельзя: две записи с общим началом получили бы одну + // координату и одна из них исчезла бы. Пропускаем со + // счётчиком — тело остаётся в архиве. + res.SkippedBadEnd++ + continue } + end = e } _, offset := start.Zone() - g.points = append(g.points, Point{ + p := Point{ Metric: m.Name, Units: m.Units, Start: start.UTC(), End: end.UTC(), OffsetSeconds: offset, Raw: raw, - }) - } + local: start, + } - if len(g.points) == 0 { - return g - } - - // Разделение схем выполняется ДО вывода слоя: у суточной сводки слой - // назначен, и её полуночные метки не должны участвовать в голосовании. - if m.Name == "sleep_analysis" { - splitSleep(&g, m.Data) + // Под именем sleep_analysis HAE шлёт две несовместимые схемы: + // поэпизодную (start/end/value/qty) и суточную сводку + // (totalSleep/core/rem/deep/awake с меткой на местной полуночи). Имя + // sleep_analysis_summary — наше; инвариант «форма Apple не + // транслируется» это не нарушает: поля внутри точки не + // переименовываются, разделяются только имена метрик, под которыми + // HAE смешал две схемы. + if m.Name == sleepMetric && head.TotalSleep != nil { + p.Metric = sleepSummaryMetric + g.summaries = append(g.summaries, p) + continue + } + g.points = append(g.points, p) } g.dense = len(g.points) >= denseThreshold @@ -303,53 +339,11 @@ func decodeGroup(m metricEnvelope, res *Result) group { return g } -// splitSleep разводит суточную сводку сна на собственное имя метрики. -// -// Под именем sleep_analysis HAE шлёт две несовместимые схемы: поэпизодную -// (start/end/value/qty) и суточную сводку (totalSleep/core/rem/deep/awake с -// меткой на местной полуночи). Имя sleep_analysis_summary — наше; инвариант -// «форма Apple не транслируется» это не нарушает, потому что поля внутри точки -// не переименовываются, разделяются только имена метрик, под которыми HAE -// смешал две схемы. -// -// Если в группе оказались обе схемы, сводки переезжают в отдельную группу, -// возвращаемую через g.spillover. Живьём смеси не наблюдалось, но полагаться -// на это нельзя: HAE меняется между версиями приложения. -func splitSleep(g *group, raws []json.RawMessage) { - summaries := make([]Point, 0) - episodes := make([]Point, 0, len(g.points)) - - for i, p := range g.points { - if i < len(raws) && isSleepSummary(raws[i]) { - p.Metric = sleepSummaryMetric - summaries = append(summaries, p) - continue - } - episodes = append(episodes, p) - } - - if len(summaries) == 0 { - return - } - g.spillover = &group{ - metric: sleepSummaryMetric, - units: g.units, - points: summaries, - layer: LayerDay, - fixed: true, - } - g.points = episodes -} - -const sleepSummaryMetric = "sleep_analysis_summary" - -func isSleepSummary(raw json.RawMessage) bool { - var head pointHead - if err := json.Unmarshal(raw, &head); err != nil { - return false - } - return head.TotalSleep != nil -} +// Имена метрик сна: пришедшее от HAE и наше для суточной сводки. +const ( + sleepMetric = "sleep_analysis" + sleepSummaryMetric = "sleep_analysis_summary" +) // parseTime разбирает метку точки. Начало берётся из start, а при его // отсутствии — из date. Измерено: start, когда он есть, всегда совпадает с @@ -380,10 +374,16 @@ func parseTime(date, start string) (time.Time, bool) { func finestAlignment(points []Point) Layer { finest := LayerHour for _, p := range points { + // По МЕСТНОЙ метке, а не по UTC: HAE строит сетку в зоне телефона. + // В зонах с получасовым смещением (+0530 Индия, +0545 Непал, +0930 + // Аделаида) ровный местный час превращается в UTC-метку на половине, + // и вся часовая выгрузка уехала бы в слой `minute` — а там столкнулась + // бы с настоящей минутной автоматизацией на координате hh:30 и завысила + // сумму минутного слоя вдвое. switch { - case p.Start.Second() != 0 || p.Start.Nanosecond() != 0: + case p.local.Second() != 0 || p.local.Nanosecond() != 0: return LayerRaw - case p.Start.Minute() != 0: + case p.local.Minute() != 0: finest = LayerMinute } } diff --git a/internal/hae/layer.go b/internal/hae/layer.go index fb1541e..d0e530d 100644 --- a/internal/hae/layer.go +++ b/internal/hae/layer.go @@ -31,7 +31,7 @@ func assignLayers(groups []group, meta Meta, res *Result) error { if delivery == "" { delivery = res.HeaderLayer } - if delivery == "" { + if delivery == "" && needsDerivedLayer(groups) { return ErrLayerUnknown } } @@ -56,6 +56,22 @@ func assignLayers(groups []group, meta Meta, res *Result) error { return nil } +// needsDerivedLayer говорит, осталась ли в доставке хоть одна группа, которой +// слой действительно надо вывести. +// +// Доставка, состоящая из одних суточных сводок сна, слоя не требует — он у них +// назначен схемой. Отказ ради неё был бы необратим: исход детерминирован, и +// пересборка из архива повторила бы его вечно, то есть автоматизация, +// настроенная только на сон, не донесла бы данных никогда. +func needsDerivedLayer(groups []group) bool { + for _, g := range groups { + if !g.fixed { + return true + } + } + return false +} + // deliveryLayer — самый мелкий слой среди плотных метрик доставки. Пустой, // если плотных метрик нет. // diff --git a/internal/hae/regress_test.go b/internal/hae/regress_test.go new file mode 100644 index 0000000..7f10930 --- /dev/null +++ b/internal/hae/regress_test.go @@ -0,0 +1,147 @@ +package hae_test + +import ( + "errors" + "testing" + + "git.vakhrushev.me/av/healthlog/internal/hae" +) + +// Схема точки определяется по ней самой, а не по её месту в исходном массиве. +// +// Соответствие «разобранная точка ↔ исходный элемент» по индексу было неверным: +// список разобранных отфильтрован пропусками, и одна пропущенная точка сдвигала +// разметку всех последующих — эпизод сна уезжал под имя суточной сводки со +// слоем `day`, сводка оставалась под именем эпизода. Метрика и слой входят в +// координату точки, поэтому ошибка необратима: пересборка воспроизвела бы её, +// а точки из объекта не удаляются никогда. +func TestParseСхемаСнаНеЗависитОтПропущенныхТочек(t *testing.T) { + t.Parallel() + + const body = `{"data":{"metrics":[{"name":"sleep_analysis","units":"hr","data":[ + null, + {"date":"","qty":1}, + {"date":"2025-06-06 00:00:00 +0300","totalSleep":7.5,"core":4.0,"deep":1.0,"rem":2.5}, + {"date":"2025-06-05 22:04:00 +0300","start":"2025-06-05 22:04:00 +0300","end":"2025-06-05 23:16:00 +0300","value":"Во сне","qty":1.2} + ]}]}}` + + res, err := hae.Parse([]byte(body), hae.Meta{Aggregation: "Minutes"}) + if err != nil { + t.Fatalf("разбор: %v", err) + } + + got := map[string]hae.Layer{} + for _, p := range res.Points { + got[p.Metric] = p.Layer + } + + if got["sleep_analysis_summary"] != hae.LayerDay { + t.Errorf("сводка получила метрику/слой %v, ожидался sleep_analysis_summary/day", got) + } + if got["sleep_analysis"] == hae.LayerDay { + t.Errorf("эпизод получил слой day — разметка съехала: %v", got) + } + if res.SkippedNoTime+res.SkippedMalformed != 2 { + t.Errorf("пропущено %d+%d точек, ожидалось 2", + res.SkippedNoTime, res.SkippedMalformed) + } +} + +// Доставка, состоящая из одних суточных сводок, слоя не требует — он у них +// назначен схемой. Отказ ради неё был бы необратим: исход детерминирован, и +// пересборка повторяла бы его вечно, то есть автоматизация, настроенная только +// на сон, не донесла бы данных никогда. +func TestParseДоставкаИзОднихСводокСохраняется(t *testing.T) { + t.Parallel() + + body := `{"data":{"metrics":[{"name":"sleep_analysis","units":"hr","data":[` + for i := range 12 { + if i > 0 { + body += "," + } + body += `{"date":"2025-06-0` + string(rune('1'+i%9)) + ` 00:00:00 +0300","totalSleep":7.5}` + } + body += `]}]}}` + + res, err := hae.Parse([]byte(body), hae.Meta{Aggregation: "Default"}) + if err != nil { + t.Fatalf("разбор: %v", err) + } + if len(res.Points) != 12 { + t.Fatalf("точек %d, ожидалось 12", len(res.Points)) + } + for _, p := range res.Points { + if p.Metric != "sleep_analysis_summary" || p.Layer != hae.LayerDay { + t.Fatalf("точка %s/%s, ожидалось sleep_analysis_summary/day", p.Metric, p.Layer) + } + } +} + +// Слой считается по МЕСТНОЙ метке: HAE строит сетку в зоне телефона. По UTC +// ровный местный час в зоне со смещением на половину превращается в метку +// hh:30, и вся часовая выгрузка уехала бы в слой minute — где столкнулась бы с +// настоящей минутной автоматизацией и завысила сумму минутного слоя вдвое. +func TestParseСлойВЗонеСПолучасовымСмещением(t *testing.T) { + t.Parallel() + + for _, offset := range []string{"+0300", "+0530", "+0545", "+0930", "-0330"} { + t.Run(offset, func(t *testing.T) { + t.Parallel() + + body := `{"data":{"metrics":[{"name":"step_count","units":"count","data":[` + for h := range 12 { + if h > 0 { + body += "," + } + body += `{"date":"2025-06-05 ` + string(rune('0'+h/10)) + string(rune('0'+h%10)) + + `:00:00 ` + offset + `","qty":1}` + } + body += `]}]}}` + + res, err := hae.Parse([]byte(body), hae.Meta{Aggregation: "Default"}) + if err != nil { + t.Fatalf("разбор: %v", err) + } + if res.Layer != hae.LayerHour { + t.Errorf("слой %q при смещении %s, ожидался hour", res.Layer, offset) + } + }) + } +} + +// Интервал, конец которого не читается, вырождать в мгновенную точку нельзя: +// две записи с общим началом получили бы одну координату, и одна исчезла бы +// молча. Пропускаем со счётчиком — тело остаётся в архиве. +func TestParseНеразобранныйКонецПропускаетТочку(t *testing.T) { + t.Parallel() + + const body = `{"data":{"metrics":[{"name":"sleep_analysis","units":"hr","data":[ + {"date":"2025-06-05 22:04:00 +0300","start":"2025-06-05 22:04:00 +0300","end":"сломано","value":"Во сне","qty":1}, + {"date":"2025-06-05 22:04:00 +0300","start":"2025-06-05 22:04:00 +0300","end":"2025-06-06 02:21:00 +0300","value":"В кровати","qty":4} + ]}]}}` + + res, err := hae.Parse([]byte(body), hae.Meta{Aggregation: "Minutes"}) + if err != nil { + t.Fatalf("разбор: %v", err) + } + if res.SkippedBadEnd != 1 { + t.Errorf("точек с нечитаемым концом %d, ожидалась 1", res.SkippedBadEnd) + } + if len(res.Points) != 1 { + t.Errorf("точек %d, ожидалась 1: вырожденная точка схлопнулась бы с соседней", len(res.Points)) + } +} + +// Проверка, что отказ по неопределимому слою остался ровно там, где он нужен: +// у доставки с обычными метриками и без всякой опоры. +func TestParseОтказПоСлоюОсталсяДляОбычныхМетрик(t *testing.T) { + t.Parallel() + + const body = `{"data":{"metrics":[{"name":"step_count","units":"count","data":[ + {"date":"2025-06-05 10:00:00 +0300","qty":1} + ]}]}}` + + if _, err := hae.Parse([]byte(body), hae.Meta{Aggregation: "Default"}); !errors.Is(err, hae.ErrLayerUnknown) { + t.Fatalf("ошибка %v, ожидалась ErrLayerUnknown", err) + } +} diff --git a/internal/httpapi/httpapi_test.go b/internal/httpapi/httpapi_test.go index 8ede0f4..8df16fc 100644 --- a/internal/httpapi/httpapi_test.go +++ b/internal/httpapi/httpapi_test.go @@ -179,6 +179,46 @@ func TestIngestRejectsOversizedBody(t *testing.T) { } } +// Граница обязана стоять на РАСПАКОВАННОМ теле, а не только на сжатом. +// +// Измерено на прежней версии: 400 КиБ сжатого тела разворачивались в 400 МиБ, +// принимались со статусом 200 и выедали гигабайт. При потолке степени сжатия +// gzip около 1030:1 и штатных 64 МиБ речь о десятках гигабайт на запрос, а +// параллельные складываются. Цена отказа здесь наивысшая в проекте: приём — +// единственное место, где поток существует, и доставка, не попавшая в архив, +// не попадает в журнал вовсе. +func TestIngestRejectsGzipBomb(t *testing.T) { + h, st := newAPI(t, nil) + + // Лимит в тесте — 1 МиБ (см. newAPI). Тело валидное, просто очень длинное + // и сжимается почти нацело. + payload := `{"data":{"metrics":[]},"pad":"` + strings.Repeat(" ", 8<<20) + `"}` + + var buf bytes.Buffer + gz := gzip.NewWriter(&buf) + if _, err := gz.Write([]byte(payload)); err != nil { + t.Fatalf("подготовка gzip: %v", err) + } + if err := gz.Close(); err != nil { + t.Fatalf("подготовка gzip: %v", err) + } + if buf.Len() >= 1<<20 { + t.Fatalf("сжатое тело %d байт — уже больше лимита, тест проверяет не то", buf.Len()) + } + + req := httptest.NewRequest(http.MethodPost, "/api/v1/ingest", &buf) + req.Header.Set("Content-Encoding", "gzip") + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + + if rec.Code != http.StatusRequestEntityTooLarge { + t.Fatalf("статус = %d, ожидался 413: распакованное тело больше лимита", rec.Code) + } + if n, _ := st.CountDeliveries(t.Context()); n != 0 { //nolint:errcheck // проверяем только счётчик + t.Errorf("доставок = %d, ожидалось 0", n) + } +} + func TestIngestChecksTokenWhenConfigured(t *testing.T) { h, _ := newAPI(t, []string{"верный-токен"}) @@ -236,7 +276,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, fold.New(arch, st, log), log), + Ingest: ingest.New(arch, st, fold.New(arch, st, 0, log), log), Log: log, WriteTokens: writeTokens, MaxBodyMB: 1, diff --git a/internal/httpapi/ingest.go b/internal/httpapi/ingest.go index 1d0dcbd..53bd133 100644 --- a/internal/httpapi/ingest.go +++ b/internal/httpapi/ingest.go @@ -105,6 +105,18 @@ func (a *api) safeHeaders(r *http.Request) map[string][]string { return out } +// readBody читает тело запроса, при необходимости распаковывая gzip. +// +// Граница стоит на РАСПАКОВАННОМ потоке, а не только на сжатом. Иначе +// `MaxBytesReader` ограничивает то, что приехало по сети, а в память попадает +// то, что из этого развернулось: измерено — 400 КиБ сжатого тела дают 400 МиБ, +// принимаются со статусом 200 и выедают гигабайт. Потолок степени сжатия gzip +// около 1030:1, так что при штатных 64 МиБ речь о десятках гигабайт, и +// параллельные запросы складываются. +// +// Цена отказа здесь наивысшая в проекте: приём — единственное место, где поток +// вообще существует. Доставка, не попавшая в архив, не попадает и в журнал — +// телефон её не перешлёт. func readBody(w http.ResponseWriter, r *http.Request, maxBody int64) ([]byte, error) { var src io.Reader = http.MaxBytesReader(w, r.Body, maxBody) @@ -114,22 +126,30 @@ func readBody(w http.ResponseWriter, r *http.Request, maxBody int64) ([]byte, er return nil, errBadGzip } defer func() { _ = gz.Close() }() - src = gz + // Лишний байт — способ отличить «ровно по границе» от «больше неё», + // не читая всё до конца. + src = io.LimitReader(gz, maxBody+1) } body, err := io.ReadAll(src) if err != nil { return nil, err //nolint:wrapcheck // разбирается writeReadError по типу } + if int64(len(body)) > maxBody { + return nil, errTooLargeUnpacked + } return body, nil } -var errBadGzip = errors.New("тело не распаковывается как gzip") +var ( + errBadGzip = errors.New("тело не распаковывается как gzip") + errTooLargeUnpacked = errors.New("распакованное тело больше допустимого размера") +) func writeReadError(w http.ResponseWriter, err error) { var tooLarge *http.MaxBytesError switch { - case errors.As(err, &tooLarge): + case errors.As(err, &tooLarge), errors.Is(err, errTooLargeUnpacked): writeError(w, http.StatusRequestEntityTooLarge, "пакет больше допустимого размера") case errors.Is(err, errBadGzip): writeError(w, http.StatusBadRequest, "тело не распаковывается как gzip") diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index be62fc9..9a3d7d4 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -204,5 +204,5 @@ func newService(t *testing.T) (*ingest.Service, *archive.Archive, *store.Store) } log := slog.New(slog.DiscardHandler) - return ingest.New(arch, st, fold.New(arch, st, log), log), arch, st + return ingest.New(arch, st, fold.New(arch, st, 0, log), log), arch, st } diff --git a/internal/store/bucket.go b/internal/store/bucket.go index 4a56041..94930a4 100644 --- a/internal/store/bucket.go +++ b/internal/store/bucket.go @@ -46,39 +46,111 @@ type Bucket struct { // это единственное наблюдение, по которому можно судить, работает ли правило // разрешения столкновений. type MergeStats struct { - Buckets int - Points int + Buckets int + // Stored — сколько точек ДОБАВИЛОСЬ в объекты. Не то же, что число точек в + // доставке: точные повторы схлопываются, и счётчик присланных завышал бы + // содержимое витрины ровно тогда, когда точки начнут теряться. + Stored int Unchanged int Overwrites int SealedHits int + // UnitsConflicts — объекты, куда точки приехали в единицах, отличных от + // сохранённых. Молчаливая замена недопустима: внутри точки единиц нет, и + // у ранних точек не остаётся ничего, по чему их единицы восстановимы. + UnitsConflicts int + // Collisions — координаты первых столкновений, для записи в лог. Без них + // счётчик перезаписей не говорит, какая метрика и какой час пострадали. + Collisions []Collision } +// Collision — координаты объекта, где содержимое точки было перезаписано. +// Значений точек не несёт: данные о здоровье чувствительнее токенов. +type Collision struct { + Metric string + Layer string + HourUTC time.Time +} + +// maxCollisionsReported — сколько координат столкновений попадает в лог. +// Больше горсти не нужно: они нужны как зацепка для разбора, а не как отчёт. +const maxCollisionsReported = 5 + // MergePoints раскладывает точки по часовым объектам и сливает их с // сохранёнными. // // Час берётся по НАЧАЛУ точки: интервал пересекает границы часов, и любой // другой выбор сделал бы принадлежность объекту зависящей от длительности. +// Вся доставка сворачивается ОДНОЙ транзакцией, а не транзакцией на объект. +// +// Транзакция на объект давала недетерминированное частичное состояние: обход +// карты групп рандомизирован, и при отказе посреди доставки набор уже +// закоммиченных объектов каждый раз другой (измерено: восемь прогонов одной +// доставки — семь разных состояний). Это ломает инвариант «состояние +// пересобираемо»: пересборка из архива давала бы не то, что живой приём, а +// хеш-детектор после расхождения переписывал бы «неизменившееся». +// +// Заодно снимается стоимость: отдельный коммит на объект стоил около 0.7 мс, +// то есть 8.5 мс на килобайт тела. func (s *Store) MergePoints(ctx context.Context, in []IncomingPoint, deliveryID string) (MergeStats, error) { - var stats MergeStats + groups := groupByHour(in) + keys := sortedKeys(groups) - for key, group := range groupByHour(in) { - res, err := s.mergeBucket(ctx, key, group, deliveryID) - if err != nil { - return stats, err - } - stats.Buckets++ - stats.Points += len(group.points) - stats.Overwrites += res.overwrites - if res.unchanged { - stats.Unchanged++ - } - if res.sealed { - stats.SealedHits++ + var stats MergeStats + err := s.inTx(ctx, func(tx *sql.Tx) error { + // Счётчики обнуляются на каждой попытке: повтор транзакции начинает + // слияние заново, и накопленное от прошлой попытки посчиталось бы дважды. + stats = MergeStats{} + + for _, key := range keys { + group := groups[key] + res, err := mergeBucket(ctx, tx, key, group, deliveryID) + if err != nil { + return err + } + stats.Buckets++ + stats.Stored += res.stored + stats.Overwrites += res.overwrites + if res.unchanged { + stats.Unchanged++ + } + if res.sealed { + stats.SealedHits++ + } + if res.unitsConflict { + stats.UnitsConflicts++ + } + if res.overwrites > 0 && len(stats.Collisions) < maxCollisionsReported { + stats.Collisions = append(stats.Collisions, Collision{ + Metric: key.metric, Layer: key.layer, HourUTC: key.hourUTC, + }) + } } + return nil + }) + if err != nil { + return MergeStats{}, err } return stats, nil } +// sortedKeys задаёт детерминированный порядок обхода объектов доставки. +func sortedKeys(groups map[bucketKey]*pointGroup) []bucketKey { + keys := make([]bucketKey, 0, len(groups)) + for k := range groups { + keys = append(keys, k) + } + sort.Slice(keys, func(i, j int) bool { + if keys[i].metric != keys[j].metric { + return keys[i].metric < keys[j].metric + } + if keys[i].layer != keys[j].layer { + return keys[i].layer < keys[j].layer + } + return keys[i].hourUTC.Before(keys[j].hourUTC) + }) + return keys +} + // IncomingPoint — точка, пришедшая на запись. Отдельный тип от Point потому, // что несёт координаты объекта (метрику, слой, единицы), а Point живёт уже // внутри объекта и их не дублирует. @@ -119,60 +191,62 @@ func groupByHour(in []IncomingPoint) map[bucketKey]*pointGroup { } type mergeResult struct { - overwrites int - unchanged bool - sealed bool + overwrites int + stored int + unchanged bool + sealed bool + unitsConflict bool } -// mergeBucket выполняет чтение, слияние и запись одного объекта одной -// транзакцией. +// mergeBucket выполняет чтение, слияние и запись одного объекта. // -// Тройка обязана быть атомарной целиком: между чтением и записью может -// вклиниться другая доставка того же часа, и её точки пропали бы молча. -func (s *Store) mergeBucket(ctx context.Context, key bucketKey, group *pointGroup, deliveryID string) (mergeResult, error) { +// Тройка обязана быть атомарной: между чтением и записью может вклиниться +// другая доставка того же часа, и её точки пропали бы молча. Атомарность даёт +// транзакция вызывающего — она общая на всю доставку. +func mergeBucket(ctx context.Context, tx *sql.Tx, key bucketKey, group *pointGroup, deliveryID string) (mergeResult, error) { var res mergeResult - err := s.inTx(ctx, func(tx *sql.Tx) error { - res = mergeResult{} - - stored, found, err := readBucket(ctx, tx, key) - if err != nil { - return err - } - - merged, overwrites := mergePoints(stored.Points, group.points) - res.overwrites = overwrites - - hash, err := hashPoints(merged) - if err != nil { - return err - } - if found && hash == stored.Hash { - // Хеш — детектор изменений: совпал, значит писать нечего. Именно - // это делает широкие проходы синхронизации дешёвыми. - res.unchanged = true - return nil - } - res.sealed = found && stored.Sealed - - now := Now() - b := Bucket{ - Metric: key.metric, - Layer: key.layer, - HourUTC: key.hourUTC, - Units: firstNonEmpty(group.units, stored.Units), - Points: merged, - Hash: hash, - Sealed: stored.Sealed, - FirstTS: merged[0].Start, - LastTS: merged[len(merged)-1].Start, - Delivery: firstNonEmpty(stored.Delivery, deliveryID), - } - return writeBucket(ctx, tx, b, now) - }) + stored, found, err := readBucket(ctx, tx, key) if err != nil { return res, err } + + merged, overwrites := mergePoints(stored.Points, group.points) + res.overwrites = overwrites + // Считаем сохранённые точки, а не присланные: точные повторы внутри + // доставки схлопываются, и счётчик присланных систематически завышал бы + // содержимое витрины — расхождение «прислали 1000, лежит 700» было бы + // невидимо ровно тогда, когда точки начнут теряться по-настоящему. + res.stored = len(merged) - len(stored.Points) + res.unitsConflict = stored.Units != "" && group.units != "" && stored.Units != group.units + + hash, err := hashPoints(merged) + if err != nil { + return res, err + } + if found && hash == stored.Hash { + // Хеш — детектор изменений: совпал, значит писать нечего. Именно это + // делает широкие проходы синхронизации дешёвыми. + res.unchanged = true + return res, nil + } + res.sealed = found && stored.Sealed + + b := Bucket{ + Metric: key.metric, + Layer: key.layer, + HourUTC: key.hourUTC, + Units: firstNonEmpty(stored.Units, group.units), + Points: merged, + Hash: hash, + Sealed: stored.Sealed, + FirstTS: merged[0].Start, + LastTS: merged[len(merged)-1].Start, + Delivery: firstNonEmpty(stored.Delivery, deliveryID), + } + if err := writeBucket(ctx, tx, b, Now()); err != nil { + return res, err + } return res, nil } @@ -196,22 +270,26 @@ func mergePoints(stored, incoming []Point) ([]Point, int) { } byCoord := make(map[coord]Point, len(stored)+len(incoming)) - order := make([]coord, 0, len(stored)+len(incoming)) put := func(p Point) int { c := coord{start: p.Start.UnixNano(), end: p.End.UnixNano()} old, ok := byCoord[c] if !ok { byCoord[c] = p - order = append(order, c) return 0 } - if bytes.Equal(old.Raw, p.Raw) { + // Побайтовое равенство — только быстрый путь. Столкновением считается + // расхождение КАНОНИЧЕСКИХ форм: канонизация и заведена потому, что + // байты нестабильны. Из 81952 повторно приехавших точек 67534 + // различаются лишь порядком ключей, ещё 63% — последним разрядом + // double. Считай мы по байтам, счётчик перезаписей давал бы тысячи + // ложных срабатываний на каждом глубоком проходе, и настоящий отказ + // правила слияния стал бы неотличим от нормы. + if bytes.Equal(old.Raw, p.Raw) || canon.Equal(old.Raw, p.Raw) { return 0 } - winner := resolve(old, p) - byCoord[c] = winner + byCoord[c] = resolve(old, p) return 1 } @@ -223,9 +301,13 @@ func mergePoints(stored, incoming []Point) ([]Point, int) { overwrites += put(p) } - out := make([]Point, 0, len(order)) - for _, c := range order { - out = append(out, byCoord[c]) + // Порядок точек в объекте канонический и входит в хеш: его задаёт ТОЛЬКО + // сортировка ниже. Порядок обхода карты не специфицирован, и полагаться на + // него значило бы получать разные хеши для одного содержимого — тогда + // «неизменившийся» объект переписывался бы каждым глубоким проходом. + out := make([]Point, 0, len(byCoord)) + for _, p := range byCoord { + out = append(out, p) } sort.Slice(out, func(i, j int) bool { if !out[i].Start.Equal(out[j].Start) { @@ -378,14 +460,21 @@ func encodePayload(points []Point) ([]byte, error) { }) } - raw, err := json.Marshal(sp) - if err != nil { + // Через Encoder с выключенным HTML-экранированием, а не json.Marshal: + // Marshal превращает `&`, `<` и `>` внутри содержимого точки в + // \u0026, \u003c, \u003e. Порчи значений это не даёт, но обещание + // «точка хранится дословно» перестаёт быть правдой, а сравнение байтов + // при следующей доставке той же точки начинает промахиваться навсегда. + var raw bytes.Buffer + enc := json.NewEncoder(&raw) + enc.SetEscapeHTML(false) + if err := enc.Encode(sp); err != nil { return nil, fmt.Errorf("сериализация точек: %w", err) } var buf bytes.Buffer gz := gzip.NewWriter(&buf) - if _, err := gz.Write(raw); err != nil { + if _, err := gz.Write(raw.Bytes()); err != nil { return nil, fmt.Errorf("сжатие точек: %w", err) } if err := gz.Close(); err != nil { diff --git a/internal/store/delivery.go b/internal/store/delivery.go index eb78d36..b0c47f8 100644 --- a/internal/store/delivery.go +++ b/internal/store/delivery.go @@ -97,12 +97,18 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) { // Слой сохраняется здесь же, потому что он нужен следующей доставке той же // автоматизации: без плотных метрик выводить его не из чего, и наследовать // приходится от прошлого раза. +// +// Пустой layer означает «не трогать»: у неудачной свёртки слоя нет, а +// затирание оборвало бы цепочку наследования — в том числе у доставки, которая +// раньше свернулась успешно, — и результат пересборки журнала изменился бы. 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 = ? + UPDATE delivery + SET parse_status = ?, points = ?, + derived_layer = CASE WHEN ? = '' THEN derived_layer ELSE ? END WHERE id = ?` - res, err := s.db.ExecContext(ctx, q, status, points, layer, id) + res, err := s.db.ExecContext(ctx, q, status, points, layer, layer, id) if err != nil { return fmt.Errorf("update parse status: %w", err) } diff --git a/internal/store/regress_test.go b/internal/store/regress_test.go new file mode 100644 index 0000000..c0af9b3 --- /dev/null +++ b/internal/store/regress_test.go @@ -0,0 +1,232 @@ +package store_test + +import ( + "context" + "strings" + "testing" + "time" + + "git.vakhrushev.me/av/healthlog/internal/store" +) + +// Столкновением считается расхождение КАНОНИЧЕСКИХ форм, а не байтов. +// +// Канонизация и заведена потому, что байты нестабильны: из 81 952 повторно +// приехавших точек 67 534 различаются лишь порядком ключей, ещё 63% — последним +// разрядом double. Побайтовое сравнение давало тысячи ложных перезаписей на +// каждом глубоком проходе, и настоящий отказ правила слияния становился +// неотличим от нормы — при том что счётчик перезаписей объявлен единственным +// наблюдением за этим правилом. +func TestMergePointsДребезгНеСчитаетсяСтолкновением(t *testing.T) { + t.Parallel() + + cases := map[string][2]string{ + "порядок ключей": { + `{"date":"a","qty":1.5,"source":"Device A"}`, + `{"source":"Device A","qty":1.5,"date":"a"}`, + }, + "последний разряд double": { + `{"qty":0.09523182962471353}`, + `{"qty":0.09523182962471352}`, + }, + } + + for name, pair := range cases { + t.Run(name, func(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + first := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", pair[0]) + second := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", pair[1]) + + if _, err := st.MergePoints(ctx, []store.IncomingPoint{first}, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + stats, err := st.MergePoints(ctx, []store.IncomingPoint{second}, "d2") + if err != nil { + t.Fatalf("второе слияние: %v", err) + } + + if stats.Overwrites != 0 { + t.Errorf("перезаписей %d, ожидалось 0: содержимое то же с точностью до канонической формы", + stats.Overwrites) + } + if len(stats.Collisions) != 0 { + t.Errorf("координат столкновений %d, ожидалось 0", len(stats.Collisions)) + } + }) + } +} + +// Настоящее столкновение обязано оставить след с координатами объекта: одно +// число `overwrites` не говорит, какая метрика и какой час пострадали, и +// расследовать перезапись по нему нечем. +func TestMergePointsСтолкновениеОставляетКоординаты(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + rich := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", + `{"Avg":60,"context":"Отдых"}`) + poor := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", + `{"Avg":61}`) + + if _, err := st.MergePoints(ctx, []store.IncomingPoint{rich}, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + stats, err := st.MergePoints(ctx, []store.IncomingPoint{poor}, "d2") + if err != nil { + t.Fatalf("второе слияние: %v", err) + } + + if len(stats.Collisions) != 1 { + t.Fatalf("координат столкновений %d, ожидалась 1", len(stats.Collisions)) + } + c := stats.Collisions[0] + if c.Metric != "heart_rate" || c.Layer != "minute" { + t.Errorf("координаты столкновения %s/%s, ожидались heart_rate/minute", c.Metric, c.Layer) + } + if c.HourUTC.Hour() != 10 { + t.Errorf("час столкновения %v", c.HourUTC) + } +} + +// «Точка хранится дословно» читается буквально: байты как пришли. json.Marshal +// экранирует `&`, `<` и `>` внутри содержимого — порчи значений это не даёт, но +// обещание перестаёт быть правдой, а сравнение байтов при следующей доставке +// той же точки начинает промахиваться навсегда. +func TestMergePointsХранитУгловыеСкобкиИАмперсанд(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + raw := `{"qty":1,"source":"Anton & Co \"x\""}` + in := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", raw) + + if _, err := st.MergePoints(ctx, []store.IncomingPoint{in}, "d1"); err != nil { + t.Fatalf("слияние: %v", err) + } + + b, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + got := string(b.Points[0].Raw) + if got != raw { + t.Errorf("точка изменилась при хранении:\n было %s\n стало %s", raw, got) + } + for _, escaped := range []string{"\\u0026", "\\u003c", "\\u003e"} { + if strings.Contains(got, escaped) { + t.Errorf("содержимое точки экранировано как HTML (%s): %s", escaped, got) + } + } +} + +// Смена единиц не имеет права молча переподписать уже сохранённые точки: +// внутри точки единиц нет, и у ранних точек не остаётся ничего, по чему их +// единицы восстановимы. Правило «первое непустое побеждает» плюс счётчик. +func TestMergePointsСменаЕдиницНеПерезаписываетМолча(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + km := store.IncomingPoint{ + Metric: "walking_running_distance", Layer: "minute", Units: "km", + Point: store.Point{ + Start: ts(t, "2025-06-05T10:00:00Z"), End: ts(t, "2025-06-05T10:00:00Z"), + Raw: []byte(`{"qty":1}`), + }, + } + mi := km + mi.Units = "mi" + mi.Start = ts(t, "2025-06-05T10:30:00Z") + mi.End = mi.Start + mi.Raw = []byte(`{"qty":2}`) + + if _, err := st.MergePoints(ctx, []store.IncomingPoint{km}, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + stats, err := st.MergePoints(ctx, []store.IncomingPoint{mi}, "d2") + if err != nil { + t.Fatalf("второе слияние: %v", err) + } + + if stats.UnitsConflicts != 1 { + t.Errorf("расхождений единиц %d, ожидалось 1", stats.UnitsConflicts) + } + + b, err := st.Bucket(ctx, "walking_running_distance", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if b.Units != "km" { + t.Errorf("единицы объекта %q: сохранённые точки переподписаны чужими", b.Units) + } +} + +// Счётчик считает СОХРАНЁННЫЕ точки, а не присланные: точные повторы внутри +// доставки схлопываются, и счётчик присланных систематически завышал бы +// содержимое витрины — расхождение «прислали 1000, лежит 700» было бы невидимо +// ровно тогда, когда точки начнут теряться по-настоящему. +func TestMergePointsСчитаетСохранённые(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + p := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`) + stats, err := st.MergePoints(ctx, []store.IncomingPoint{p, p, p}, "d1") + if err != nil { + t.Fatalf("слияние: %v", err) + } + if stats.Stored != 1 { + t.Errorf("сохранено %d точек, ожидалась 1 (три точных повтора)", stats.Stored) + } + + stats, err = st.MergePoints(ctx, []store.IncomingPoint{p}, "d2") + if err != nil { + t.Fatalf("повтор: %v", err) + } + if stats.Stored != 0 { + t.Errorf("повтор добавил %d точек, ожидалось 0", stats.Stored) + } +} + +// Доставка сворачивается одной транзакцией, поэтому отказ посреди неё не +// оставляет частичного состояния — а значит и не зависит от порядка обхода. +// Раньше транзакция была на объект, и восемь прогонов одной доставки давали +// семь разных наборов записанных объектов. +func TestMergePointsОтказНеОставляетЧастиОбъектов(t *testing.T) { + t.Parallel() + + st := open(t) + + in := make([]store.IncomingPoint, 0, 40) + for i := range 40 { + at := ts(t, "2025-06-05T00:00:00Z").Add(time.Duration(i) * time.Hour) + in = append(in, store.IncomingPoint{ + Metric: "step_count", Layer: "minute", Units: "count", + Point: store.Point{Start: at, End: at, Raw: []byte(`{"qty":1}`)}, + }) + } + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + if _, err := st.MergePoints(ctx, in, "d1"); err == nil { + t.Fatal("слияние на отменённом контексте прошло успешно") + } + + n, err := st.CountBuckets(context.Background()) + if err != nil { + t.Fatalf("счёт объектов: %v", err) + } + if n != 0 { + t.Errorf("записано %d объектов из 40: отказ оставил частичное состояние", n) + } +} 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 242e534..691a096 100644 --- a/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md +++ b/openspec/changes/razbor-metrik-v-obekty/specs/parsing/spec.md @@ -34,6 +34,15 @@ - **THEN** точка не попадает в хранилище - **AND** факт учитывается в итоге разбора доставки +#### Scenario: Точка с нечитаемым концом интервала пропускается + +- **WHEN** точка несёт `end`, который не разбирается +- **THEN** точка не попадает в хранилище +- **AND** факт учитывается отдельным счётчиком + +Вырождать такую точку в мгновенную нельзя: две записи с общим началом получили +бы одну координату, и одна исчезла бы молча. Тело остаётся в архиве. + ### Requirement: Вывод слоя гранулярности Система SHALL выводить слой точки из **выравнивания меток времени**, а не из @@ -45,6 +54,12 @@ но **не выводится** — он назначается схемам с фиксированной гранулярностью (см. «Разделение схем под одним именем метрики»). +Выравнивание SHALL считаться по метке в **исходной зоне**, а не по метке в UTC. +HAE строит сетку по местному времени; ровный местный час в зоне со смещением на +половину (`+0530`, `+0545`, `+0930`) даёт UTC-метку на середине часа, и часовая +выгрузка целиком уехала бы в слой `minute` — где столкнулась бы с настоящей +минутной автоматизацией и завысила сумму минутного слоя вдвое. + Слоем метрики SHALL становиться **самое мелкое** выравнивание среди её меток, а не преобладающее. Измерено: у плотных метрик выравнивания перемешаны (`active_energy` — 1320 минутных меток и 21 часовая, `heart_rate` — 654 @@ -67,6 +82,11 @@ - **WHEN** в доставке у метрики не меньше десяти точек с метками - **THEN** слой определяется выравниванием её собственных меток +#### Scenario: Зона с получасовым смещением не делает часовую выгрузку минутной + +- **WHEN** метки стоят на ровном местном часе, а смещение зоны равно `+0530` +- **THEN** слой метрики `hour` + #### Scenario: Одна метка на середине часа делает метрику минутной - **WHEN** у плотной метрики десять меток стоят ровно на часе, а одна — на @@ -79,6 +99,16 @@ - **THEN** она получает самый мелкий слой среди плотных метрик этой доставки - **AND** её собственное выравнивание во внимание не принимается +#### Scenario: Доставка, где всем метрикам слой назначен схемой + +- **WHEN** в доставке нет метрик, которым слой надо выводить, — все точки + принадлежат схемам с фиксированной гранулярностью +- **THEN** доставка сохраняется, а слой доставки не выводится и не требуется + +Иначе доставка автоматизации, настроенной только на сон, отвергалась бы +целиком — и необратимо: исход детерминирован, и пересборка повторяла бы его +вечно. + #### Scenario: В доставке нет плотных метрик - **WHEN** ни у одной метрики доставки нет десяти точек @@ -163,6 +193,20 @@ Export шлёт под одним именем, чтобы одно имя оз (`totalSleep`/`core`/`rem`/`deep`/`awake` с меткой на местной полуночи). Общих полей, кроме `date` и `source`, у них нет. +Схема точки SHALL определяться по самой точке, а не по её месту в исходном +массиве. Список разобранных точек отфильтрован пропусками, и соответствие по +индексу съезжало бы от одной пропущенной точки: эпизод уезжал бы под имя +суточной сводки со слоем `day`, сводка — под имя эпизода. Метрика и слой входят +в координату, поэтому ошибка необратима — точки из объекта не удаляются, и +пересборка воспроизвела бы её. + +#### Scenario: Пропущенная точка не сдвигает разметку схем + +- **WHEN** в метрике `sleep_analysis` перед суточной сводкой стоит точка без + разбираемой метки +- **THEN** сводка всё равно сохраняется под именем `sleep_analysis_summary` со + слоем `day` + #### Scenario: Поэпизодная запись сна - **WHEN** точка `sleep_analysis` содержит поле `value` diff --git a/openspec/changes/razbor-metrik-v-obekty/specs/storage/spec.md b/openspec/changes/razbor-metrik-v-obekty/specs/storage/spec.md index c5780e1..0aa7c12 100644 --- a/openspec/changes/razbor-metrik-v-obekty/specs/storage/spec.md +++ b/openspec/changes/razbor-metrik-v-obekty/specs/storage/spec.md @@ -88,13 +88,26 @@ - **THEN** исход определяется детерминированно и не зависит от порядка воспроизведения доставок +Столкновением SHALL считаться расхождение **канонических форм**, а не байтов. +Байты нестабильны — ради этого канонизация и заведена: из 81 952 повторно +приехавших точек 67 534 различаются лишь порядком ключей, ещё 63% — последним +разрядом double. Побайтовое сравнение давало бы тысячи ложных срабатываний на +каждом глубоком проходе, и настоящий отказ правила стал бы неотличим от нормы. + #### Scenario: Столкновение с различием содержимого оставляет след - **WHEN** по одним координатам сохраняется точка, каноническая форма которой отличается от уже сохранённой - **THEN** система пишет запись уровня `WARN` без значений точки +- **AND** запись несёт координаты объекта: метрику, слой и час - **AND** увеличивает счётчик перезаписей в итоге разбора доставки +#### Scenario: Дребезг сериализации столкновением не считается + +- **WHEN** та же точка приезжает с другим порядком ключей или отличаясь + последним разрядом числа +- **THEN** счётчик перезаписей не растёт и `WARN` не пишется + Без этого следа допущение «меньше полей не значит новее» не получит ни одного наблюдения, а отказ правила будет неотличим от нормальной работы до сверки с экспортом Apple — то есть месяцами. @@ -117,7 +130,32 @@ слияний. Запись — чтение объекта, слияние точек, запись обратно. Точки из объекта -MUST NOT удаляться. +MUST NOT удаляться. Содержимое объекта SHALL сериализоваться без +HTML-экранирования: `&`, `<` и `>` внутри точки обязаны храниться теми же +байтами, какими пришли, иначе «точка хранится дословно» перестаёт быть правдой, +а сравнение с последующей доставкой той же точки промахивается навсегда. + +Доставка SHALL сворачиваться **одной транзакцией**. Транзакция на объект давала +недетерминированное частичное состояние: обход групп рандомизирован, и при +отказе посреди доставки набор уже записанных объектов каждый раз другой +(измерено: восемь прогонов одной доставки — семь разных состояний). Это ломает +инвариант «состояние пересобираемо». + +Единицы измерения MUST NOT переписываться молча: при расхождении сохранённых и +пришедших единиц остаётся сохранённое значение, факт учитывается счётчиком и +попадает в запись уровня `WARN`. Внутри точки единиц нет, и у ранее сохранённых +точек не остаётся ничего, по чему их единицы восстановимы. + +#### Scenario: Отказ посреди доставки не оставляет части объектов + +- **WHEN** свёртка доставки прерывается на середине +- **THEN** не записывается ни один объект этой доставки + +#### Scenario: Смена единиц не переподписывает сохранённые точки + +- **WHEN** в объект приезжают точки в единицах, отличных от сохранённых +- **THEN** единицы объекта остаются прежними +- **AND** факт учитывается счётчиком и записью `WARN` #### Scenario: Точки за один час ложатся в один объект diff --git a/openspec/changes/razbor-metrik-v-obekty/tasks.md b/openspec/changes/razbor-metrik-v-obekty/tasks.md index e8a5d99..a11626e 100644 --- a/openspec/changes/razbor-metrik-v-obekty/tasks.md +++ b/openspec/changes/razbor-metrik-v-obekty/tasks.md @@ -36,7 +36,7 @@ - [x] 4.4 Ключ точки — `метрика + слой + начало + конец` одной формы для всех точек: начало из `start`, иначе из `date`; конец из `end`, иначе равен началу. Час объекта — по началу. Ветвления по «классу метрики» быть не должно - [x] 4.4a Тест на фикстуре 1.3a: три записи с одной меткой и разными интервалами дают три точки, а не одну; повтор той же тройки следующей доставкой не задваивает; точка без `end` кладётся вырожденным интервалом - [x] 4.5 Хеш объекта как детектор изменений: совпал — записи нет -- [x] 4.6 `_txlock=immediate` в DSN; повтор оборачивает **всю тройку** чтение-слияние-запись; путь «хеш совпал» — под `TxOptions{ReadOnly: true}` +- [x] 4.6 `_txlock=immediate` в DSN; повтор оборачивает **всю доставку** целиком. Путь «хеш совпал» под `TxOptions{ReadOnly: true}` **не сделан** — вынесен блокером `otvet-i-svyortka` вместе с остальной ценой синхронной свёртки (измерено: 52 мс на неизменившийся плотный час под write-lock) - [x] 4.7 Распознавание занятости — `errors.As` на `*sqlite.Error`, коды 5 и 517, обёрнуто в `store` - [x] 4.8 Тест конкурентной записи: N горутин × M слияний в один `hour_utc`, проверка **суммы** точек, под `-race` - [x] 4.9 Тест идемпотентности: повторное слияние того же набора не меняет ни содержимое, ни хеш @@ -74,3 +74,22 @@ - [x] 7.11 Слой выводится из данных, а не из заголовка; неопределимый слой имеет явную судьбу - [x] 7.12 Парсер детерминирован и чист: ни `time.Now`, ни генерации id; два вызова на одном входе равны - [x] 7.13 (добавлено после снятия блокера) Точка-интервал не схлопывается по метке: прогон архива даёт 174 координаты сна, а не 170 + +## 8. Отработка ревью кода (профиль `deep`) + +Триаж свёл 62 сырые находки к 33 причинам: 3 блокера, 4 «сейчас», 2 развилки. + +- [x] 8.1 Схема точки сна определяется по самой точке, а не по индексу в исходном массиве (подтверждено шестью проходами) +- [x] 8.2 Доставка сворачивается одной транзакцией: частичное состояние было недетерминированным (8 прогонов — 7 состояний) +- [x] 8.3 Граница размера на распакованном теле: 400 КиБ gzip разворачивались в 400 МиБ мимо лимита +- [x] 8.4 Столкновение — расхождение канонических форм, а не байтов; `WARN` с координатами объекта +- [x] 8.5 `encodePayload` без HTML-экранирования: `&`, `<`, `>` хранятся дословно +- [x] 8.6 Доставка из одних суточных сводок больше не отвергается целиком +- [x] 8.7 Выравнивание считается по местной метке: получасовые зоны уводили часовую выгрузку в `minute` +- [x] 8.8 Нечитаемый `end` пропускает точку со счётчиком, а не вырождает интервал +- [x] 8.9 Единицы не переписываются молча: сохранённое побеждает, расхождение — счётчик и `WARN` +- [x] 8.10 Счётчик считает сохранённые точки, а не присланные +- [x] 8.11 Все выходы `Fold` с ошибкой логируются; исход пишется на контексте, переживающем отмену; слой не затирается +- [x] 8.12 Признаки — атрибутами всегда, уровень выбирается отдельно; `WARN` на отброшенных всех точках +- [x] 8.13 Регрессионные тесты на каждое исправление; тест «значения точек не попадают в лог» с перехватывающим handler +- [x] 8.14 Развилки вынесены блокерами: `otvet-i-svyortka`, `pravilo-sliyaniya-tochek`, `edinicy-metriki-v-razreze`, `nerazobrannye-sekcii-dostavki`