package store import ( "bytes" "compress/gzip" "context" "database/sql" "encoding/json" "errors" "fmt" "io" "sort" "time" "git.vakhrushev.me/av/healthlog/internal/canon" ) // Point — точка внутри часового объекта. // // Координаты — метрика, слой, начало и конец; у точки-измерения конец равен // началу. Содержимое хранится исходными байтами: пересборка повторной // сериализацией теряет литерал, и потеря не видна тестам на фикстурах. type Point struct { Start time.Time End time.Time OffsetSeconds int Raw json.RawMessage } // Bucket — часовой объект: точки одной метрики одного слоя за один час UTC. type Bucket struct { Metric string Layer string HourUTC time.Time Units string Points []Point Hash string Sealed bool FirstTS time.Time LastTS time.Time Delivery string } // MergeStats — что случилось при слиянии. Счётчики нужны логу на границе: // перезапись означает, что по одним координатам приехали разные содержимые, и // это единственное наблюдение, по которому можно судить, работает ли правило // разрешения столкновений. type MergeStats struct { Buckets int // Stored — сколько точек ДОБАВИЛОСЬ в объекты. Не то же, что число точек в // доставке: точные повторы схлопываются, и счётчик присланных завышал бы // содержимое витрины ровно тогда, когда точки начнут теряться. Stored int Unchanged int Overwrites int SealedHits int // UnitsConflicts — объекты, куда точки приехали в единицах, отличных от // сохранённых. Молчаливая замена недопустима: внутри точки единиц нет, и // у ранних точек не остаётся ничего, по чему их единицы восстановимы. UnitsConflicts int // Collisions — координаты первых столкновений, для записи в лог. Без них // счётчик перезаписей не говорит, какая метрика и какой час пострадали. Collisions []Collision } // Collision — координаты объекта, где содержимое точки было перезаписано. // Значений точек не несёт: данные о здоровье чувствительнее токенов. type Collision struct { Metric string Layer string HourUTC time.Time } // maxCollisionsReported — сколько координат столкновений попадает в лог. // Больше горсти не нужно: они нужны как зацепка для разбора, а не как отчёт. const maxCollisionsReported = 5 // MergePoints раскладывает точки по часовым объектам и сливает их с // сохранёнными. // // Час берётся по НАЧАЛУ точки: интервал пересекает границы часов, и любой // другой выбор сделал бы принадлежность объекту зависящей от длительности. // Вся доставка сворачивается ОДНОЙ транзакцией, а не транзакцией на объект. // // Транзакция на объект давала недетерминированное частичное состояние: обход // карты групп рандомизирован, и при отказе посреди доставки набор уже // закоммиченных объектов каждый раз другой (измерено: восемь прогонов одной // доставки — семь разных состояний). Это ломает инвариант «состояние // пересобираемо»: пересборка из архива давала бы не то, что живой приём, а // хеш-детектор после расхождения переписывал бы «неизменившееся». // // Заодно снимается стоимость: отдельный коммит на объект стоил около 0.7 мс, // то есть 8.5 мс на килобайт тела. func (s *Store) MergePoints(ctx context.Context, in []IncomingPoint, deliveryID string) (MergeStats, error) { groups := groupByHour(in) keys := sortedKeys(groups) var stats MergeStats err := s.inTx(ctx, func(tx *sql.Tx) error { // Счётчики обнуляются на каждой попытке: повтор транзакции начинает // слияние заново, и накопленное от прошлой попытки посчиталось бы дважды. stats = MergeStats{} for _, key := range keys { group := groups[key] res, err := mergeBucket(ctx, tx, key, group, deliveryID) if err != nil { return err } stats.Buckets++ stats.Stored += res.stored stats.Overwrites += res.overwrites if res.unchanged { stats.Unchanged++ } if res.sealed { stats.SealedHits++ } if res.unitsConflict { stats.UnitsConflicts++ } if res.overwrites > 0 && len(stats.Collisions) < maxCollisionsReported { stats.Collisions = append(stats.Collisions, Collision{ Metric: key.metric, Layer: key.layer, HourUTC: key.hourUTC, }) } } 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 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 := mergePoints(stored.Points, group.points) res.overwrites = overwrites // Считаем сохранённые точки, а не присланные: точные повторы внутри // доставки схлопываются, и счётчик присланных систематически завышал бы // содержимое витрины — расхождение «прислали 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 сливает сохранённые точки с пришедшими по координатному ключу. // // При столкновении выигрывает БОЛЕЕ ПОЛНАЯ точка, а не последняя пришедшая: // 0.66% координат различаются набором полей при одинаковом значении, и правило // «последняя победила» стирало бы у сохранённой точки поля, которых новая не // несёт. // // При равной полноте исход определяется порядком канонических форм, а не // порядком доставок: у сохранённой точки нет провенанса, а четверть доставок // несёт столкновения внутри себя, где время приёма общее. Свёртка по журналу // обязана давать то же состояние, что приём в реальном времени. // // Точки из объекта не удаляются никогда. func mergePoints(stored, incoming []Point) ([]Point, int) { type coord struct { start int64 end int64 } byCoord := make(map[coord]Point, len(stored)+len(incoming)) put := func(p Point) int { c := coord{start: p.Start.UnixNano(), end: p.End.UnixNano()} old, ok := byCoord[c] if !ok { byCoord[c] = p return 0 } // Побайтовое равенство — только быстрый путь. Столкновением считается // расхождение КАНОНИЧЕСКИХ форм: канонизация и заведена потому, что // байты нестабильны. Из 81952 повторно приехавших точек 67534 // различаются лишь порядком ключей, ещё 63% — последним разрядом // double. Считай мы по байтам, счётчик перезаписей давал бы тысячи // ложных срабатываний на каждом глубоком проходе, и настоящий отказ // правила слияния стал бы неотличим от нормы. if bytes.Equal(old.Raw, p.Raw) || canon.Equal(old.Raw, p.Raw) { return 0 } byCoord[c] = resolve(old, p) return 1 } overwrites := 0 for _, p := range stored { put(p) } for _, p := range incoming { overwrites += put(p) } // Порядок точек в объекте канонический и входит в хеш: его задаёт ТОЛЬКО // сортировка ниже. Порядок обхода карты не специфицирован, и полагаться на // него значило бы получать разные хеши для одного содержимого — тогда // «неизменившийся» объект переписывался бы каждым глубоким проходом. out := make([]Point, 0, len(byCoord)) for _, p := range byCoord { out = append(out, p) } sort.Slice(out, func(i, j int) bool { if !out[i].Start.Equal(out[j].Start) { return out[i].Start.Before(out[j].Start) } return out[i].End.Before(out[j].End) }) return out, overwrites } // resolve выбирает победителя столкновения: сперва по полноте, затем // детерминированно по канонической форме. func resolve(a, b Point) Point { ca, cb := canon.Completeness(a.Raw), canon.Completeness(b.Raw) switch { case ca > cb: return a case cb > ca: return b case canon.Less(a.Raw, b.Raw): return a default: return b } } func hashPoints(points []Point) (string, error) { raws := make([][]byte, 0, len(points)) for _, p := range points { raws = append(raws, p.Raw) } h, err := canon.HashAll(raws) if err != nil { return "", fmt.Errorf("хеш объекта: %w", err) } return h, nil } func firstNonEmpty(a, b string) string { if a != "" { return a } return b } func readBucket(ctx context.Context, tx *sql.Tx, key bucketKey) (Bucket, bool, error) { const q = ` SELECT units, payload, content_hash, sealed, first_delivery_id, first_ts, last_ts FROM bucket WHERE metric = ? AND layer = ? AND hour_utc = ?` var ( units string payload []byte hash string sealed int delivery string firstTS string lastTS string ) err := tx.QueryRowContext(ctx, q, key.metric, key.layer, FormatTime(key.hourUTC)). Scan(&units, &payload, &hash, &sealed, &delivery, &firstTS, &lastTS) if errors.Is(err, sql.ErrNoRows) { return Bucket{}, false, nil } if err != nil { return Bucket{}, false, fmt.Errorf("select bucket: %w", err) } points, err := decodePayload(payload) if err != nil { return Bucket{}, false, err } first, err := ParseTime(firstTS) if err != nil { return Bucket{}, false, err } last, err := ParseTime(lastTS) if err != nil { return Bucket{}, false, err } return Bucket{ Metric: key.metric, Layer: key.layer, HourUTC: key.hourUTC, Units: units, Points: points, Hash: hash, Sealed: sealed != 0, Delivery: delivery, FirstTS: first, LastTS: last, }, true, nil } func writeBucket(ctx context.Context, tx *sql.Tx, b Bucket, now time.Time) error { payload, err := encodePayload(b.Points) if err != nil { return err } const q = ` INSERT INTO bucket (metric, layer, hour_utc, units, payload, content_hash, points, first_ts, last_ts, first_delivery_id, sealed, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (metric, layer, hour_utc) DO UPDATE SET units = excluded.units, payload = excluded.payload, content_hash = excluded.content_hash, points = excluded.points, first_ts = excluded.first_ts, last_ts = excluded.last_ts, updated_at = excluded.updated_at` _, err = tx.ExecContext(ctx, q, b.Metric, b.Layer, FormatTime(b.HourUTC), b.Units, payload, b.Hash, len(b.Points), FormatTime(b.FirstTS), FormatTime(b.LastTS), b.Delivery, boolToInt(b.Sealed), FormatTime(now), FormatTime(now)) if err != nil { return fmt.Errorf("upsert bucket: %w", err) } return nil } func boolToInt(b bool) int { if b { return 1 } return 0 } // storedPoint — форма точки внутри payload. Содержимое лежит сырым сообщением, // координаты — рядом: разжимать и перечитывать метки из содержимого пришлось // бы на каждом слиянии, а формат метки принадлежит HAE, не хранилищу. type storedPoint struct { Start string `json:"s"` End string `json:"e"` Offset int `json:"o"` Raw json.RawMessage `json:"p"` } func encodePayload(points []Point) ([]byte, error) { sp := make([]storedPoint, 0, len(points)) for _, p := range points { sp = append(sp, storedPoint{ Start: FormatTime(p.Start), End: FormatTime(p.End), Offset: p.OffsetSeconds, Raw: p.Raw, }) } // Через 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) } var buf bytes.Buffer gz := gzip.NewWriter(&buf) if _, err := gz.Write(raw.Bytes()); err != nil { return nil, fmt.Errorf("сжатие точек: %w", err) } if err := gz.Close(); err != nil { return nil, fmt.Errorf("закрытие gzip: %w", err) } return buf.Bytes(), nil } func decodePayload(payload []byte) ([]Point, error) { gz, err := gzip.NewReader(bytes.NewReader(payload)) if err != nil { return nil, fmt.Errorf("распаковка точек: %w", err) } defer func() { _ = gz.Close() }() raw, err := io.ReadAll(gz) if err != nil { return nil, fmt.Errorf("чтение точек: %w", err) } var sp []storedPoint if err := json.Unmarshal(raw, &sp); err != nil { return nil, fmt.Errorf("разбор точек: %w", err) } out := make([]Point, 0, len(sp)) for _, p := range sp { start, err := ParseTime(p.Start) if err != nil { return nil, err } end, err := ParseTime(p.End) if err != nil { return nil, err } out = append(out, Point{ Start: start, End: end, OffsetSeconds: p.Offset, Raw: p.Raw, }) } return out, nil } // MarkSealed помечает час запечатанным — тем, в который досчёта уже не // ожидается. // // Правило, по которому час переводится в это состояние, пока не определено: // порог глубины досчёта ставится по наблюдениям, которых нет (наблюдалось до // 22 минут). Метод существует раньше правила намеренно — без него признак // `sealed` и сигнал о его нарушении нечем проверить, а неопровержимый инвариант // не отличается от отсутствующего. func (s *Store) MarkSealed(ctx context.Context, metric, layer string, hour time.Time, sealed bool) error { const q = `UPDATE bucket SET sealed = ? WHERE metric = ? AND layer = ? AND hour_utc = ?` res, err := s.db.ExecContext(ctx, q, boolToInt(sealed), metric, layer, FormatTime(hour.UTC().Truncate(time.Hour))) if err != nil { return fmt.Errorf("update sealed: %w", err) } n, err := res.RowsAffected() if err != nil { return fmt.Errorf("update sealed: %w", err) } if n == 0 { return ErrNotFound } return nil } // 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 } // 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 }