отработаны находки ревью кода

Триаж свёл 62 сырые находки девяти проходов к 33 причинам: 3 блокера,
4 «сейчас», 2 развилки. Все закрыты регрессионными тестами.

- схема точки сна определяется по самой точке, а не по индексу в исходном
  массиве: одна пропущенная точка меняла эпизод и сводку местами
- доставка сворачивается одной транзакцией: частичное состояние было
  недетерминированным (восемь прогонов — семь состояний)
- граница размера на распакованном теле: 400 КиБ gzip разворачивались в
  400 МиБ мимо max_body_mb
- столкновение — расхождение канонических форм, а не байтов; WARN с
  координатами объекта; payload без HTML-экранирования
- выравнивание по местной метке: получасовые зоны уводили часовую выгрузку
  в minute
- доставка из одних суточных сводок больше не отвергается целиком
- единицы не переписываются молча; счётчик считает сохранённые точки
- все выходы Fold логируются, исход пишется на переживающем отмену контексте

Четыре развилки вынесены блокерами в беклог.
This commit is contained in:
av
2026-08-01 19:02:05 +03:00
parent 32f21044e4
commit cd7a4c1493
23 changed files with 1340 additions and 175 deletions
+22
View File
@@ -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 — мера полноты точки: сколько значащих полей она несёт.
//
// Нужна правилу разрешения столкновений: по одним координатам приезжают точки
+123 -18
View File
@@ -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
}
+1 -1
View File
@@ -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 кладёт тело в архив и заводит доставку — ровно то, что делает приём.
+115
View File
@@ -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())
}
}
+68 -68
View File
@@ -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
}
}
+17 -1
View File
@@ -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 — самый мелкий слой среди плотных метрик доставки. Пустой,
// если плотных метрик нет.
//
+147
View File
@@ -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)
}
}
+41 -1
View File
@@ -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,
+23 -3
View File
@@ -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")
+1 -1
View File
@@ -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
}
+162 -73
View File
@@ -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 {
+8 -2
View File
@@ -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)
}
+232
View File
@@ -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 <iPhone> \"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)
}
}