- отношение победы было нетранзитивным: полнота (частичный порядок) плюс тай-брейк (тотальный) в попарной свёртке давали цикл, из-за которого одна и та же доставка меняла содержимое объекта при каждой пересборке - надмножество побеждает только при совпадении значений общих содержательных ключей: иначе точка без единого измерения вытесняла измерение - Less стал тотальным, isEmpty не материализует значение, имя метрики в координате столкновения обрезается, отпечаток витрины включает units и sealed - на живом архиве строгий no-op: 1737 объектов, содержимое совпало побайтово
311 lines
14 KiB
Go
311 lines
14 KiB
Go
// Package fold — свёртка доставки в часовые объекты.
|
|
//
|
|
// Хранилище — свёртка по журналу: состояние пересобирается как
|
|
// `import(экспорт) + replay(доставки по received_at)`. Поэтому свёртка
|
|
// принимает идентификатор доставки, а тело читает из архива — тем же кодом,
|
|
// каким его прочитает пересборка. Разбор «из памяти, раз тело всё равно в
|
|
// руках» дал бы второй путь, который разошёлся бы с первым молча.
|
|
package fold
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"log/slog"
|
|
"strings"
|
|
"time"
|
|
|
|
"git.vakhrushev.me/av/healthlog/internal/archive"
|
|
"git.vakhrushev.me/av/healthlog/internal/hae"
|
|
"git.vakhrushev.me/av/healthlog/internal/store"
|
|
)
|
|
|
|
// defaultMaxBodyBytes — граница размера тела при чтении из архива, когда
|
|
// вызывающий свою не задал.
|
|
//
|
|
// Наблюдалось 42 МиБ. Граница явная, потому что молчаливое «сколько дадут» —
|
|
// это отказ, который проявится только на пике потока: битый или враждебный
|
|
// файл архива выест память процесса.
|
|
const defaultMaxBodyBytes = 64 << 20
|
|
|
|
// Service сворачивает доставки в часовые объекты.
|
|
type Service struct {
|
|
arch *archive.Archive
|
|
store *store.Store
|
|
log *slog.Logger
|
|
|
|
// maxBody — та же граница, что у приёма, и это существенно. Тело между
|
|
// границами приёма и свёртки было бы принято со статусом 200, легло бы в
|
|
// архив и потом вечно валилось бы при каждой пересборке.
|
|
maxBody int64
|
|
}
|
|
|
|
// New собирает свёртку. maxBody — граница размера РАСПАКОВАННОГО тела; ноль
|
|
// означает умолчание.
|
|
func New(arch *archive.Archive, st *store.Store, maxBody int64, log *slog.Logger) *Service {
|
|
if maxBody <= 0 {
|
|
maxBody = defaultMaxBodyBytes
|
|
}
|
|
return &Service{
|
|
arch: arch,
|
|
store: st,
|
|
maxBody: maxBody,
|
|
log: log.With("capability", "fold"),
|
|
}
|
|
}
|
|
|
|
// finishTimeout — сколько отводится записи исхода свёртки.
|
|
//
|
|
// Исход пишется на контексте, ПЕРЕЖИВАЮЩЕМ отмену исходного: иначе при
|
|
// срабатывании дедлайна свёртки запись статуса гарантированно провалится, и
|
|
// доставка навсегда останется `pending` — притом что часть объектов уже
|
|
// записана. То есть ровно в том случае, ради которого дедлайн и заведён,
|
|
// учёт разошёлся бы с содержимым витрины.
|
|
const finishTimeout = 10 * time.Second
|
|
|
|
// Stats — итог свёртки одной доставки.
|
|
type Stats struct {
|
|
Metrics int
|
|
Points int
|
|
Stored int
|
|
Buckets int
|
|
Unchanged int
|
|
Overwrites int
|
|
Incomparable int
|
|
SealedHits int
|
|
UnitsConflicts int
|
|
SkippedNoTime int
|
|
SkippedMalformed int
|
|
SkippedBadEnd int
|
|
Layer string
|
|
LayerMismatch bool
|
|
Collisions []store.Collision
|
|
IncomparableAt []store.Collision
|
|
}
|
|
|
|
// Fold разбирает тело доставки и раскладывает точки по часовым объектам.
|
|
//
|
|
// Это единственный логирующий чекпоинт свёртки: транспорт и приём исход
|
|
// разбора не логируют. Значения точек и имена устройств в лог не попадают —
|
|
// данные о здоровье чувствительнее токенов.
|
|
func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) {
|
|
var stats Stats
|
|
|
|
d, err := s.store.DeliveryForParse(ctx, deliveryID)
|
|
if err != nil {
|
|
// Ни один выход с ошибкой не молчит: вызывающий эту ошибку сознательно
|
|
// отбрасывает, и молчащий путь означал бы доставку вообще без событий
|
|
// в журнале — её не найти ни по какому запросу.
|
|
s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID)
|
|
return stats, err
|
|
}
|
|
|
|
body, err := s.readBody(d.RawPath)
|
|
if err != nil {
|
|
s.fail(ctx, deliveryID, err)
|
|
return stats, err
|
|
}
|
|
|
|
fallback, err := s.store.LastDerivedLayer(ctx, d.AutomationID, d.ReceivedAt, d.ID)
|
|
if err != nil {
|
|
s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID)
|
|
return stats, err
|
|
}
|
|
|
|
parsed, err := hae.Parse(body, hae.Meta{
|
|
Aggregation: d.Aggregation,
|
|
FallbackLayer: hae.Layer(fallback),
|
|
})
|
|
if err != nil {
|
|
s.fail(ctx, deliveryID, err)
|
|
return stats, err
|
|
}
|
|
|
|
stats.Metrics = parsed.Metrics
|
|
stats.Points = len(parsed.Points)
|
|
stats.SkippedNoTime = parsed.SkippedNoTime
|
|
stats.SkippedMalformed = parsed.SkippedMalformed
|
|
stats.SkippedBadEnd = parsed.SkippedBadEnd
|
|
stats.Layer = string(parsed.Layer)
|
|
stats.LayerMismatch = parsed.LayerMismatch
|
|
|
|
merge, err := s.store.MergePoints(ctx, toIncoming(parsed.Points), deliveryID)
|
|
if err != nil {
|
|
s.fail(ctx, deliveryID, err)
|
|
return stats, err
|
|
}
|
|
|
|
stats.Stored = merge.Stored
|
|
stats.Buckets = merge.Buckets
|
|
stats.Unchanged = merge.Unchanged
|
|
stats.Overwrites = merge.Overwrites
|
|
stats.Incomparable = merge.Incomparable
|
|
stats.SealedHits = merge.SealedHits
|
|
stats.UnitsConflicts = merge.UnitsConflicts
|
|
stats.Collisions = merge.Collisions
|
|
stats.IncomparableAt = merge.IncomparableAt
|
|
|
|
if err := s.finish(ctx, deliveryID, store.ParseDone, int64(stats.Points), stats.Layer); err != nil {
|
|
s.log.ErrorContext(ctx, "delivery fold failed", "error", err, "delivery_id", deliveryID)
|
|
return stats, err
|
|
}
|
|
|
|
s.logResult(ctx, deliveryID, stats)
|
|
return stats, nil
|
|
}
|
|
|
|
// logResult — единственный логирующий чекпоинт свёртки.
|
|
//
|
|
// Все признаки идут АТРИБУТАМИ всегда, а уровень выбирается отдельно. Раньше
|
|
// признаки жили только в тексте сообщения, и `switch` их терял: доставка,
|
|
// одновременно задевшая запечатанный час и разошедшаяся с заголовком, не
|
|
// оставляла следа о слое вовсе — а это единственный индикатор того, что слой
|
|
// выводится неправильно.
|
|
//
|
|
// Значений точек и имён устройств здесь нет и быть не может: данные о здоровье
|
|
// чувствительнее токенов. Координаты столкновений — метрика, слой, час — не
|
|
// значения.
|
|
func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) {
|
|
skipped := st.SkippedNoTime + st.SkippedMalformed + st.SkippedBadEnd
|
|
|
|
attrs := []any{
|
|
"delivery_id", deliveryID,
|
|
"metrics", st.Metrics,
|
|
"points", st.Points,
|
|
"stored", st.Stored,
|
|
"buckets", st.Buckets,
|
|
"unchanged", st.Unchanged,
|
|
"overwrites", st.Overwrites,
|
|
"incomparable", st.Incomparable,
|
|
"units_conflicts", st.UnitsConflicts,
|
|
"sealed_hits", st.SealedHits,
|
|
"skipped", skipped,
|
|
"skipped_no_time", st.SkippedNoTime,
|
|
"skipped_malformed", st.SkippedMalformed,
|
|
"skipped_bad_end", st.SkippedBadEnd,
|
|
"layer", st.Layer,
|
|
"layer_mismatch", st.LayerMismatch,
|
|
}
|
|
if len(st.Collisions) > 0 {
|
|
attrs = append(attrs, "collisions", formatCollisions(st.Collisions))
|
|
}
|
|
if len(st.IncomparableAt) > 0 {
|
|
attrs = append(attrs, "incomparable_at", formatCollisions(st.IncomparableAt))
|
|
}
|
|
|
|
// Доставка, у которой отброшены ВСЕ точки, — это сломавшийся формат, а не
|
|
// штатная работа. Без этого условия смена формата метки выглядела бы как
|
|
// здоровый поток: 200, parsed, INFO, points=0.
|
|
allSkipped := st.Points == 0 && skipped > 0
|
|
|
|
switch {
|
|
case st.Incomparable > 0:
|
|
// Выше перезаписей намеренно: несравнимый набор полей — событие реже и
|
|
// информативнее, на живом потоке не случавшееся ни разу. Признаки при
|
|
// этом идут атрибутами всегда, так что выбор ветви ничего не прячет.
|
|
s.log.WarnContext(ctx, "delivery folded, incomparable point fields", attrs...)
|
|
case st.Overwrites > 0:
|
|
// Единственное наблюдение, по которому проверяется правило слияния.
|
|
// В INFO оно тонуло: поток идёт раз в пять минут.
|
|
s.log.WarnContext(ctx, "delivery folded, points overwritten", attrs...)
|
|
case st.UnitsConflicts > 0:
|
|
s.log.WarnContext(ctx, "delivery folded, units differ from stored", attrs...)
|
|
case st.SealedHits > 0:
|
|
// Досчёт часа, в который его уже не ждали: единственное наблюдение, по
|
|
// которому вообще можно судить о глубине досчёта.
|
|
s.log.WarnContext(ctx, "delivery folded, sealed hour changed", attrs...)
|
|
case allSkipped:
|
|
s.log.WarnContext(ctx, "delivery folded, all points skipped", attrs...)
|
|
case st.LayerMismatch:
|
|
// Расхождение сверяется только с надёжным заголовком: `Default` не
|
|
// означает режима, и сравнение с ним давало бы WARN на каждой доставке.
|
|
s.log.WarnContext(ctx, "delivery folded, layer differs from header", attrs...)
|
|
default:
|
|
s.log.InfoContext(ctx, "delivery folded", attrs...)
|
|
}
|
|
}
|
|
|
|
// formatCollisions превращает координаты столкновений в строку для лога.
|
|
func formatCollisions(cs []store.Collision) string {
|
|
parts := make([]string, 0, len(cs))
|
|
for _, c := range cs {
|
|
parts = append(parts, c.Metric+"/"+c.Layer+"@"+store.FormatTime(c.HourUTC))
|
|
}
|
|
return strings.Join(parts, " ")
|
|
}
|
|
|
|
// keepLayer — значение слоя, означающее «оставить как было».
|
|
const keepLayer = ""
|
|
|
|
// finish записывает исход разбора на контексте, переживающем отмену исходного.
|
|
func (s *Service) finish(ctx context.Context, deliveryID, status string, points int64, layer string) error {
|
|
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), finishTimeout)
|
|
defer cancel()
|
|
|
|
if err := s.store.FinishParse(ctx, deliveryID, status, points, layer); err != nil {
|
|
return fmt.Errorf("запись исхода разбора: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// fail отмечает доставку неразобранной. Тело остаётся в архиве, и её подберёт
|
|
// пересборка — приём при этом не затрагивается: сохранили значит приняли.
|
|
func (s *Service) fail(ctx context.Context, deliveryID string, cause error) {
|
|
level := slog.LevelError
|
|
switch {
|
|
case errors.Is(cause, hae.ErrLayerUnknown):
|
|
// Слой не определился — это не поломка, а ожидаемый исход для доставки
|
|
// без плотных метрик. Тело ждёт пересборки.
|
|
level = slog.LevelWarn
|
|
case errors.Is(cause, hae.ErrMalformed):
|
|
// Непонятое содержимое от отправителя — норма жизни, разбирать нечего.
|
|
// На границе приёма такой же отказ уходит в DEBUG; два разных уровня у
|
|
// одного класса ошибки давали бы постоянный ERROR-шум.
|
|
level = slog.LevelWarn
|
|
}
|
|
s.log.Log(ctx, level, "delivery fold failed", "error", cause, "delivery_id", deliveryID)
|
|
|
|
// Слой НЕ затирается: доставка могла свернуться успешно раньше, и пустая
|
|
// строка здесь оборвала бы цепочку наследования, то есть изменила бы
|
|
// результат пересборки журнала.
|
|
if err := s.finish(ctx, deliveryID, store.ParseFailed, 0, keepLayer); err != nil {
|
|
s.log.ErrorContext(ctx, "delivery parse status not recorded", "error", err, "delivery_id", deliveryID)
|
|
}
|
|
}
|
|
|
|
func (s *Service) readBody(rawPath string) ([]byte, error) {
|
|
r, err := s.arch.Open(rawPath)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer func() { _ = r.Close() }()
|
|
|
|
body, err := io.ReadAll(io.LimitReader(r, s.maxBody+1))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("чтение тела из архива: %w", err)
|
|
}
|
|
if int64(len(body)) > s.maxBody {
|
|
return nil, fmt.Errorf("тело больше %d байт", s.maxBody)
|
|
}
|
|
return body, nil
|
|
}
|
|
|
|
func toIncoming(points []hae.Point) []store.IncomingPoint {
|
|
out := make([]store.IncomingPoint, 0, len(points))
|
|
for _, p := range points {
|
|
out = append(out, store.IncomingPoint{
|
|
Metric: p.Metric,
|
|
Layer: string(p.Layer),
|
|
Units: p.Units,
|
|
Point: store.Point{
|
|
Start: p.Start,
|
|
End: p.End,
|
|
OffsetSeconds: p.OffsetSeconds,
|
|
Raw: p.Raw,
|
|
},
|
|
})
|
|
}
|
|
return out
|
|
}
|