Files
healthlog/internal/store/bucket_test.go
T
av f8200f7f80 feat: разбор и хранение тренировок и состояния разума
- секции `workouts` и `stateOfMind` покрыты разбором: тренировка лежит одной
  строкой вместе с маршрутом и внутренними рядами, запись — по ключу `род + id`;
  миграция 00007 заводит обе таблицы и возвращает в очередь `partial`-доставки
  с этими ключами
- сущность заменяется целиком, но условно: приехавшая побеждает, если не теряет
  содержания сохранённой (множество ключей и длины верхнеуровневых массивов), а
  при равном содержании выигрывает версия из более поздней доставки ЖУРНАЛА —
  «побеждает приехавшая» было бы функцией порядка свёртки, и живая витрина
  расходилась бы с пересборкой молча
- отпечаток витрины покрывает тренировки и записи и снимается одним снимком
  базы; отчёт `reindex` считает «было и стало» по каждой единице хранения
2026-08-02 13:05:16 +03:00

860 lines
34 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package store_test
import (
"context"
"encoding/json"
"errors"
"fmt"
"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 TestMergeКладётЧасОднимОбъектом(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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 TestMergeДозаписьНеТеряетСохранённое(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.Merge(ctx, store.Incoming{Points: first}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
if _, err := st.Merge(ctx, store.Incoming{Points: second}, store.DeliveryRef{ID: "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 TestMergeПовторНеПишет(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "d2"})
if err != nil {
t.Fatalf("повторное слияние: %v", err)
}
if stats.Unchanged != 1 {
t.Errorf("объектов без изменений %d, ожидался 1", stats.Unchanged)
}
// Повтор — не столкновение: точка сравнивается сама с собой, и ни один
// счётчик правила слияния расти не должен.
if stats.Overwrites != 0 || stats.Incomparable != 0 {
t.Errorf("повтор посчитан столкновением: перезаписей %d, несравнимых %d",
stats.Overwrites, stats.Incomparable)
}
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 TestMergeИдемпотентенКПорядкуТочек(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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.Merge(ctx, store.Incoming{Points: reversed}, store.DeliveryRef{ID: "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 TestMergeЗаписиСОднойМеткойНеСхлопываются(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "d2"})
if err != nil {
t.Fatalf("повтор: %v", err)
}
if stats.Unchanged != 1 {
t.Errorf("повтор записей засчитан изменением: %+v", stats)
}
}
// Эпизод, пересекающий границу часа, ложится в объект по НАЧАЛУ: любой другой
// выбор сделал бы принадлежность объекту зависящей от длительности.
func TestMergeЧасПоНачалуИнтервала(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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 TestMergeБеднаяТочкаНеСтираетБогатую(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.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{rich}}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
stats, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{poor}}, store.DeliveryRef{ID: "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)
}
}
// Полнота — множество ключей, а не их число. Счётчик значащих полей давал
// сохранённой точке 5 против 2 и стирал настоящее измерение безвозвратно:
// восстановить его можно было бы только из сырого архива, пока он жив.
func TestMergeПоляБезСодержанияНеСтираютИзмерение(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
hollow := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"qty":0,"a":0,"b":0,"c":{},"d":[]}`)
measured := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"date":"2026-07-31 12:00:00 +0300","qty":123.4}`)
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{hollow}}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
stats, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{measured}}, store.DeliveryRef{ID: "d2"})
if 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 got := string(b.Points[0].Raw); got != string(measured.Raw) {
t.Errorf("измерение стёрто точкой без содержания: %s", got)
}
if stats.Incomparable != 0 {
t.Errorf("несравнимых %d: наборы сравнимы, содержательных ключей у первой нет",
stats.Incomparable)
}
}
// Поле с нулевым значением содержания не несёт, но и теряться не должно: при
// равном множестве содержательных ключей выигрывает точка со всеми ключами.
func TestMergeРавноеСодержаниеНеТеряетПоля(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
wide := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"date":"d","qty":10,"Min":0,"Max":0}`)
narrow := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"date":"d","qty":12}`)
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{wide}}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{narrow}}, store.DeliveryRef{ID: "d2"}); 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 got := string(b.Points[0].Raw); got != string(wide.Raw) {
t.Errorf("ключи Min и Max потеряны: %s", got)
}
}
// Несравнимые наборы полей на живом потоке не встретились ни разу (0 из 2 897
// столкновений), поэтому объединение полей не реализовано. Взамен — наблюдение:
// счётчик и координаты объекта, по которым событие можно будет разобрать.
func TestMergeНесравнимыеНаборыСчитаются(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
a := point(t, "blood_glucose", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"qty":5.1}`)
b := point(t, "blood_glucose", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z",
`{"mealTime":"До еды"}`)
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{a}}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
stats, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{b}}, store.DeliveryRef{ID: "d2"})
if err != nil {
t.Fatalf("второе слияние: %v", err)
}
if stats.Incomparable != 1 {
t.Fatalf("несравнимых %d, ожидался 1", stats.Incomparable)
}
// Несравнимость — частный случай столкновения: иначе сумма перезаписей за
// период перестала бы быть сравнимой с прежней.
if stats.Overwrites != 1 {
t.Errorf("перезаписей %d, ожидалась 1", stats.Overwrites)
}
if len(stats.IncomparableAt) != 1 {
t.Fatalf("координат %d, ожидалась 1", len(stats.IncomparableAt))
}
if got := stats.IncomparableAt[0]; got.Metric != "blood_glucose" || got.Layer != "minute" {
t.Errorf("координаты объекта не те: %s/%s", got.Metric, got.Layer)
}
got, err := st.Bucket(ctx, "blood_glucose", "minute", ts(t, "2025-06-05T10:00:00Z"))
if err != nil {
t.Fatalf("чтение: %v", err)
}
if len(got.Points) != 1 {
t.Errorf("точек %d, ожидалась 1: правило обязано выбрать одну", len(got.Points))
}
}
// Список координат упирается в потолок, счётчик — нет: обрезанный список
// остаётся зацепкой для разбора, а масштаб события считает счётчик.
func TestMergeСчётчикРастётПослеПотолкаКоординат(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
const hours = 8
var first, second []store.IncomingPoint
for h := range hours {
at := fmt.Sprintf("2025-06-05T%02d:00:00Z", h)
first = append(first, point(t, "blood_glucose", "minute", at, at, `{"qty":5.1}`))
second = append(second, point(t, "blood_glucose", "minute", at, at, `{"mealTime":"До еды"}`))
}
if _, err := st.Merge(ctx, store.Incoming{Points: first}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
stats, err := st.Merge(ctx, store.Incoming{Points: second}, store.DeliveryRef{ID: "d2"})
if err != nil {
t.Fatalf("второе слияние: %v", err)
}
if stats.Incomparable != hours {
t.Errorf("несравнимых %d, ожидалось %d: счётчик остановился вместе со списком",
stats.Incomparable, hours)
}
if len(stats.IncomparableAt) >= hours {
t.Errorf("координат %d — список не обрезан потолком", len(stats.IncomparableAt))
}
}
// source нестабилен: то же измерение приезжает то с одним именем устройства,
// то с другим. Он не входит в ключ и не считается полнотой.
func TestMergeСменаИсточникаНеСоздаётВторуюТочку(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.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{a}}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{b}}, store.DeliveryRef{ID: "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 TestMergeРавнаяПолнотаРазрешаетсяДетерминированно(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.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{p}}, store.DeliveryRef{ID: "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 TestMergeХранитТочкуДословно(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.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "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 TestMergeКонкурентноеСлияниеНеТеряетТочки(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.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{p}}, store.DeliveryRef{ID: "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 TestMergeИзменениеЗапечатанногоЧаса(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.Merge(ctx, store.Incoming{Points: first}, store.DeliveryRef{ID: "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.Merge(ctx, store.Incoming{Points: late}, store.DeliveryRef{ID: "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 TestMergeОтменаНеОставляетПоловины(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.Merge(context.Background(), store.Incoming{Points: base}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("первое слияние: %v", err)
}
ctx, cancel := context.WithCancel(context.Background())
cancel()
// Сущности идут в ТОЙ ЖЕ транзакции и пишутся ПОСЛЕ объектов, то есть на
// той половине, где отмена вероятнее. Вынесение их во вторую транзакцию —
// напрашивающаяся правка при жалобе на длину транзакции с маршрутом, и она
// прошла бы зелёной, сломав «состояние пересобираемо»: доставка получила бы
// failed при записанной тренировке.
more := store.Incoming{
Points: []store.IncomingPoint{
point(t, "step_count", "minute", "2025-06-05T10:30:00Z", "2025-06-05T10:30:00Z", `{"qty":2}`),
},
Workouts: []store.IncomingEntity{workout(t, "w-отменённая", `{"id":"w-отменённая","qty":1}`)},
}
if _, err := st.Merge(ctx, more, store.DeliveryRef{ID: "d2"}); err == nil {
t.Fatal("слияние на отменённом контексте прошло успешно")
}
if _, err := st.Workout(context.Background(), "w-отменённая"); !errors.Is(err, store.ErrNotFound) {
t.Errorf("тренировка отменённой доставки осталась в витрине: %v", err)
}
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))
}
}
// Отношение победы обязано быть функцией МНОЖЕСТВА точек, а не порядка их
// поступления. Попарная свёртка этого не давала: полнота — частичный порядок,
// тай-брейк — тотальный, и вместе они образовывали цикл (A ⊃ B по ключам,
// B бьёт C тай-брейком, C бьёт A тай-брейком). Из-за цикла одна и та же
// доставка, свёрнутая дважды, давала два состояния витрины поочерёдно —
// хранилище переставало быть свёрткой журнала.
//
// Тройка ниже — ровно такая: её нашёл враждебный проход ревью.
const (
cyclеB = `{"a":1,"b":1}`
cyclеC = `{"a":2,"d":1}`
cyclеA = `{"a":3,"b":1,"c":1}`
)
func cycleTriple(t *testing.T) []store.IncomingPoint {
t.Helper()
return []store.IncomingPoint{
point(t, "heart_rate", "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", cyclеB),
point(t, "heart_rate", "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", cyclеC),
point(t, "heart_rate", "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", cyclеA),
}
}
func TestMergeПовторнаяСвёрткаНеМеняетСостояние(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
in := cycleTriple(t)
state := func() (string, string) {
t.Helper()
if _, err := st.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "d1"}); err != nil {
t.Fatalf("слияние: %v", err)
}
b, err := st.Bucket(ctx, "heart_rate", "raw", ts(t, "2025-06-05T10:00:00Z"))
if err != nil {
t.Fatalf("чтение: %v", err)
}
fp, err := st.Fingerprint(ctx)
if err != nil {
t.Fatalf("отпечаток: %v", err)
}
return string(b.Points[0].Raw), fp
}
wantRaw, wantFP := state()
for pass := 2; pass <= 4; pass++ {
gotRaw, gotFP := state()
if gotRaw != wantRaw || gotFP != wantFP {
t.Fatalf("свёртка %d той же доставки изменила витрину:\n было %s / %s\n стало %s / %s",
pass, wantRaw, wantFP[:16], gotRaw, gotFP[:16])
}
}
}
func TestMergeИсходНеЗависитОтПерестановки(t *testing.T) {
t.Parallel()
ctx := context.Background()
in := cycleTriple(t)
// Все шесть перестановок тройки, и вдобавок разбиение на разные доставки:
// в живом приёме точки приходят порознь, в пересборке — вместе.
perms := [][]int{{0, 1, 2}, {0, 2, 1}, {1, 0, 2}, {1, 2, 0}, {2, 0, 1}, {2, 1, 0}}
winner := func(order []int, split bool) string {
st := open(t)
if split {
for _, i := range order {
if _, err := st.Merge(ctx, store.Incoming{Points: []store.IncomingPoint{in[i]}}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
}
} else {
batch := make([]store.IncomingPoint, 0, len(order))
for _, i := range order {
batch = append(batch, in[i])
}
if _, err := st.Merge(ctx, store.Incoming{Points: batch}, store.DeliveryRef{ID: "d"}); err != nil {
t.Fatalf("слияние: %v", err)
}
}
b, err := st.Bucket(ctx, "heart_rate", "raw", ts(t, "2025-06-05T10:00:00Z"))
if err != nil {
t.Fatalf("чтение: %v", err)
}
return string(b.Points[0].Raw)
}
want := winner(perms[0], false)
for _, p := range perms {
for _, split := range []bool{false, true} {
if got := winner(p, split); got != want {
t.Errorf("перестановка %v (порознь=%v) дала %s, ожидалось %s", p, split, got, want)
}
}
}
}
// Точка, ни одно значение которой не несёт измерения, не должна вытеснять
// настоящее измерение. Раньше вытесняла: множества содержательных ключей
// равны, и решал второй разряд — по ключам, а не по содержанию.
func TestMergeПадингНеВытесняетИзмерение(t *testing.T) {
t.Parallel()
ctx := context.Background()
const real = `{"date":"d","qty":72.5}`
cases := []struct {
name string
junk string
}{
{"падинг из null", `{"date":"d","qty":0.001,"p1":null,"p2":null,"p3":null}`},
{"падинг из false", `{"date":"d","qty":false,"Min":false,"Max":false}`},
}
for _, c := range cases {
t.Run(c.name, func(t *testing.T) {
t.Parallel()
st := open(t)
in := []store.IncomingPoint{
point(t, "heart_rate", "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", real),
point(t, "heart_rate", "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", c.junk),
}
stats, err := st.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "d1"})
if err != nil {
t.Fatalf("слияние: %v", err)
}
// Правило полноты обязано молчать: это столкновение равнополных,
// а не победа надмножества.
if stats.Incomparable != 0 {
t.Errorf("несравнимых %d, ожидалось 0: наборы содержательных ключей равны", stats.Incomparable)
}
})
}
}
// Имя метрики приходит из тела доставки дословно и ничем не ограничено.
// Без обрезки одна доставка порождает WARN-строку в десятки мегабайт и
// вытесняет из ротации логов всю недавнюю историю.
func TestMergeИмяМетрикиВКоординатеОбрезано(t *testing.T) {
t.Parallel()
st := open(t)
ctx := context.Background()
huge := strings.Repeat("м", 5000)
in := []store.IncomingPoint{
point(t, huge, "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`),
point(t, huge, "raw", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":2}`),
}
stats, err := st.Merge(ctx, store.Incoming{Points: in}, store.DeliveryRef{ID: "d1"})
if err != nil {
t.Fatalf("слияние: %v", err)
}
if len(stats.Collisions) == 0 {
t.Fatal("столкновение не зафиксировано")
}
if got := len(stats.Collisions[0].Metric); got > 128 {
t.Errorf("имя метрики в координате %d байт — оно уедет в лог как есть", got)
}
}