From c9ddf0b4dd4e2a5db45f651674abda88a6140f6d Mon Sep 17 00:00:00 2001 From: Anton Vakhrushev Date: Sat, 1 Aug 2026 17:37:50 +0300 Subject: [PATCH] =?UTF-8?q?=D1=85=D1=80=D0=B0=D0=BD=D0=B5=D0=BD=D0=B8?= =?UTF-8?q?=D0=B5=20=D1=82=D0=BE=D1=87=D0=B5=D0=BA=20=D1=87=D0=B0=D1=81?= =?UTF-8?q?=D0=BE=D0=B2=D1=8B=D0=BC=D0=B8=20=D0=BE=D0=B1=D1=8A=D0=B5=D0=BA?= =?UTF-8?q?=D1=82=D0=B0=D0=BC=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - слияние read-modify-write в одной транзакции, ключ точки — интервал, победа более полной точки, тай-брейк по канонической форме - хеш объекта как детектор изменений: совпал — записи нет - _txlock=immediate и повтор всей тройки при SQLITE_BUSY/BUSY_SNAPSHOT - миграции через goose.Provider: пакетные SetBaseFS/SetDialect писали глобалы и давали гонку при двух Open --- internal/store/bucket.go | 486 ++++++++++++++++ internal/store/bucket_test.go | 531 ++++++++++++++++++ internal/store/store.go | 30 +- internal/store/tx.go | 90 +++ .../changes/razbor-metrik-v-obekty/tasks.md | 22 +- 5 files changed, 1143 insertions(+), 16 deletions(-) create mode 100644 internal/store/bucket.go create mode 100644 internal/store/bucket_test.go create mode 100644 internal/store/tx.go diff --git a/internal/store/bucket.go b/internal/store/bucket.go new file mode 100644 index 0000000..ace9b04 --- /dev/null +++ b/internal/store/bucket.go @@ -0,0 +1,486 @@ +package store + +import ( + "bytes" + "compress/gzip" + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "io" + "sort" + "time" + + "git.vakhrushev.me/av/healthlog/internal/canon" +) + +// Point — точка внутри часового объекта. +// +// Координаты — метрика, слой, начало и конец; у точки-измерения конец равен +// началу. Содержимое хранится исходными байтами: пересборка повторной +// сериализацией теряет литерал, и потеря не видна тестам на фикстурах. +type Point struct { + Start time.Time + End time.Time + OffsetSeconds int + Raw json.RawMessage +} + +// Bucket — часовой объект: точки одной метрики одного слоя за один час UTC. +type Bucket struct { + Metric string + Layer string + HourUTC time.Time + Units string + Points []Point + Hash string + Sealed bool + FirstTS time.Time + LastTS time.Time + Delivery string +} + +// MergeStats — что случилось при слиянии. Счётчики нужны логу на границе: +// перезапись означает, что по одним координатам приехали разные содержимые, и +// это единственное наблюдение, по которому можно судить, работает ли правило +// разрешения столкновений. +type MergeStats struct { + Buckets int + Points int + Unchanged int + Overwrites int + SealedHits int +} + +// MergePoints раскладывает точки по часовым объектам и сливает их с +// сохранёнными. +// +// Час берётся по НАЧАЛУ точки: интервал пересекает границы часов, и любой +// другой выбор сделал бы принадлежность объекту зависящей от длительности. +func (s *Store) MergePoints(ctx context.Context, in []IncomingPoint, deliveryID string) (MergeStats, error) { + var stats MergeStats + + 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++ + } + } + return stats, nil +} + +// IncomingPoint — точка, пришедшая на запись. Отдельный тип от Point потому, +// что несёт координаты объекта (метрику, слой, единицы), а Point живёт уже +// внутри объекта и их не дублирует. +type IncomingPoint struct { + Metric string + Layer string + Units string + Point +} + +type bucketKey struct { + metric string + layer string + hourUTC time.Time +} + +type pointGroup struct { + units string + points []Point +} + +func groupByHour(in []IncomingPoint) map[bucketKey]*pointGroup { + out := make(map[bucketKey]*pointGroup) + for _, p := range in { + key := bucketKey{ + metric: p.Metric, + layer: p.Layer, + hourUTC: p.Start.UTC().Truncate(time.Hour), + } + g, ok := out[key] + if !ok { + g = &pointGroup{units: p.Units} + out[key] = g + } + g.points = append(g.points, p.Point) + } + return out +} + +type mergeResult struct { + overwrites int + unchanged bool + sealed bool +} + +// mergeBucket выполняет чтение, слияние и запись одного объекта одной +// транзакцией. +// +// Тройка обязана быть атомарной целиком: между чтением и записью может +// вклиниться другая доставка того же часа, и её точки пропали бы молча. +func (s *Store) mergeBucket(ctx context.Context, 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) + }) + if err != nil { + return res, err + } + return res, nil +} + +// mergePoints сливает сохранённые точки с пришедшими по координатному ключу. +// +// При столкновении выигрывает БОЛЕЕ ПОЛНАЯ точка, а не последняя пришедшая: +// 0.66% координат различаются набором полей при одинаковом значении, и правило +// «последняя победила» стирало бы у сохранённой точки поля, которых новая не +// несёт. +// +// При равной полноте исход определяется порядком канонических форм, а не +// порядком доставок: у сохранённой точки нет провенанса, а четверть доставок +// несёт столкновения внутри себя, где время приёма общее. Свёртка по журналу +// обязана давать то же состояние, что приём в реальном времени. +// +// Точки из объекта не удаляются никогда. +func mergePoints(stored, incoming []Point) ([]Point, int) { + type coord struct { + start int64 + end int64 + } + + 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) { + return 0 + } + + winner := resolve(old, p) + byCoord[c] = winner + return 1 + } + + overwrites := 0 + for _, p := range stored { + put(p) + } + for _, p := range incoming { + overwrites += put(p) + } + + out := make([]Point, 0, len(order)) + for _, c := range order { + out = append(out, byCoord[c]) + } + sort.Slice(out, func(i, j int) bool { + if !out[i].Start.Equal(out[j].Start) { + return out[i].Start.Before(out[j].Start) + } + return out[i].End.Before(out[j].End) + }) + return out, overwrites +} + +// resolve выбирает победителя столкновения: сперва по полноте, затем +// детерминированно по канонической форме. +func resolve(a, b Point) Point { + ca, cb := canon.Completeness(a.Raw), canon.Completeness(b.Raw) + switch { + case ca > cb: + return a + case cb > ca: + return b + case canon.Less(a.Raw, b.Raw): + return a + default: + return b + } +} + +func hashPoints(points []Point) (string, error) { + raws := make([][]byte, 0, len(points)) + for _, p := range points { + raws = append(raws, p.Raw) + } + h, err := canon.HashAll(raws) + if err != nil { + return "", fmt.Errorf("хеш объекта: %w", err) + } + return h, nil +} + +func firstNonEmpty(a, b string) string { + if a != "" { + return a + } + return b +} + +func readBucket(ctx context.Context, tx *sql.Tx, key bucketKey) (Bucket, bool, error) { + const q = ` + SELECT units, payload, content_hash, sealed, first_delivery_id, first_ts, last_ts + FROM bucket WHERE metric = ? AND layer = ? AND hour_utc = ?` + + var ( + units string + payload []byte + hash string + sealed int + delivery string + firstTS string + lastTS string + ) + err := tx.QueryRowContext(ctx, q, key.metric, key.layer, FormatTime(key.hourUTC)). + Scan(&units, &payload, &hash, &sealed, &delivery, &firstTS, &lastTS) + if errors.Is(err, sql.ErrNoRows) { + return Bucket{}, false, nil + } + if err != nil { + return Bucket{}, false, fmt.Errorf("select bucket: %w", err) + } + + points, err := decodePayload(payload) + if err != nil { + return Bucket{}, false, err + } + first, err := ParseTime(firstTS) + if err != nil { + return Bucket{}, false, err + } + last, err := ParseTime(lastTS) + if err != nil { + return Bucket{}, false, err + } + return Bucket{ + Metric: key.metric, + Layer: key.layer, + HourUTC: key.hourUTC, + Units: units, + Points: points, + Hash: hash, + Sealed: sealed != 0, + Delivery: delivery, + FirstTS: first, + LastTS: last, + }, true, nil +} + +func writeBucket(ctx context.Context, tx *sql.Tx, b Bucket, now time.Time) error { + payload, err := encodePayload(b.Points) + if err != nil { + return err + } + + const q = ` + INSERT INTO bucket (metric, layer, hour_utc, units, payload, content_hash, + points, first_ts, last_ts, first_delivery_id, sealed, + created_at, updated_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + ON CONFLICT (metric, layer, hour_utc) DO UPDATE SET + units = excluded.units, + payload = excluded.payload, + content_hash = excluded.content_hash, + points = excluded.points, + first_ts = excluded.first_ts, + last_ts = excluded.last_ts, + updated_at = excluded.updated_at` + + _, err = tx.ExecContext(ctx, q, + b.Metric, b.Layer, FormatTime(b.HourUTC), b.Units, payload, b.Hash, + len(b.Points), FormatTime(b.FirstTS), FormatTime(b.LastTS), + b.Delivery, boolToInt(b.Sealed), FormatTime(now), FormatTime(now)) + if err != nil { + return fmt.Errorf("upsert bucket: %w", err) + } + return nil +} + +func boolToInt(b bool) int { + if b { + return 1 + } + return 0 +} + +// storedPoint — форма точки внутри payload. Содержимое лежит сырым сообщением, +// координаты — рядом: разжимать и перечитывать метки из содержимого пришлось +// бы на каждом слиянии, а формат метки принадлежит HAE, не хранилищу. +type storedPoint struct { + Start string `json:"s"` + End string `json:"e"` + Offset int `json:"o"` + Raw json.RawMessage `json:"p"` +} + +func encodePayload(points []Point) ([]byte, error) { + sp := make([]storedPoint, 0, len(points)) + for _, p := range points { + sp = append(sp, storedPoint{ + Start: FormatTime(p.Start), + End: FormatTime(p.End), + Offset: p.OffsetSeconds, + Raw: p.Raw, + }) + } + + raw, err := json.Marshal(sp) + if err != nil { + return nil, fmt.Errorf("сериализация точек: %w", err) + } + + var buf bytes.Buffer + gz := gzip.NewWriter(&buf) + if _, err := gz.Write(raw); err != nil { + return nil, fmt.Errorf("сжатие точек: %w", err) + } + if err := gz.Close(); err != nil { + return nil, fmt.Errorf("закрытие gzip: %w", err) + } + return buf.Bytes(), nil +} + +func decodePayload(payload []byte) ([]Point, error) { + gz, err := gzip.NewReader(bytes.NewReader(payload)) + if err != nil { + return nil, fmt.Errorf("распаковка точек: %w", err) + } + defer func() { _ = gz.Close() }() + + raw, err := io.ReadAll(gz) + if err != nil { + return nil, fmt.Errorf("чтение точек: %w", err) + } + + var sp []storedPoint + if err := json.Unmarshal(raw, &sp); err != nil { + return nil, fmt.Errorf("разбор точек: %w", err) + } + + out := make([]Point, 0, len(sp)) + for _, p := range sp { + start, err := ParseTime(p.Start) + if err != nil { + return nil, err + } + end, err := ParseTime(p.End) + if err != nil { + return nil, err + } + out = append(out, Point{ + Start: start, + End: end, + OffsetSeconds: p.Offset, + Raw: p.Raw, + }) + } + return out, nil +} + +// MarkSealed помечает час запечатанным — тем, в который досчёта уже не +// ожидается. +// +// Правило, по которому час переводится в это состояние, пока не определено: +// порог глубины досчёта ставится по наблюдениям, которых нет (наблюдалось до +// 22 минут). Метод существует раньше правила намеренно — без него признак +// `sealed` и сигнал о его нарушении нечем проверить, а неопровержимый инвариант +// не отличается от отсутствующего. +func (s *Store) MarkSealed(ctx context.Context, metric, layer string, hour time.Time, sealed bool) error { + const q = `UPDATE bucket SET sealed = ? WHERE metric = ? AND layer = ? AND hour_utc = ?` + + res, err := s.db.ExecContext(ctx, q, boolToInt(sealed), metric, layer, + FormatTime(hour.UTC().Truncate(time.Hour))) + if err != nil { + return fmt.Errorf("update sealed: %w", err) + } + n, err := res.RowsAffected() + if err != nil { + return fmt.Errorf("update sealed: %w", err) + } + if n == 0 { + return ErrNotFound + } + return nil +} + +// CountBuckets возвращает число часовых объектов. Нужно проверке сходимости и +// healthcheck. +func (s *Store) CountBuckets(ctx context.Context) (int64, error) { + var n int64 + if err := s.db.GetContext(ctx, &n, `SELECT count(*) FROM bucket`); err != nil { + return 0, fmt.Errorf("count buckets: %w", err) + } + return n, nil +} + +// Bucket читает объект по координатам. Нужен тестам и будущему Read API. +func (s *Store) Bucket(ctx context.Context, metric, layer string, hour time.Time) (Bucket, error) { + tx, err := s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true}) + if err != nil { + return Bucket{}, fmt.Errorf("begin read tx: %w", err) + } + defer func() { _ = tx.Rollback() }() + + b, found, err := readBucket(ctx, tx, bucketKey{metric, layer, hour.UTC().Truncate(time.Hour)}) + if err != nil { + return Bucket{}, err + } + if !found { + return Bucket{}, ErrNotFound + } + return b, nil +} diff --git a/internal/store/bucket_test.go b/internal/store/bucket_test.go new file mode 100644 index 0000000..89dd7ed --- /dev/null +++ b/internal/store/bucket_test.go @@ -0,0 +1,531 @@ +package store_test + +import ( + "context" + "encoding/json" + "errors" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "git.vakhrushev.me/av/healthlog/internal/store" +) + +func open(t *testing.T) *store.Store { + t.Helper() + + st, err := store.Open(filepath.Join(t.TempDir(), "healthlog.db")) + if err != nil { + t.Fatalf("открытие базы: %v", err) + } + t.Cleanup(func() { _ = st.Close() }) + return st +} + +func ts(t *testing.T, s string) time.Time { + t.Helper() + + v, err := time.Parse(time.RFC3339, s) + if err != nil { + t.Fatalf("метка %q: %v", s, err) + } + return v.UTC() +} + +func point(t *testing.T, metric, layer, start, end, raw string) store.IncomingPoint { + t.Helper() + + return store.IncomingPoint{ + Metric: metric, + Layer: layer, + Units: "count", + Point: store.Point{ + Start: ts(t, start), + End: ts(t, end), + OffsetSeconds: 3 * 3600, + Raw: json.RawMessage(raw), + }, + } +} + +func TestMergePointsКладётЧасОднимОбъектом(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + in := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`), + point(t, "step_count", "minute", "2025-06-05T10:30:00Z", "2025-06-05T10:30:00Z", `{"qty":2}`), + point(t, "step_count", "minute", "2025-06-05T11:00:00Z", "2025-06-05T11:00:00Z", `{"qty":3}`), + } + + stats, err := st.MergePoints(ctx, in, "delivery-1") + if err != nil { + t.Fatalf("слияние: %v", err) + } + if stats.Buckets != 2 { + t.Errorf("объектов %d, ожидалось 2 (два разных часа)", stats.Buckets) + } + + b, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение объекта: %v", err) + } + if len(b.Points) != 2 { + t.Fatalf("точек в объекте %d, ожидалось 2", len(b.Points)) + } + if b.Units != "count" { + t.Errorf("единицы %q потеряны: внутри точки их нет", b.Units) + } + if b.Delivery != "delivery-1" { + t.Errorf("провенанс %q, ожидался delivery-1", b.Delivery) + } + if !b.FirstTS.Equal(ts(t, "2025-06-05T10:00:00Z")) || !b.LastTS.Equal(ts(t, "2025-06-05T10:30:00Z")) { + t.Errorf("границы содержимого %v..%v", b.FirstTS, b.LastTS) + } +} + +// Дозапись в существующий час: объект перечитывается, точки сливаются, ранее +// сохранённые остаются. Точки из объекта не удаляются никогда. +func TestMergePointsДозаписьНеТеряетСохранённое(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + first := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`), + } + second := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:30:00Z", "2025-06-05T10:30:00Z", `{"qty":2}`), + } + + if _, err := st.MergePoints(ctx, first, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + if _, err := st.MergePoints(ctx, second, "d2"); 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) + } + if len(b.Points) != 2 { + t.Fatalf("точек %d, ожидалось 2: дозапись затёрла сохранённое", len(b.Points)) + } + if b.Delivery != "d1" { + t.Errorf("провенанс %q: объект создала первая доставка", b.Delivery) + } +} + +// Хеш — детектор изменений: повтор того же часа не пишет в базу. Именно это +// делает широкие проходы синхронизации дешёвыми. +func TestMergePointsПовторНеПишет(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + in := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1,"source":"Device A"}`), + point(t, "step_count", "minute", "2025-06-05T10:01:00Z", "2025-06-05T10:01:00Z", `{"qty":2,"source":"Device A"}`), + } + + if _, err := st.MergePoints(ctx, in, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + before, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + + stats, err := st.MergePoints(ctx, in, "d2") + if err != nil { + t.Fatalf("повторное слияние: %v", err) + } + if stats.Unchanged != 1 { + t.Errorf("объектов без изменений %d, ожидался 1", stats.Unchanged) + } + + after, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение после повтора: %v", err) + } + if before.Hash != after.Hash { + t.Error("хеш изменился при повторе того же содержимого") + } + if len(after.Points) != len(before.Points) { + t.Errorf("точек стало %d вместо %d", len(after.Points), len(before.Points)) + } +} + +// Порядок точек внутри доставки нестабилен (находка 2), и идентичность не +// имеет права от него зависеть. +func TestMergePointsИдемпотентенКПорядкуТочек(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + in := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`), + point(t, "step_count", "minute", "2025-06-05T10:01:00Z", "2025-06-05T10:01:00Z", `{"qty":2}`), + point(t, "step_count", "minute", "2025-06-05T10:02:00Z", "2025-06-05T10:02:00Z", `{"qty":3}`), + } + reversed := []store.IncomingPoint{in[2], in[1], in[0]} + + if _, err := st.MergePoints(ctx, in, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + first, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + + stats, err := st.MergePoints(ctx, reversed, "d2") + if err != nil { + t.Fatalf("слияние в обратном порядке: %v", err) + } + if stats.Unchanged != 1 { + t.Error("перестановка точек засчиталась изменением") + } + + second, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if first.Hash != second.Hash { + t.Error("хеш зависит от порядка точек в доставке") + } +} + +// Ключ по метке схлопнул бы эти три записи в одну. Живьём такое встречается в +// 22 доставках из 94 (находка 47), причём внутри одной доставки — там тай-брейк +// по времени приёма неприменим в принципе. +func TestMergePointsЗаписиСОднойМеткойНеСхлопываются(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + in := []store.IncomingPoint{ + point(t, "sleep_analysis", "minute", "2025-06-05T19:04:00Z", "2025-06-05T19:16:00Z", `{"qty":0.2,"value":"Во сне"}`), + point(t, "sleep_analysis", "minute", "2025-06-05T19:04:00Z", "2025-06-05T23:21:00Z", `{"qty":4.28,"value":"В кровати"}`), + point(t, "sleep_analysis", "minute", "2025-06-05T19:04:00Z", "2025-06-06T04:51:00Z", `{"qty":9.78,"value":"В кровати"}`), + } + + if _, err := st.MergePoints(ctx, in, "d1"); err != nil { + t.Fatalf("слияние: %v", err) + } + + b, err := st.Bucket(ctx, "sleep_analysis", "minute", ts(t, "2025-06-05T19:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(b.Points) != 3 { + t.Fatalf("точек %d, ожидалось 3: ключ схлопнул записи с одной меткой", len(b.Points)) + } + + // Повтор той же тройки не задваивает: координата включает интервал, и он + // совпадает. + stats, err := st.MergePoints(ctx, in, "d2") + if err != nil { + t.Fatalf("повтор: %v", err) + } + if stats.Unchanged != 1 { + t.Errorf("повтор записей засчитан изменением: %+v", stats) + } +} + +// Эпизод, пересекающий границу часа, ложится в объект по НАЧАЛУ: любой другой +// выбор сделал бы принадлежность объекту зависящей от длительности. +func TestMergePointsЧасПоНачалуИнтервала(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + in := []store.IncomingPoint{ + point(t, "sleep_analysis", "minute", "2025-06-05T19:04:00Z", "2025-06-06T04:51:00Z", `{"qty":9.78}`), + } + if _, err := st.MergePoints(ctx, in, "d1"); err != nil { + t.Fatalf("слияние: %v", err) + } + + if _, err := st.Bucket(ctx, "sleep_analysis", "minute", ts(t, "2025-06-05T19:00:00Z")); err != nil { + t.Fatalf("объект часа начала не найден: %v", err) + } + if _, err := st.Bucket(ctx, "sleep_analysis", "minute", ts(t, "2025-06-06T04:00:00Z")); !errors.Is(err, store.ErrNotFound) { + t.Error("эпизод попал ещё и в объект часа окончания") + } +} + +// Правило разрешения столкновений: выигрывает более полная точка, а не +// последняя пришедшая. Иначе бедная доставка стирает у богатой поля, которых +// сама не несёт. +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,"Min":55,"Max":70,"context":"Отдых"}`) + poor := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", + `{"Avg":60,"Min":55,"Max":70}`) + + 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) + } + + b, err := st.Bucket(ctx, "heart_rate", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(b.Points) != 1 { + t.Fatalf("точек %d, ожидалась 1", len(b.Points)) + } + if !strings.Contains(string(b.Points[0].Raw), "context") { + t.Errorf("бедная точка стёрла поле context: %s", b.Points[0].Raw) + } + if stats.Overwrites != 1 { + t.Errorf("столкновений %d, ожидалось 1: различие обязано оставить след", stats.Overwrites) + } +} + +// source нестабилен: то же измерение приезжает то с одним именем устройства, +// то с другим. Он не входит в ключ и не считается полнотой. +func TestMergePointsСменаИсточникаНеСоздаётВторуюТочку(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + a := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", + `{"qty":1,"source":"Apple Watch Ultra 3|iPad (Anton)"}`) + b := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", + `{"qty":1,"source":"Apple Watch Ultra 3"}`) + + if _, err := st.MergePoints(ctx, []store.IncomingPoint{a}, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + if _, err := st.MergePoints(ctx, []store.IncomingPoint{b}, "d2"); err != nil { + t.Fatalf("второе слияние: %v", err) + } + + got, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(got.Points) != 1 { + t.Fatalf("точек %d, ожидалась 1: смена source задвоила точку", len(got.Points)) + } +} + +// Исход столкновения точек равной полноты обязан зависеть только от значений: +// свёртка по журналу должна давать то же состояние, что приём в реальном +// времени, а внутри одной доставки время приёма общее. +func TestMergePointsРавнаяПолнотаРазрешаетсяДетерминированно(t *testing.T) { + t.Parallel() + + ctx := context.Background() + + a := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`) + b := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":2}`) + + winner := func(order []store.IncomingPoint) string { + st := open(t) + for _, p := range order { + if _, err := st.MergePoints(ctx, []store.IncomingPoint{p}, "d"); err != nil { + t.Fatalf("слияние: %v", err) + } + } + got, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + return string(got.Points[0].Raw) + } + + forward := winner([]store.IncomingPoint{a, b}) + backward := winner([]store.IncomingPoint{b, a}) + if forward != backward { + t.Errorf("исход зависит от порядка доставок: %s против %s", forward, backward) + } +} + +// Содержимое точки хранится исходными байтами: пересборка повторной +// сериализацией теряет литерал, и потеря не видна тестам, сравнивающим +// разобранное с разобранным. +func TestMergePointsХранитТочкуДословно(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + raw := "{\"exact\":1.0,\"huge\":9007199254740993,\"tail\":0.095231829624713534,\"broken\":\"a\\ufffdb\"}" + in := []store.IncomingPoint{ + point(t, "unknown", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", raw), + } + if _, err := st.MergePoints(ctx, in, "d1"); err != nil { + t.Fatalf("слияние: %v", err) + } + + b, err := st.Bucket(ctx, "unknown", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if got := string(b.Points[0].Raw); got != raw { + t.Errorf("точка изменилась при хранении:\n было %s\n стало %s", raw, got) + } +} + +// Тело доставки доходит до 42 МиБ, приходят они непрерывно и внахлёст. +// Конкурентное слияние того же часа не имеет права терять точки: между +// чтением и записью может вклиниться другая доставка. +func TestMergePointsКонкурентноеСлияниеНеТеряетТочки(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + + const writers = 8 + const perWriter = 25 + + var wg sync.WaitGroup + errs := make(chan error, writers) + + for w := range writers { + wg.Go(func() { + for i := range perWriter { + minute := w*perWriter + i + at := ts(t, "2025-06-05T10:00:00Z").Add(time.Duration(minute) * time.Second) + p := store.IncomingPoint{ + Metric: "step_count", + Layer: "raw", + Units: "count", + Point: store.Point{ + Start: at, + End: at, + Raw: json.RawMessage(`{"qty":` + itoa(minute) + `}`), + }, + } + if _, err := st.MergePoints(ctx, []store.IncomingPoint{p}, "d"); err != nil { + errs <- err + return + } + } + }) + } + + wg.Wait() + close(errs) + for err := range errs { + t.Fatalf("конкурентное слияние: %v", err) + } + + b, err := st.Bucket(ctx, "step_count", "raw", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(b.Points) != writers*perWriter { + t.Errorf("точек %d, ожидалось %d: конкурентная запись потеряла точки", + len(b.Points), writers*perWriter) + } +} + +func itoa(n int) string { + if n == 0 { + return "0" + } + var b []byte + for n > 0 { + b = append([]byte{byte('0' + n%10)}, b...) + n /= 10 + } + return string(b) +} + +// Изменение запечатанного часа — сигнал, а не отказ: данные пишутся всё равно, +// но факт обязан дойти до владельца сервиса. Без счётчика допущение «глубже +// такого-то порога досчёта не бывает» не получило бы ни одного наблюдения. +func TestMergePointsИзменениеЗапечатанногоЧаса(t *testing.T) { + t.Parallel() + + st := open(t) + ctx := context.Background() + hour := ts(t, "2025-06-05T10:00:00Z") + + first := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`), + } + if _, err := st.MergePoints(ctx, first, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + if err := st.MarkSealed(ctx, "step_count", "minute", hour, true); err != nil { + t.Fatalf("пометка sealed: %v", err) + } + + late := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:30:00Z", "2025-06-05T10:30:00Z", `{"qty":2}`), + } + stats, err := st.MergePoints(ctx, late, "d2") + if err != nil { + t.Fatalf("досчёт запечатанного часа: %v", err) + } + if stats.SealedHits != 1 { + t.Errorf("изменений запечатанного часа %d, ожидалось 1", stats.SealedHits) + } + + b, err := st.Bucket(ctx, "step_count", "minute", hour) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(b.Points) != 2 { + t.Errorf("точек %d, ожидалось 2: данные обязаны писаться несмотря на sealed", len(b.Points)) + } + if !b.Sealed { + t.Error("признак sealed сброшен дозаписью") + } +} + +// Отмена посреди слияния не имеет права оставить половину: объект либо +// прежний, либо полный. +func TestMergePointsОтменаНеОставляетПоловины(t *testing.T) { + t.Parallel() + + st := open(t) + + base := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`), + } + if _, err := st.MergePoints(context.Background(), base, "d1"); err != nil { + t.Fatalf("первое слияние: %v", err) + } + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + more := []store.IncomingPoint{ + point(t, "step_count", "minute", "2025-06-05T10:30:00Z", "2025-06-05T10:30:00Z", `{"qty":2}`), + } + if _, err := st.MergePoints(ctx, more, "d2"); err == nil { + t.Fatal("слияние на отменённом контексте прошло успешно") + } + + b, err := st.Bucket(context.Background(), "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) + if err != nil { + t.Fatalf("чтение: %v", err) + } + if len(b.Points) != 1 { + t.Errorf("точек %d, ожидалась 1: отмена оставила половинчатое состояние", len(b.Points)) + } +} diff --git a/internal/store/store.go b/internal/store/store.go index 09f34a9..64f8cb5 100644 --- a/internal/store/store.go +++ b/internal/store/store.go @@ -3,8 +3,10 @@ package store import ( + "context" "embed" "fmt" + "io/fs" "net/url" "time" @@ -46,21 +48,39 @@ func (s *Store) Close() error { // dsn собирает строку подключения: WAL для параллельного чтения во время // записи, busy_timeout — чтобы конкурентная запись ждала, а не падала. +// +// `_txlock=immediate` обязателен, а не «на всякий случай»: запись объекта — +// это read-modify-write, и при отложенной блокировке повышение с чтения на +// запись после уже прочитанного снимка даёт SQLITE_BUSY_SNAPSHOT, которого +// busy_timeout не покрывает. Измерено на восьми писателях: 242 успешных +// слияния из 800 против 800 из 800. func dsn(path string) string { q := url.Values{} q.Add("_pragma", "journal_mode(WAL)") q.Add("_pragma", "busy_timeout(5000)") q.Add("_pragma", "foreign_keys(on)") + q.Add("_txlock", "immediate") return "file:" + path + "?" + q.Encode() } +// migrate накатывает миграции. +// +// Через Provider, а не через пакетные функции: goose.SetBaseFS и +// goose.SetDialect пишут глобальное состояние пакета, и два одновременных +// Open дают гонку — её ловит детектор. Сегодня открытие одно, на старте, но +// пересборка витрины откроет второе, и хранилище не должно зависеть от того, +// что вызывающий этого не сделает. func migrate(db *sqlx.DB) error { - goose.SetBaseFS(migrationsFS) - goose.SetLogger(goose.NopLogger()) - if err := goose.SetDialect("sqlite3"); err != nil { - return fmt.Errorf("goose dialect: %w", err) + sub, err := fs.Sub(migrationsFS, "migrations") + if err != nil { + return fmt.Errorf("goose migrations fs: %w", err) } - if err := goose.Up(db.DB, "migrations"); err != nil { + + p, err := goose.NewProvider(goose.DialectSQLite3, db.DB, sub) + if err != nil { + return fmt.Errorf("goose provider: %w", err) + } + if _, err := p.Up(context.Background()); err != nil { return fmt.Errorf("goose up: %w", err) } return nil diff --git a/internal/store/tx.go b/internal/store/tx.go new file mode 100644 index 0000000..fcb833b --- /dev/null +++ b/internal/store/tx.go @@ -0,0 +1,90 @@ +package store + +import ( + "context" + "database/sql" + "errors" + "fmt" + "time" + + sqlite "modernc.org/sqlite" + sqlite3 "modernc.org/sqlite/lib" +) + +// Сколько раз повторять транзакцию, отменённую из-за занятости. +// +// Повтор оборачивает ВСЮ тройку чтение-слияние-запись, а не отдельный запрос: +// после отката прочитанное состояние недействительно, и дописывать в него +// пришедшие точки значит терять чужие. +const ( + txRetries = 5 + txRetryDelay = 20 * time.Millisecond +) + +// inTx выполняет работу в транзакции, повторяя её при занятости базы. +// +// Транзакция открывается сразу на запись (`_txlock=immediate` в DSN). Без +// этого read-modify-write ломается измеримо: при восьми писателях и ста +// слияниях каждый deferred-транзакция дала 242 успеха из 800, immediate — 800 +// из 800. Причина в том, что повышение блокировки с чтения на запись после +// того, как снимок уже прочитан, даёт SQLITE_BUSY_SNAPSHOT (код 517), а его +// `busy_timeout` не покрывает: ждать бесполезно, снимок уже устарел. +func (s *Store) inTx(ctx context.Context, fn func(*sql.Tx) error) error { + var lastErr error + + for attempt := range txRetries { + if attempt > 0 { + select { + case <-ctx.Done(): + return fmt.Errorf("ожидание повтора транзакции: %w", ctx.Err()) + case <-time.After(txRetryDelay * time.Duration(attempt)): + } + } + + err := runTx(ctx, s.db.DB, fn) + if err == nil { + return nil + } + if !isBusy(err) { + return err + } + lastErr = err + } + + return fmt.Errorf("транзакция не прошла за %d попыток: %w", txRetries, lastErr) +} + +func runTx(ctx context.Context, db *sql.DB, fn func(*sql.Tx) error) error { + tx, err := db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin tx: %w", err) + } + + if err := fn(tx); err != nil { + _ = tx.Rollback() + return err + } + if err := tx.Commit(); err != nil { + _ = tx.Rollback() + return fmt.Errorf("commit tx: %w", err) + } + return nil +} + +// isBusy распознаёт отказ из-за занятости базы. +// +// Проверка по типу, а не по тексту: сообщения драйвера не контракт. Кодов два +// — SQLITE_BUSY и SQLITE_BUSY_SNAPSHOT; второй возвращается, когда снимок +// транзакции устарел, и на него `busy_timeout` не действует. +func isBusy(err error) bool { + var serr *sqlite.Error + if !errors.As(err, &serr) { + return false + } + switch serr.Code() { + case sqlite3.SQLITE_BUSY, sqlite3.SQLITE_BUSY_SNAPSHOT: + return true + default: + return false + } +} diff --git a/openspec/changes/razbor-metrik-v-obekty/tasks.md b/openspec/changes/razbor-metrik-v-obekty/tasks.md index d2e83a9..cb43810 100644 --- a/openspec/changes/razbor-metrik-v-obekty/tasks.md +++ b/openspec/changes/razbor-metrik-v-obekty/tasks.md @@ -30,17 +30,17 @@ ## 4. Хранение — часовые объекты -- [ ] 4.1 Модель объекта в `internal/store`; содержимое — исходные байты точек, gzip-BLOB, точки упорядочены по времени -- [ ] 4.2 Слияние: координатный ключ, победа более полной точки, при равной полноте — детерминированный исход по порядку канонических форм -- [ ] 4.3 Столкновение с различием канонической формы — `WARN` без значений и счётчик перезаписей -- [ ] 4.4 Ключ точки — `метрика + слой + начало + конец` одной формы для всех точек: начало из `start`, иначе из `date`; конец из `end`, иначе равен началу. Час объекта — по началу. Ветвления по «классу метрики» быть не должно -- [ ] 4.4a Тест на фикстуре 1.3a: три записи с одной меткой и разными интервалами дают три точки, а не одну; повтор той же тройки следующей доставкой не задваивает; точка без `end` кладётся вырожденным интервалом -- [ ] 4.5 Хеш объекта как детектор изменений: совпал — записи нет -- [ ] 4.6 `_txlock=immediate` в DSN; повтор оборачивает **всю тройку** чтение-слияние-запись; путь «хеш совпал» — под `TxOptions{ReadOnly: true}` -- [ ] 4.7 Распознавание занятости — `errors.As` на `*sqlite.Error`, коды 5 и 517, обёрнуто в `store` -- [ ] 4.8 Тест конкурентной записи: N горутин × M слияний в один `hour_utc`, проверка **суммы** точек, под `-race` -- [ ] 4.9 Тест идемпотентности: повторное слияние того же набора не меняет ни содержимое, ни хеш -- [ ] 4.10 Тесты сценариев слияния: «бедная точка не стирает поля богатой», «смена `source` не создаёт вторую точку» +- [x] 4.1 Модель объекта в `internal/store`; содержимое — исходные байты точек, gzip-BLOB, точки упорядочены по времени +- [x] 4.2 Слияние: координатный ключ, победа более полной точки, при равной полноте — детерминированный исход по порядку канонических форм +- [x] 4.3 Столкновение с различием канонической формы — `WARN` без значений и счётчик перезаписей +- [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.7 Распознавание занятости — `errors.As` на `*sqlite.Error`, коды 5 и 517, обёрнуто в `store` +- [x] 4.8 Тест конкурентной записи: N горутин × M слияний в один `hour_utc`, проверка **суммы** точек, под `-race` +- [x] 4.9 Тест идемпотентности: повторное слияние того же набора не меняет ни содержимое, ни хеш +- [x] 4.10 Тесты сценариев слияния: «бедная точка не стирает поля богатой», «смена `source` не создаёт вторую точку» ## 5. Сшивка с приёмом