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 Points int Unchanged int Overwrites int SealedHits int } // MergePoints раскладывает точки по часовым объектам и сливает их с // сохранёнными. // // Час берётся по НАЧАЛУ точки: интервал пересекает границы часов, и любой // другой выбор сделал бы принадлежность объекту зависящей от длительности. func (s *Store) MergePoints(ctx context.Context, in []IncomingPoint, deliveryID string) (MergeStats, error) { var stats MergeStats for key, group := range groupByHour(in) { res, err := s.mergeBucket(ctx, key, group, deliveryID) if err != nil { return stats, err } stats.Buckets++ stats.Points += len(group.points) stats.Overwrites += res.overwrites if res.unchanged { stats.Unchanged++ } if res.sealed { stats.SealedHits++ } } return stats, nil } // 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 unchanged bool sealed bool } // mergeBucket выполняет чтение, слияние и запись одного объекта одной // транзакцией. // // Тройка обязана быть атомарной целиком: между чтением и записью может // вклиниться другая доставка того же часа, и её точки пропали бы молча. func (s *Store) mergeBucket(ctx context.Context, key bucketKey, group *pointGroup, deliveryID string) (mergeResult, error) { var res mergeResult err := s.inTx(ctx, func(tx *sql.Tx) error { res = mergeResult{} stored, found, err := readBucket(ctx, tx, key) if err != nil { return err } merged, overwrites := mergePoints(stored.Points, group.points) res.overwrites = overwrites hash, err := hashPoints(merged) if err != nil { return err } if found && hash == stored.Hash { // Хеш — детектор изменений: совпал, значит писать нечего. Именно // это делает широкие проходы синхронизации дешёвыми. res.unchanged = true return nil } res.sealed = found && stored.Sealed now := Now() b := Bucket{ Metric: key.metric, Layer: key.layer, HourUTC: key.hourUTC, Units: firstNonEmpty(group.units, stored.Units), Points: merged, Hash: hash, Sealed: stored.Sealed, FirstTS: merged[0].Start, LastTS: merged[len(merged)-1].Start, Delivery: firstNonEmpty(stored.Delivery, deliveryID), } return writeBucket(ctx, tx, b, now) }) if 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)) order := make([]coord, 0, 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 order = append(order, c) return 0 } if bytes.Equal(old.Raw, p.Raw) { return 0 } winner := resolve(old, p) byCoord[c] = winner return 1 } overwrites := 0 for _, p := range stored { put(p) } for _, p := range incoming { overwrites += put(p) } out := make([]Point, 0, len(order)) for _, c := range order { out = append(out, byCoord[c]) } 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, }) } raw, err := json.Marshal(sp) if err != nil { return nil, fmt.Errorf("сериализация точек: %w", err) } var buf bytes.Buffer gz := gzip.NewWriter(&buf) if _, err := gz.Write(raw); 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 } // 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 }