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 mergeCategories(ctx, tx, in.Categories, from) }) 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 } // writeBucket пишет часовой объект. // // `first_delivery_id` в `DO UPDATE` НЕ входит намеренно: провенанс объекта — // «кто создал строку», то есть функция ПОРЯДКА СВЁРТКИ, а не журнала. Это // допустимо ровно потому, что в отпечаток витрины он не идёт // (см. fingerprintBuckets): расхождение по нему ненаблюдаемо и решений по нему // не принимают. // // Единице хранения, чей провенанс входит в отпечаток, такого правила МАЛО — там // нужен явный минимум по журналу, иначе живой приём и пересборка разойдутся при // одинаковом журнале. Образец — mergeCategories в category.go. Сказано здесь, // потому что копировать будут отсюда: этот upsert старше и проще. 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 по координатам // и хешам всех её сущностей в детерминированном порядке. // // Нужен проверке сходимости на живом архиве. Число объектов и число точек к // правилу разрешения столкновений нечувствительны: на координате всегда лежит // ровно одна точка, и правило выбирает, КАКАЯ это будет точка, а не сколько их. // Значит «объектов столько же» совпадёт и при заведомо сломанном правиле, а // отпечаток — нет. // // Покрывает ВСЕ единицы хранения — часовые объекты, тренировки, записи и реестр // категориальных значений. Отпечаток одних объектов давал бы «состояние // сошлось» при разъехавшихся тренировках, то есть ломался бы молча тем самым // изменением, которое добавило данные. // // От реестра берётся НАБЛЮДЕНИЕ, но не выведенный код: см. fingerprintCategories. // // Все разделы читаются ОДНИМ снимком базы: отпечаток рабочей витрины снимается // под живым приёмом, и запросы вне общей транзакции дали бы смесь «объекты до» // и «тренировки после» — ложное расхождение у единственного оракула. // // Значений он не раскрывает: содержимое участвует только своим хешем. 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 } if err := fingerprintCategories(ctx, tx, h); err != nil { return "", err } return hex.EncodeToString(h.Sum(nil)), nil } // Признак раздела впереди строки: без него строка одного раздела может совпасть // со строкой другого, и два разных состояния витрины дали бы один отпечаток. const ( fpBucket = "b" fpWorkout = "w" fpRecord = "r" fpCategory = "c" ) // fingerprintCategories добавляет в отпечаток реестр категориальных значений — // его НАБЛЮДЕНИЕ, но не выведенный код. // // Ключ и провенанс — функция журнала: те же тела в том же порядке дают их // побайтно. Код — функция журнала И версии словаря в бинаре. Включи его сюда, и // отпечаток перестал бы отвечать на свой единственный вопрос («дал ли повтор // журнала то же состояние») ровно тогда, когда его задают: всякое пополнение // словаря давало бы расхождение при побайтно совпавшем журнале, а человек, // принимающий по отпечатку необратимое решение о подмене базы, читал бы это как // дефект. Правильность вывода кода проверяют тесты словаря — это другой вопрос, // и смешение обесценило бы оракул. // // Значений наружу отпечаток не раскрывает: он хеш. func fingerprintCategories(ctx context.Context, tx *sql.Tx, h io.Writer) error { const q = ` SELECT metric, field, value, first_seen_utc, first_delivery_id FROM category_value ORDER BY metric, field, value` rows, err := tx.QueryContext(ctx, q) if err != nil { return fmt.Errorf("select category values: %w", err) } defer func() { _ = rows.Close() }() for rows.Next() { var metric, field, value, firstSeen, deliveryID string if err := rows.Scan(&metric, &field, &value, &firstSeen, &deliveryID); err != nil { return fmt.Errorf("scan category value: %w", err) } // Длина впереди каждого поля переменной длины: значение приходит из тела // дословно и может содержать что угодно, включая признак раздела и // разделители. Без длины пара (`a`, `b|c`) дала бы ту же строку, что // (`a|b`, `c`), — то есть два разных состояния витрины сошлись бы // отпечатком. Тот же приём в fingerprintBuckets и по той же причине. fmt.Fprintf(h, "%s|%d:%s|%d:%s|%d:%s|%s|%d:%s\n", fpCategory, len(metric), metric, len(field), field, len(value), value, firstSeen, len(deliveryID), deliveryID) } if err := rows.Err(); err != nil { return fmt.Errorf("select category values: %w", err) } return nil } 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 }