package store_test import ( "context" "strings" "testing" "time" "git.vakhrushev.me/av/healthlog/internal/store" ) // Столкновением считается расхождение КАНОНИЧЕСКИХ форм, а не байтов. // // Канонизация и заведена потому, что байты нестабильны: из 81 952 повторно // приехавших точек 67 534 различаются лишь порядком ключей, ещё 63% — последним // разрядом double. Побайтовое сравнение давало тысячи ложных перезаписей на // каждом глубоком проходе, и настоящий отказ правила слияния становился // неотличим от нормы — при том что счётчик перезаписей объявлен единственным // наблюдением за этим правилом. func TestMergePointsДребезгНеСчитаетсяСтолкновением(t *testing.T) { t.Parallel() cases := map[string][2]string{ "порядок ключей": { `{"date":"a","qty":1.5,"source":"Device A"}`, `{"source":"Device A","qty":1.5,"date":"a"}`, }, "последний разряд double": { `{"qty":0.09523182962471353}`, `{"qty":0.09523182962471352}`, }, } for name, pair := range cases { t.Run(name, func(t *testing.T) { t.Parallel() st := open(t) ctx := context.Background() first := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", pair[0]) second := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", pair[1]) if _, err := st.MergePoints(ctx, []store.IncomingPoint{first}, "d1"); err != nil { t.Fatalf("первое слияние: %v", err) } stats, err := st.MergePoints(ctx, []store.IncomingPoint{second}, "d2") if err != nil { t.Fatalf("второе слияние: %v", err) } if stats.Overwrites != 0 { t.Errorf("перезаписей %d, ожидалось 0: содержимое то же с точностью до канонической формы", stats.Overwrites) } if len(stats.Collisions) != 0 { t.Errorf("координат столкновений %d, ожидалось 0", len(stats.Collisions)) } }) } } // Настоящее столкновение обязано оставить след с координатами объекта: одно // число `overwrites` не говорит, какая метрика и какой час пострадали, и // расследовать перезапись по нему нечем. func TestMergePointsСтолкновениеОставляетКоординаты(t *testing.T) { t.Parallel() st := open(t) ctx := context.Background() rich := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"Avg":60,"context":"Отдых"}`) poor := point(t, "heart_rate", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"Avg":61}`) if _, err := st.MergePoints(ctx, []store.IncomingPoint{rich}, "d1"); err != nil { t.Fatalf("первое слияние: %v", err) } stats, err := st.MergePoints(ctx, []store.IncomingPoint{poor}, "d2") if err != nil { t.Fatalf("второе слияние: %v", err) } if len(stats.Collisions) != 1 { t.Fatalf("координат столкновений %d, ожидалась 1", len(stats.Collisions)) } c := stats.Collisions[0] if c.Metric != "heart_rate" || c.Layer != "minute" { t.Errorf("координаты столкновения %s/%s, ожидались heart_rate/minute", c.Metric, c.Layer) } if c.HourUTC.Hour() != 10 { t.Errorf("час столкновения %v", c.HourUTC) } } // «Точка хранится дословно» читается буквально: байты как пришли. json.Marshal // экранирует `&`, `<` и `>` внутри содержимого — порчи значений это не даёт, но // обещание перестаёт быть правдой, а сравнение байтов при следующей доставке // той же точки начинает промахиваться навсегда. func TestMergePointsХранитУгловыеСкобкиИАмперсанд(t *testing.T) { t.Parallel() st := open(t) ctx := context.Background() raw := `{"qty":1,"source":"Anton & Co \"x\""}` in := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", raw) if _, err := st.MergePoints(ctx, []store.IncomingPoint{in}, "d1"); err != nil { t.Fatalf("слияние: %v", err) } b, err := st.Bucket(ctx, "step_count", "minute", ts(t, "2025-06-05T10:00:00Z")) if err != nil { t.Fatalf("чтение: %v", err) } got := string(b.Points[0].Raw) if got != raw { t.Errorf("точка изменилась при хранении:\n было %s\n стало %s", raw, got) } for _, escaped := range []string{"\\u0026", "\\u003c", "\\u003e"} { if strings.Contains(got, escaped) { t.Errorf("содержимое точки экранировано как HTML (%s): %s", escaped, got) } } } // Смена единиц не имеет права молча переподписать уже сохранённые точки: // внутри точки единиц нет, и у ранних точек не остаётся ничего, по чему их // единицы восстановимы. Правило «первое непустое побеждает» плюс счётчик. func TestMergePointsСменаЕдиницНеПерезаписываетМолча(t *testing.T) { t.Parallel() st := open(t) ctx := context.Background() km := store.IncomingPoint{ Metric: "walking_running_distance", Layer: "minute", Units: "km", Point: store.Point{ Start: ts(t, "2025-06-05T10:00:00Z"), End: ts(t, "2025-06-05T10:00:00Z"), Raw: []byte(`{"qty":1}`), }, } mi := km mi.Units = "mi" mi.Start = ts(t, "2025-06-05T10:30:00Z") mi.End = mi.Start mi.Raw = []byte(`{"qty":2}`) if _, err := st.MergePoints(ctx, []store.IncomingPoint{km}, "d1"); err != nil { t.Fatalf("первое слияние: %v", err) } stats, err := st.MergePoints(ctx, []store.IncomingPoint{mi}, "d2") if err != nil { t.Fatalf("второе слияние: %v", err) } if stats.UnitsConflicts != 1 { t.Errorf("расхождений единиц %d, ожидалось 1", stats.UnitsConflicts) } b, err := st.Bucket(ctx, "walking_running_distance", "minute", ts(t, "2025-06-05T10:00:00Z")) if err != nil { t.Fatalf("чтение: %v", err) } if b.Units != "km" { t.Errorf("единицы объекта %q: сохранённые точки переподписаны чужими", b.Units) } } // Счётчик считает СОХРАНЁННЫЕ точки, а не присланные: точные повторы внутри // доставки схлопываются, и счётчик присланных систематически завышал бы // содержимое витрины — расхождение «прислали 1000, лежит 700» было бы невидимо // ровно тогда, когда точки начнут теряться по-настоящему. func TestMergePointsСчитаетСохранённые(t *testing.T) { t.Parallel() st := open(t) ctx := context.Background() p := point(t, "step_count", "minute", "2025-06-05T10:00:00Z", "2025-06-05T10:00:00Z", `{"qty":1}`) stats, err := st.MergePoints(ctx, []store.IncomingPoint{p, p, p}, "d1") if err != nil { t.Fatalf("слияние: %v", err) } if stats.Stored != 1 { t.Errorf("сохранено %d точек, ожидалась 1 (три точных повтора)", stats.Stored) } stats, err = st.MergePoints(ctx, []store.IncomingPoint{p}, "d2") if err != nil { t.Fatalf("повтор: %v", err) } if stats.Stored != 0 { t.Errorf("повтор добавил %d точек, ожидалось 0", stats.Stored) } } // Доставка сворачивается одной транзакцией, поэтому отказ посреди неё не // оставляет частичного состояния — а значит и не зависит от порядка обхода. // Раньше транзакция была на объект, и восемь прогонов одной доставки давали // семь разных наборов записанных объектов. func TestMergePointsОтказНеОставляетЧастиОбъектов(t *testing.T) { t.Parallel() st := open(t) in := make([]store.IncomingPoint, 0, 40) for i := range 40 { at := ts(t, "2025-06-05T00:00:00Z").Add(time.Duration(i) * time.Hour) in = append(in, store.IncomingPoint{ Metric: "step_count", Layer: "minute", Units: "count", Point: store.Point{Start: at, End: at, Raw: []byte(`{"qty":1}`)}, }) } ctx, cancel := context.WithCancel(context.Background()) cancel() if _, err := st.MergePoints(ctx, in, "d1"); err == nil { t.Fatal("слияние на отменённом контексте прошло успешно") } n, err := st.CountBuckets(context.Background()) if err != nil { t.Fatalf("счёт объектов: %v", err) } if n != 0 { t.Errorf("записано %d объектов из 40: отказ оставил частичное состояние", n) } }