- internal/fold — свёртка по идентификатору доставки, тело из архива: тот же код, каким пойдёт пересборка витрины - приём зовёт свёртку на context.WithoutCancel с собственным дедлайном; исход разбора на код ответа не влияет - слой доставки хранится в delivery.derived_layer и наследуется от ПРЕДШЕСТВУЮЩЕЙ доставки автоматизации: без границы по времени свёртка переставала быть функцией от префикса журнала (1737 объектов против 1742) - task verify:archive — сходимость на живом архиве, 99 доставок из 99
511 lines
15 KiB
Go
511 lines
15 KiB
Go
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
|
|
}
|
|
|
|
// 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
|
|
}
|