- секции `workouts` и `stateOfMind` покрыты разбором: тренировка лежит одной строкой вместе с маршрутом и внутренними рядами, запись — по ключу `род + id`; миграция 00007 заводит обе таблицы и возвращает в очередь `partial`-доставки с этими ключами - сущность заменяется целиком, но условно: приехавшая побеждает, если не теряет содержания сохранённой (множество ключей и длины верхнеуровневых массивов), а при равном содержании выигрывает версия из более поздней доставки ЖУРНАЛА — «побеждает приехавшая» было бы функцией порядка свёртки, и живая витрина расходилась бы с пересборкой молча - отпечаток витрины покрывает тренировки и записи и снимается одним снимком базы; отчёт `reindex` считает «было и стало» по каждой единице хранения
456 lines
23 KiB
Go
456 lines
23 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 {
|
||
store.MergeStats
|
||
|
||
Metrics int
|
||
Points int
|
||
SkippedNoTime int
|
||
SkippedMalformed int
|
||
SkippedBadEnd int
|
||
// Счётчики пропуска сущностей: у каждого класса свой, потому что тело в
|
||
// архиве остаётся, а вернуть сущность может только пересборка.
|
||
SkippedNoID int
|
||
SkippedEntityNoTime int
|
||
SkippedEntityMalformed int
|
||
Layer string
|
||
LayerMismatch bool
|
||
// Uncovered — верхнеуровневые ключи `data`, которых разбор не покрывает.
|
||
// Это ответ на вопрос «что останется потерянным, если тело удалить».
|
||
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.SkippedNoID = parsed.SkippedNoID
|
||
stats.SkippedEntityNoTime = parsed.SkippedEntityNoTime
|
||
stats.SkippedEntityMalformed = parsed.SkippedEntityMalformed
|
||
stats.Layer = string(parsed.Layer)
|
||
stats.LayerMismatch = parsed.LayerMismatch
|
||
|
||
merge, err := s.store.Merge(ctx, toIncoming(parsed), store.DeliveryRef{
|
||
ID: d.ID,
|
||
ReceivedAt: d.ReceivedAt,
|
||
})
|
||
if err != nil {
|
||
s.fail(ctx, deliveryID, err, parsed.Uncovered)
|
||
return stats, err
|
||
}
|
||
stats.MergeStats = merge
|
||
|
||
// Источник истины — список; статус производен от него и от факта отказа.
|
||
// Приоритет назван явно, иначе два будущих читателя (ретеншен и /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
|
||
skippedEntities := st.SkippedNoID + st.SkippedEntityNoTime + st.SkippedEntityMalformed
|
||
|
||
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,
|
||
// Сущности с собственным `id`: пришло, легло и удержано. Координаты
|
||
// (род и идентификатор) разрешены, содержимое — нет: маршрут
|
||
// тренировки это геотрек до дома, а метки состояния разума —
|
||
// измерение душевного состояния.
|
||
"workouts", st.Workouts,
|
||
"workouts_written", st.WorkoutsWritten,
|
||
"records", st.Records,
|
||
"records_written", st.RecordsWritten,
|
||
"entities_held", st.EntitiesHeld,
|
||
"skipped_entities", skippedEntities,
|
||
"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))
|
||
}
|
||
if len(st.HeldAt) > 0 {
|
||
attrs = append(attrs, "held_at", formatEntityRefs(st.HeldAt))
|
||
}
|
||
|
||
// Доставка, у которой отброшены ВСЕ точки, — это сломавшийся формат, а не
|
||
// штатная работа. Без этого условия смена формата метки выглядела бы как
|
||
// здоровый поток: 200, parsed, INFO, points=0.
|
||
allSkipped := st.Points == 0 && skipped > 0
|
||
// То же для сущностей: доставка из одних тренировок, у которой не осталось
|
||
// ни одной, — это сменившийся формат, а не пустая секция.
|
||
allEntitiesSkipped := st.Workouts == 0 && st.Records == 0 && skippedEntities > 0
|
||
|
||
switch {
|
||
case st.EntitiesHeld > 0:
|
||
// Приехавшая версия сущности отклонена как теряющая содержание. Плата
|
||
// за отказ объединять поля: событие обязано быть видно, потому что на
|
||
// живом потоке оно не наступало ни разу и правило держится на этом.
|
||
s.log.WarnContext(ctx, "delivery folded, poorer entity version held", attrs...)
|
||
case allEntitiesSkipped:
|
||
s.log.WarnContext(ctx, "delivery folded, all entities skipped", attrs...)
|
||
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, " ")
|
||
}
|
||
|
||
// formatEntityRefs превращает координаты сущностей в строку для лога.
|
||
// Содержимого не несёт: род и идентификатор — координаты, а не измерение.
|
||
func formatEntityRefs(refs []store.EntityRef) string {
|
||
parts := make([]string, 0, len(refs))
|
||
for _, r := range refs {
|
||
parts = append(parts, r.Kind+"/"+r.ID)
|
||
}
|
||
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(parsed hae.Result) store.Incoming {
|
||
points := make([]store.IncomingPoint, 0, len(parsed.Points))
|
||
for _, p := range parsed.Points {
|
||
points = append(points, 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 store.Incoming{
|
||
Points: points,
|
||
Workouts: toEntities(parsed.Workouts),
|
||
Records: toEntities(parsed.Records),
|
||
}
|
||
}
|
||
|
||
func toEntities(in []hae.Entity) []store.IncomingEntity {
|
||
if len(in) == 0 {
|
||
return nil
|
||
}
|
||
out := make([]store.IncomingEntity, 0, len(in))
|
||
for _, e := range in {
|
||
out = append(out, store.IncomingEntity{
|
||
ID: e.ID,
|
||
Kind: e.Kind,
|
||
Name: e.Name,
|
||
Start: e.Start,
|
||
End: e.End,
|
||
OffsetSeconds: e.OffsetSeconds,
|
||
Duration: e.Duration,
|
||
Raw: e.Raw,
|
||
})
|
||
}
|
||
return out
|
||
}
|