Files
healthlog/internal/store/bucket.go
T
av cd7a4c1493 отработаны находки ревью кода
Триаж свёл 62 сырые находки девяти проходов к 33 причинам: 3 блокера,
4 «сейчас», 2 развилки. Все закрыты регрессионными тестами.

- схема точки сна определяется по самой точке, а не по индексу в исходном
  массиве: одна пропущенная точка меняла эпизод и сводку местами
- доставка сворачивается одной транзакцией: частичное состояние было
  недетерминированным (восемь прогонов — семь состояний)
- граница размера на распакованном теле: 400 КиБ gzip разворачивались в
  400 МиБ мимо max_body_mb
- столкновение — расхождение канонических форм, а не байтов; WARN с
  координатами объекта; payload без HTML-экранирования
- выравнивание по местной метке: получасовые зоны уводили часовую выгрузку
  в minute
- доставка из одних суточных сводок больше не отвергается целиком
- единицы не переписываются молча; счётчик считает сохранённые точки
- все выходы Fold логируются, исход пишется на переживающем отмену контексте

Четыре развилки вынесены блокерами в беклог.
2026-08-01 19:02:05 +03:00

600 lines
22 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"
"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
}