// 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, parseResidue{}) }() 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, parseResidue{}) 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, residueOf(parsed)) 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, parseResidue{ uncovered: parsed.Uncovered, skipped: skippedEntities(parsed), }) 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, SkippedEntities: skippedEntities(parsed), } 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, "entities_diverging", st.EntitiesDiverging, "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)) } if len(st.DivergingAt) > 0 { attrs = append(attrs, "diverging_at", formatEntityRefs(st.DivergingAt)) } // Доставка, у которой отброшены ВСЕ точки, — это сломавшийся формат, а не // штатная работа. Без этого условия смена формата метки выглядела бы как // здоровый поток: 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 st.EntitiesDiverging > 0: // Событие другого рода и с другим лечением: в одном теле приехали // версии одного ключа с разным содержанием. Победитель лёг в витрину // целиком, терять нечего — но корпус такого не производил, и молчать // об этом нельзя. Отдельной ветвью, а не общей с удержанием: сообщение // «удержана обеднённая версия» отправляло бы владельца искать то, чего // не случилось. s.log.WarnContext(ctx, "delivery folded, entity versions diverge in one body", 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`, тело ждёт пересборки. Приём при этом не // затрагивается: сохранили значит приняли. // parseResidue — то, что разбор успел узнать о доставке до отказа и что обязано // пережить его в учёте: список непокрытых секций и число пропущенных сущностей. // // Структурой, а не двумя параметрами: у `fail` их стало бы четыре, и следующий // счётчик неизбежно перепутали бы местами с предыдущим. Пустое значение — // «разбор до этого не дошёл», и оно честно: отказ на чтении тела ничего о // содержимом не знает. type parseResidue struct { uncovered []string // skipped — nil означает «разбор до конца не дошёл, пропусков никто не // считал». Ноль означал бы «проверено, терять нечего», а по этому числу // ретеншен принимает необратимое решение об удалении тела. skipped *int64 } func residueOf(parsed hae.Result) parseResidue { return parseResidue{uncovered: parsed.Uncovered} } // skippedEntities — сколько сущностей с собственным `id` разбор пропустил. // Сумма трёх классов, а не три колонки: ретеншен спрашивает «есть ли что // терять», а не «почему», а разбор класса живёт в логе свёртки, где все три // счётчика идут атрибутами. func skippedEntities(parsed hae.Result) *int64 { n := int64(parsed.SkippedNoID + parsed.SkippedEntityNoTime + parsed.SkippedEntityMalformed) return &n } func (s *Service) fail(ctx context.Context, deliveryID string, cause error, residue parseResidue) { 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: residue.uncovered, SkippedEntities: residue.skipped, } 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 }