Files
healthlog/internal/fold/fold.go
T
av 63bffe2865 Приём отвечает 200 до свёртки, свёртку ведёт фоновый воркер
- Очередью служит сама таблица: доставка ждёт свёртки в статусе `pending`,
  канал несёт только бит «есть работа». Переполнять нечего, падение процесса
  очередь не теряет, а подбор `pending` при старте — обычный проход воркера, а
  не отдельный код. Классификация исхода общая с пересборкой журнала.
- Исход разбора начал отражать доставку, а не обстоятельства: отмена и
  занятость базы статус не меняют (иначе конкуренция за базу выводила бы
  доставку из очереди навсегда), паника свёртки больше не валит процесс, а
  учёт доставки идёт через транзакцию с повторами.
- Длинный бюджет ответа выдан маршруту приёма, а не всему серверу:
  `write_timeout` в Go покрывает и чтение тела, и общий подъём снял бы защиту с
  остальных маршрутов.
2026-08-02 11:01:42 +03:00

399 lines
20 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 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
// Uncovered — верхнеуровневые ключи `data`, которых разбор не покрывает.
// Это ответ на вопрос «что останется потерянным, если тело удалить»:
// для stateOfMind он необратим — в экспорте Apple этой секции нет.
Uncovered []string
// UncoveredDropped — сколько имён отброшено границей списка.
UncoveredDropped int
}
// ErrPanicked — свёртка паниковала. Доставка получает `failed`: тело в архиве, и
// пересборка вернёт её, когда дефект будет исправлен.
var ErrPanicked = errors.New("свёртка паниковала")
// Fold разбирает тело доставки и раскладывает точки по часовым объектам.
//
// Это единственный логирующий чекпоинт свёртки: транспорт и приём исход
// разбора не логируют. Значения точек и имена устройств в лог не попадают —
// данные о здоровье чувствительнее токенов.
//
// Паника перехватывается ЗДЕСЬ, у той же границы, что пишет исход разбора.
// Пока свёртка шла внутри HTTP-обработчика, панику ловил middleware.Recoverer и
// она стоила одного ответа; из фоновой горутины она валит процесс целиком, а
// `restart: unless-stopped` поднимает его снова — и первый же проход берёт ту
// же доставку, то есть дефект превращается в цикл перезапуска, при котором
// приём не работает вовсе. Перехват у этой границы, а не у вызывающего,
// оставляет писателя `parse_status` единственным.
func (s *Service) Fold(ctx context.Context, deliveryID string) (stats Stats, err error) {
defer func() {
r := recover()
if r == nil {
return
}
stats = Stats{}
err = fmt.Errorf("%w: %v", ErrPanicked, r) //nolint:errorlint // причину раскрываем текстом, sentinel — для ветвления
s.fail(ctx, deliveryID, err, nil)
}()
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, nil)
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, parsed.Uncovered)
return stats, err
}
stats.Uncovered = parsed.Uncovered
stats.UncoveredDropped = parsed.UncoveredDropped
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, parsed.Uncovered)
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
// Источник истины — список; статус производен от него и от факта отказа.
// Приоритет назван явно, иначе два будущих читателя (ретеншен и /stats)
// разойдутся: один спросит parse_status, другой — непустоту списка.
status := store.ParseDone
if len(parsed.Uncovered) > 0 {
status = store.ParsePartial
}
out := store.ParseOutcome{
Status: status,
Points: int64(stats.Points),
Layer: stats.Layer,
Uncovered: parsed.Uncovered,
}
if err := s.finish(ctx, deliveryID, out); 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,
// Структурным []string, а не склейкой: JSON-кодировщик slog экранирует
// управляющие символы, поэтому имя секции из чужого тела не разрывает
// построчный разбор логов. Содержимого секций здесь нет.
"uncovered", st.Uncovered,
"uncovered_dropped", st.UncoveredDropped,
}
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.UncoveredDropped > 0:
// Не частичный разбор, а тело, не похожее на HAE: секций у HAE восемь,
// а границу выбило больше тридцати двух.
s.log.WarnContext(ctx, "delivery folded, uncovered section list truncated", attrs...)
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 string, out store.ParseOutcome) error {
ctx, cancel := context.WithTimeout(context.WithoutCancel(ctx), finishTimeout)
defer cancel()
if err := s.store.FinishParse(ctx, deliveryID, out); err != nil {
return fmt.Errorf("запись исхода разбора: %w", err)
}
return nil
}
// fail записывает исход неудачной свёртки.
//
// Исход отражает ДОСТАВКУ, а не обстоятельства. Отмена снаружи и занятость базы
// работой доставки не являются: они означают «не сделано», а не «не выходит».
// Статус в этих случаях не трогается вовсе — доставка остаётся `pending` и
// подбирается следующим проходом. Иначе конкуренция за базу (после разнесения
// ответа и свёртки она штатная) выводила бы доставку из очереди навсегда:
// `failed` возвращает только пересборка, то есть ручная операция с остановкой
// сервиса.
//
// Всё прочее — непонятое содержимое, невыводимый слой, нечитаемое или слишком
// большое тело, исчерпанный дедлайн — свойства самой доставки, и повторять их
// бесполезно: статус `failed`, тело ждёт пересборки. Приём при этом не
// затрагивается: сохранили значит приняли.
func (s *Service) fail(ctx context.Context, deliveryID string, cause error, uncovered []string) {
if store.Transient(cause) {
// WARN, а не ERROR: пройдёт само, разбирать нечего. Строка нужна, чтобы
// повтор не выглядел беспричинным.
s.log.WarnContext(ctx, "delivery fold deferred", "error", cause, "delivery_id", deliveryID)
return
}
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)
// Слой НЕ затирается: доставка могла свернуться успешно раньше, и пустая
// строка здесь оборвала бы цепочку наследования, то есть изменила бы
// результат пересборки журнала.
out := store.ParseOutcome{
Status: store.ParseFailed,
Layer: keepLayer,
Uncovered: uncovered,
}
if err := s.finish(ctx, deliveryID, out); err != nil {
s.log.ErrorContext(ctx, "delivery parse status not recorded", "error", err, "delivery_id", deliveryID)
}
}
// ReadBody читает тело из архива с той же границей размера, что и свёртка.
//
// Экспортировано ради пересборки: ей нужно прочесть тело, у которого ещё нет
// учётной записи, чтобы посчитать размер и хеш. Своей копией чтения это делать
// нельзя — граница обязана быть общей, иначе тело, принятое приёмом со `200`,
// начнёт вечно отказывать на каждой пересборке, и договорённость «предел тот
// же» ничем не проверяется.
func (s *Service) ReadBody(rawPath string) ([]byte, error) {
return s.readBody(rawPath)
}
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
}