Files
healthlog/internal/store/bucket.go
T
av 8331328134 Дозакрыты находки ревью по слиянию сущностей
- Правило покрытия получило второй разряд (условный, как у точек), запрет
  вырождения формы и счёт содержательных элементов ряда: скелет из скаляров и
  ряд из null больше не затирают маршрут. Победитель внутри доставки стал
  функцией множества версий — общим помощником с точками, — а провенанс
  поднимается и при совпавшем хеше, иначе отложенная доставка возвращала витрину
  к прежнему содержимому.
- Одно поле не того типа больше не уносит сущность, а пропуски видны в учётной
  записи доставки (миграция 00008, NULL = «не измерялось»); каноническая форма
  считается один раз и вне транзакции; откат бинаря поверх новой схемы отказывает
  на старте; текст ошибки разбора не несёт значений из тела.
- Ревью кода профилем deep (девять проходов) нашло две регрессии и обе закрыты:
  безусловный второй разряд запирал законный досчёт навсегда, а выбор победителя
  был квадратичен по числу присланных версий одного ключа.
2026-08-02 16:38:18 +03:00

938 lines
42 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
import (
"bytes"
"compress/gzip"
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"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
// Stored — сколько точек ДОБАВИЛОСЬ в объекты. Не то же, что число точек в
// доставке: точные повторы схлопываются, и счётчик присланных завышал бы
// содержимое витрины ровно тогда, когда точки начнут теряться.
Stored int
Unchanged int
Overwrites int
SealedHits int
// UnitsConflicts — объекты, куда точки приехали в единицах, отличных от
// сохранённых. Молчаливая замена недопустима: внутри точки единиц нет, и
// у ранних точек не остаётся ничего, по чему их единицы восстановимы.
UnitsConflicts int
// Incomparable — столкновения, где наборы полей несравнимы: у каждой точки
// есть содержательный ключ, которого нет у другой. Правило полноты тут
// бессильно, и вместо объединения полей — которого на живом потоке не
// потребовалось ни разу — заведено наблюдение. Несравнимость считается
// СВЕРХ Overwrites, а не вместо: иначе сумма перезаписей за период
// перестала бы быть сравнимой с прежней.
Incomparable int
// Collisions — координаты первых столкновений, для записи в лог. Без них
// счётчик перезаписей не говорит, какая метрика и какой час пострадали.
Collisions []Collision
// IncomparableAt — то же для несравнимых наборов.
IncomparableAt []Collision
// Workouts и Records — сколько сущностей пришло с доставкой.
Workouts int
Records int
// WorkoutsWritten и RecordsWritten — сколько из них действительно легло в
// витрину. Разница с пришедшими — работа хеша-детектора: тренировка
// переприсылается каждой доставкой, пока не доедет маршрут.
WorkoutsWritten int
RecordsWritten int
// EntitiesHeld — приехавшие версии, отклонённые как теряющие содержание
// сохранённой (включая несравнимые наборы). Это и есть плата за отказ
// объединять поля: событие считается, а не предотвращается молча.
//
// На него опирается ЕДИНСТВЕННЫЙ контроль того, что правило покрытия не
// стало слишком строгим: сходимость отпечатка этого не проверяет по
// построению — живой приём и пересборка пользуются одним правилом и
// одинаково сойдутся на одинаково удержанной версии. Поэтому счётчик
// обязан считать ровно удержания и ничего сверх.
EntitiesHeld int
// HeldAt — координаты первых таких сущностей, для записи в лог.
HeldAt []EntityRef
// EntitiesDiverging — версии одного ключа, приехавшие в ОДНОМ теле с разным
// содержанием. Событие другого рода: победитель ложится в витрину целиком,
// терять нечего, лечится оно не тем же. Считается отдельно от удержаний,
// иначе одно число отвечало бы на два вопроса — и число удержаний, по
// которому судят о строгости правила, стало бы неотличимо от шума.
EntitiesDiverging int
// DivergingAt — координаты первых таких сущностей.
DivergingAt []EntityRef
}
// Collision — координаты объекта, где столкновение разрешилось перезаписью
// содержимого точки. Значений точек не несёт: данные о здоровье чувствительнее
// токенов.
type Collision struct {
Metric string
Layer string
HourUTC time.Time
}
// maxCollisionsReported — сколько координат столкновений попадает в лог.
// Больше горсти не нужно: они нужны как зацепка для разбора, а не как отчёт.
const maxCollisionsReported = 5
// maxMetricInLog — предел длины имени метрики в координате столкновения.
//
// Имя приходит из тела доставки дословно и ничем не ограничено, а предел
// приёма — 64 МиБ: без обрезки одна доставка порождает WARN-строку в десятки
// мегабайт и вытесняет из ротации логов всю недавнюю историю, включая записи о
// доставках, которые действительно потерялись. Имена метрик HAE — десятки байт.
const maxMetricInLog = 64
func clipMetric(metric string) string {
if len(metric) <= maxMetricInLog {
return metric
}
return metric[:maxMetricInLog] + "…"
}
// Merge раскладывает всё, что дала доставка, по витрине: точки — по часовым
// объектам, сущности — по своим таблицам.
//
// Час берётся по НАЧАЛУ точки: интервал пересекает границы часов, и любой
// другой выбор сделал бы принадлежность объекту зависящей от длительности.
// Вся доставка сворачивается ОДНОЙ транзакцией, а не транзакцией на объект и
// не отдельной транзакцией на сущности.
//
// Транзакция на объект давала недетерминированное частичное состояние: обход
// карты групп рандомизирован, и при отказе посреди доставки набор уже
// закоммиченных объектов каждый раз другой (измерено: восемь прогонов одной
// доставки — семь разных состояний). Это ломает инвариант «состояние
// пересобираемо»: пересборка из архива давала бы не то, что живой приём, а
// хеш-детектор после расхождения переписывал бы «неизменившееся». Наблюдение
// «секции не смешиваются в одной доставке» собрано за двое суток и основанием
// для второй транзакции не является.
//
// Заодно снимается стоимость: отдельный коммит на объект стоил около 0.7 мс,
// то есть 8.5 мс на килобайт тела.
func (s *Store) Merge(ctx context.Context, in Incoming, from DeliveryRef) (MergeStats, error) {
groups := groupByHour(in.Points)
keys := sortedKeys(groups)
// Каноническая форма сущности и её хеш считаются ОДИН раз на версию и здесь
// — до входа в транзакцию. Транзакция открыта `immediate`, то есть блокирует
// запись, и повторяется до пяти раз при занятости базы: канонизация внутри
// неё умножала бы и пик кучи, и время удержания блокировки. Измерено: тело
// 40 МиБ даёт 768 МиБ пика, 63 МиБ удерживают блокировку 5.019 с при
// busy_timeout 5000, после чего конкурентный CreateDelivery исчерпывает
// повторы.
//
// Пределов остаётся два, и оба названы вслух. Первый: разбор СОХРАНЁННОЙ
// версии остаётся внутри транзакции — её содержимое читается оттуда же и
// только когда хеш разошёлся; удержание блокировки пропорционально её
// размеру. Второй: форма и множества ключей всех версий доставки
// УДЕРЖИВАЮТСЯ в памяти до конца транзакции, то есть расход пропорционален
// размеру доставки, а не самой большой её сущности. Про процессор здесь
// стало лучше, про память — хуже, и закрыть оба может лишь предел на размер
// сущности вместе с потоковым расчётом.
workouts, err := prepareEntities(in.Workouts, from)
if err != nil {
return MergeStats{}, err
}
records, err := prepareEntities(in.Records, from)
if err != nil {
return MergeStats{}, err
}
workouts, workoutsDiverging, workoutsDivergingAt, err := dedupeEntities(ctx, workouts)
if err != nil {
return MergeStats{}, err
}
records, recordsDiverging, recordsDivergingAt, err := dedupeEntities(ctx, records)
if err != nil {
return MergeStats{}, err
}
var stats MergeStats
err = s.inTx(ctx, func(tx *sql.Tx) error {
// Счётчики обнуляются на каждой попытке: повтор транзакции начинает
// слияние заново, и накопленное от прошлой попытки посчиталось бы дважды.
stats = MergeStats{
Workouts: len(in.Workouts),
Records: len(in.Records),
EntitiesDiverging: workoutsDiverging + recordsDiverging,
DivergingAt: clipRefs(append(append([]EntityRef{},
workoutsDivergingAt...), recordsDivergingAt...)),
}
for _, key := range keys {
group := groups[key]
res, err := mergeBucket(ctx, tx, key, group, from.ID)
if err != nil {
return err
}
stats.Buckets++
stats.Stored += res.stored
stats.Overwrites += res.overwrites
stats.Incomparable += res.incomparable
if res.unchanged {
stats.Unchanged++
}
if res.sealed {
stats.SealedHits++
}
if res.unitsConflict {
stats.UnitsConflicts++
}
// Список координат упирается в потолок, счётчики — нет: обрезанный
// список остаётся зацепкой для разбора, а масштаб события считает
// счётчик.
coord := Collision{Metric: clipMetric(key.metric), Layer: key.layer, HourUTC: key.hourUTC}
if res.overwrites > 0 && len(stats.Collisions) < maxCollisionsReported {
stats.Collisions = append(stats.Collisions, coord)
}
if res.incomparable > 0 && len(stats.IncomparableAt) < maxCollisionsReported {
stats.IncomparableAt = append(stats.IncomparableAt, coord)
}
}
now := Now()
written, held, heldAt, err := mergeEntities(ctx, tx, workoutTable, workouts, now)
if err != nil {
return err
}
stats.WorkoutsWritten = written
stats.EntitiesHeld += held
stats.HeldAt = clipRefs(append(stats.HeldAt, heldAt...))
written, held, heldAt, err = mergeEntities(ctx, tx, recordTable, records, now)
if err != nil {
return err
}
stats.RecordsWritten = written
stats.EntitiesHeld += held
stats.HeldAt = clipRefs(append(stats.HeldAt, heldAt...))
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 живёт уже
// внутри объекта и их не дублирует.
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
incomparable int
stored int
unchanged bool
sealed bool
unitsConflict bool
}
// mergeBucket выполняет чтение, слияние и запись одного объекта.
//
// Тройка обязана быть атомарной: между чтением и записью может вклиниться
// другая доставка того же часа, и её точки пропали бы молча. Атомарность даёт
// транзакция вызывающего — она общая на всю доставку.
func mergeBucket(ctx context.Context, tx *sql.Tx, key bucketKey, group *pointGroup, deliveryID string) (mergeResult, error) {
var res mergeResult
stored, found, err := readBucket(ctx, tx, key)
if err != nil {
return res, err
}
merged, overwrites, incomparable, err := mergePoints(ctx, stored.Points, group.points)
if err != nil {
return res, err
}
res.overwrites = overwrites
res.incomparable = incomparable
// Считаем сохранённые точки, а не присланные: точные повторы внутри
// доставки схлопываются, и счётчик присланных систематически завышал бы
// содержимое витрины — расхождение «прислали 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
}
// mergePoints сливает сохранённые точки с пришедшими по координатному ключу.
//
// При столкновении выигрывает БОЛЕЕ ПОЛНАЯ точка, а не последняя пришедшая:
// 981 столкновение из 2 897 различается набором полей, и правило «последняя
// победила» стирало бы у сохранённой точки поля, которых новая не несёт.
//
// При равной полноте исход определяется порядком канонических форм, а не
// порядком доставок: у сохранённой точки нет провенанса, а четверть доставок
// несёт столкновения внутри себя, где время приёма общее. Свёртка по журналу
// обязана давать то же состояние, что приём в реальном времени.
//
// Второй счётчик — несравнимые наборы полей. Они и есть та часть правила,
// которую заменили наблюдением: объединять поля никто не будет, пока счётчик
// не заговорит.
//
// Точки из объекта не удаляются никогда.
func mergePoints(ctx context.Context, stored, incoming []Point) (merged []Point, overwrites, incomparable int, err error) {
type coord struct {
start int64
end int64
}
byCoord := make(map[coord][]candidate, len(stored)+len(incoming))
order := make([]coord, 0, len(stored)+len(incoming))
// Кандидаты копятся, а не сворачиваются попарно. Попарная свёртка была
// НЕВЕРНА: полнота — частичный порядок, тай-брейк — тотальный, и их
// смешение даёт нетранзитивное отношение победы. Оно образует цикл
// (A ⊃ B по ключам, B бьёт C тай-брейком, C бьёт A тай-брейком), после
// чего исход зависит от того, что уже лежало в объекте: одна и та же
// доставка, свёрнутая дважды, давала два разных состояния витрины
// поочерёдно. Это ломало главный инвариант — «состояние пересобираемо».
add := func(p Point) {
c := coord{start: p.Start.UnixNano(), end: p.End.UnixNano()}
cands, seen := byCoord[c]
if !seen {
order = append(order, c)
}
cand := newCandidate(p)
// Побайтовое равенство — только быстрый путь. Тем же самым считается
// совпадение КАНОНИЧЕСКИХ форм: канонизация и заведена потому, что
// байты нестабильны. Из 81952 повторно приехавших точек 67534
// различаются лишь порядком ключей, ещё 63% — последним разрядом
// double. Считай мы по байтам, счётчик перезаписей давал бы тысячи
// ложных срабатываний на каждом глубоком проходе, и настоящий отказ
// правила слияния стал бы неотличим от нормы.
for _, q := range cands {
if bytes.Equal(q.key, cand.key) {
return
}
}
byCoord[c] = append(cands, cand)
}
for _, p := range stored {
add(p)
}
for _, p := range incoming {
add(p)
}
out := make([]Point, 0, len(order))
for _, c := range order {
cands := byCoord[c]
winner, unrelated, err := resolve(ctx, cands)
if err != nil {
return nil, 0, 0, err
}
// Перезаписей столько, сколько точек уступило: при двух кандидатах
// одна, при трёх две. Так счёт остаётся сравнимым с прежним, где
// столкновение считалось на каждую приехавшую точку.
overwrites += len(cands) - 1
if unrelated {
incomparable++
}
out = append(out, winner)
}
// Порядок точек в объекте канонический и входит в хеш: его задаёт ТОЛЬКО
// сортировка ниже. Порядок обхода карты не специфицирован, и полагаться на
// него значило бы получать разные хеши для одного содержимого — тогда
// «неизменившийся» объект переписывался бы каждым глубоким проходом.
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, incomparable, nil
}
// candidate — точка вместе с тем, что о ней нужно знать при выборе
// победителя. Разбор и канонизация делаются ОДИН раз на точку: сравнений
// квадратично по числу кандидатов, и пересчёт на каждое сравнение означал бы
// разбор точки столько раз, сколько на координате кандидатов.
type candidate struct {
pt Point
key []byte // каноническая форма: она же ключ дедупликации и порядок
fields canon.Fields
}
func newCandidate(p Point) candidate {
return candidate{pt: p, key: canon.SortKey(p.Raw), fields: canon.Analyze(p.Raw)}
}
// resolve выбирает победителя среди кандидатов одной координаты.
//
// Механизм общий с выбором версии сущности — pickBest: отбрасываем
// превзойдённых по частичному порядку, среди оставшихся берём минимум по
// тотальному. Отношения разные (полнота у точек, покрытие у сущностей), а
// рассуждение одно, и второй его экземпляр однажды уже разошёлся со стандартом
// нетранзитивностью.
//
// Победителем остаётся одна из пришедших точек ДОСЛОВНО: правило выбирает, а
// не конструирует. Каноническая форма существует только в момент сравнения, и
// вернуть её значило бы сохранить округлённое число вместо присланного.
//
// Второй возврат — остались ли непревзойдёнными несколько точек с
// несравнимыми наборами содержательных полей. На живом потоке этого не
// случилось ни разу (0 из 2 897 столкновений), поэтому объединение полей не
// реализовано: вместо него счётчик, который скажет, если событие наступит.
func resolve(ctx context.Context, cands []candidate) (Point, bool, error) {
winner, maximal, err := pickBest(ctx, cands, pointDominates, pointLess)
if err != nil {
return Point{}, false, err
}
// Несравнимость — не «осталось больше одного»: точки с одинаковыми
// наборами полей и разными значениями тоже остаются обе, и это рядовой
// тай-брейк. Считается только то, ради чего отложено объединение полей:
// у каждой из двух есть содержательный ключ, которого нет у другой.
return cands[winner].pt, hasIncomparablePair(cands, maximal), nil
}
// pointDominates — строгое превосходство по полноте. Relate возвращает
// FullnessSuperset только когда a несёт всё, что b, и сверх того, поэтому
// отношение уже строгое.
func pointDominates(a, b candidate) bool {
return a.fields.Relate(b.fields) == canon.FullnessSuperset
}
// pointLess — тотальный порядок по канонической форме. Минимум единствен:
// кандидаты с равной формой схлопываются ещё при сборе множества.
func pointLess(a, b candidate) bool {
return bytes.Compare(a.key, b.key) < 0
}
func hasIncomparablePair(cands []candidate, maximal []int) bool {
for i := range maximal {
for j := i + 1; j < len(maximal); j++ {
if cands[maximal[i]].fields.Relate(cands[maximal[j]].fields) == canon.FullnessIncomparable {
return true
}
}
}
return false
}
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,
})
}
// Через 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)
}
return gzipBytes(raw.Bytes())
}
// gzipBytes и gunzipBytes — единственное кодирование содержимого витрины,
// общее для часового объекта и для сущности с собственным `id`. Второй кадр
// упаковки разошёлся бы с первым при первой же правке — например, забытым
// Close, который оставляет усечённый блоб.
func gzipBytes(raw []byte) ([]byte, error) {
var buf bytes.Buffer
gz := gzip.NewWriter(&buf)
if _, err := gz.Write(raw); err != nil {
return nil, fmt.Errorf("сжатие содержимого: %w", err)
}
// Close дописывает хвост gzip; без него блоб читается лишь частично.
if err := gz.Close(); err != nil {
return nil, fmt.Errorf("закрытие gzip: %w", err)
}
return buf.Bytes(), nil
}
func gunzipBytes(payload []byte) ([]byte, 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)
}
return raw, nil
}
func decodePayload(payload []byte) ([]Point, error) {
raw, err := gunzipBytes(payload)
if err != nil {
return nil, 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
}
// BucketHours возвращает часы, за которые есть объекты метрики в слое, по
// возрастанию. Первое, что понадобится каталогу разрезов, и то, чем проверка
// сходимости находит объект, не зная дат в фикстуре.
func (s *Store) BucketHours(ctx context.Context, metric, layer string) ([]time.Time, error) {
const q = `
SELECT hour_utc FROM bucket
WHERE metric = ? AND layer = ? ORDER BY hour_utc`
var raw []string
if err := s.db.SelectContext(ctx, &raw, q, metric, layer); err != nil {
return nil, fmt.Errorf("select bucket hours: %w", err)
}
out := make([]time.Time, 0, len(raw))
for _, s := range raw {
t, err := ParseTime(s)
if err != nil {
return nil, err
}
out = append(out, t)
}
return out, 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
}
// Fingerprint возвращает отпечаток содержимого витрины: SHA-256 по координатам
// и хешам всех её сущностей в детерминированном порядке.
//
// Нужен проверке сходимости на живом архиве. Число объектов и число точек к
// правилу разрешения столкновений нечувствительны: на координате всегда лежит
// ровно одна точка, и правило выбирает, КАКАЯ это будет точка, а не сколько их.
// Значит «объектов столько же» совпадёт и при заведомо сломанном правиле, а
// отпечаток — нет.
//
// Покрывает ВСЕ единицы хранения — часовые объекты, тренировки и записи.
// Отпечаток одних объектов давал бы «состояние сошлось» при разъехавшихся
// тренировках, то есть ломался бы молча тем самым изменением, которое добавило
// данные.
//
// Все разделы читаются ОДНИМ снимком базы: отпечаток рабочей витрины снимается
// под живым приёмом, и запросы вне общей транзакции дали бы смесь «объекты до»
// и «тренировки после» — ложное расхождение у единственного оракула.
//
// Значений он не раскрывает: содержимое участвует только своим хешем.
func (s *Store) Fingerprint(ctx context.Context) (string, error) {
tx, err := s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true})
if err != nil {
return "", fmt.Errorf("begin read tx: %w", err)
}
defer func() { _ = tx.Rollback() }()
h := sha256.New()
if err := fingerprintBuckets(ctx, tx, h); err != nil {
return "", err
}
if err := fingerprintEntities(ctx, tx, h); err != nil {
return "", err
}
return hex.EncodeToString(h.Sum(nil)), nil
}
// Признак раздела впереди строки: без него строка одного раздела может совпасть
// со строкой другого, и два разных состояния витрины дали бы один отпечаток.
const (
fpBucket = "b"
fpWorkout = "w"
fpRecord = "r"
)
func fingerprintBuckets(ctx context.Context, tx *sql.Tx, h io.Writer) error {
const q = `
SELECT metric, layer, hour_utc, content_hash, points, units, sealed FROM bucket
ORDER BY metric, layer, hour_utc`
rows, err := tx.QueryContext(ctx, q)
if err != nil {
return fmt.Errorf("select buckets: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var metric, layer, hour, hash, units string
var points int
var sealed bool
if err := rows.Scan(&metric, &layer, &hour, &hash, &points, &units, &sealed); err != nil {
return fmt.Errorf("scan bucket: %w", err)
}
// Поля переменной длины идут с длиной впереди: разделитель, который
// может встретиться ВНУТРИ поля, даёт одну строку для разных состояний,
// а имя метрики и единицы приходят из тела доставки дословно. Ошибиться
// здесь значит получить «состояние совпало» при разошедшемся состоянии —
// то есть сломать молча ровно тот оракул, ради которого отпечаток и
// заведён. Тот же приём в canon.HashAll и по той же причине.
fmt.Fprintf(h, "%s|%d:%s|%d:%s|%s|%s|%d|%d:%s|%t\n",
fpBucket, len(metric), metric, len(layer), layer, hour, hash, points,
len(units), units, sealed)
}
if err := rows.Err(); err != nil {
return fmt.Errorf("select buckets: %w", err)
}
return nil
}
func fingerprintEntities(ctx context.Context, tx *sql.Tx, h io.Writer) error {
queries := []struct {
tag string
sql string
}{
// Род и идентификатор идут ОТДЕЛЬНЫМИ полями, каждое со своей длиной, а
// не склейкой `kind || '/' || id`: склейка выполняется до взятия длины,
// и пара (`a`, `b/c`) даёт ту же строку, что (`a/b`, `c`). Отпечаток —
// единственный оракул сходимости, по нему принимается необратимое
// решение о подмене базы; два разных состояния витрины не имеют права
// дать один отпечаток. У тренировки род один и в строку не идёт.
{fpWorkout, `SELECT '', id, start_utc, content_hash FROM workout ORDER BY id`},
{fpRecord, `SELECT kind, id, ts_utc, content_hash FROM record ORDER BY kind, id`},
}
for _, q := range queries {
if err := fingerprintRows(ctx, tx, h, q.tag, q.sql); err != nil {
return err
}
}
return nil
}
func fingerprintRows(ctx context.Context, tx *sql.Tx, h io.Writer, tag, query string) error {
rows, err := tx.QueryContext(ctx, query)
if err != nil {
return fmt.Errorf("select entities: %w", err)
}
defer func() { _ = rows.Close() }()
for rows.Next() {
var kind, id, ts, hash string
if err := rows.Scan(&kind, &id, &ts, &hash); err != nil {
return fmt.Errorf("scan entity: %w", err)
}
fmt.Fprintf(h, "%s|%d:%s|%d:%s|%s|%s\n", tag, len(kind), kind, len(id), id, ts, hash)
}
if err := rows.Err(); err != nil {
return fmt.Errorf("select entities: %w", err)
}
return 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
}