// Package replay — проигрывание журнала доставок в витрину. // // Состояние healthlog есть свёртка по журналу: // `import(снапшот экспорта) + replay(доставки по received_at)`. Здесь живёт // вторая половина формулы; пересборка из архива — её вырожденный случай с // пустым снапшотом, а не отдельная операция. Когда появится импорт родного // экспорта Apple, он добавит стадию снапшота ПЕРЕД проигрыванием и переиспользует // эту же операцию. // // Собственного разбора и собственного слияния пакет не имеет: он зовёт ту же // свёртку, что и приём, по идентификатору доставки. Второй путь разбора // разошёлся бы с первым молча — и уже расходился: прежний прогон живого архива // сортировал тела по путям, а не по времени приёма. package replay import ( "context" "crypto/sha256" "encoding/hex" "errors" "fmt" "log/slog" "sort" "git.vakhrushev.me/av/healthlog/internal/archive" "git.vakhrushev.me/av/healthlog/internal/fold" "git.vakhrushev.me/av/healthlog/internal/ident" "git.vakhrushev.me/av/healthlog/internal/store" ) // Options — что нужно проигрыванию. type Options struct { // Archive — журнал тел. Именно он источник состава: перечислять только // строки учёта значило бы пересобирать витрину из витрины. Archive *archive.Archive // Source — рабочая база, откуда берётся учёт доставок. Открывается на // чтение: заголовков доставки в архиве нет, а восстановить их неоткуда. // Пустой источник допустим — тогда весь журнал состоит из подобранных тел. Source *store.Store // Target — база назначения, куда собирается витрина. Target *store.Store // Fold — свёртка над Target. Собирается вызывающим с тем же пределом // размера тела, что и у приёма; тем же пределом пересборка читает тело при // подборе — своей копии границы у неё нет намеренно. Fold *fold.Service // Progress зовётся по мере продвижения. Нужен потому, что прогон на полном // архиве молчит минутами, и зависший неотличим от идущего. Progress func(done, total int) Log *slog.Logger } // Report — итог проигрывания. Значений точек не несёт: содержимое входит в // отчёт только отпечатком. type Report struct { // Bodies — тел в архиве, признанных телами. Bodies int // SkippedFiles — файлов, телом не являющихся (чужое расширение, остаток // прерванной записи, имя не разбирается как ULID). SkippedFiles int // Duplicates — тел, чей идентификатор уже встретился в другом каталоге. // Проигрывается первое; второе считается, а не роняет прогон. Duplicates int // Adopted — тел, у которых учётной записи не было: приём успел записать // тело и не успел строку. Adopted int // AdoptFailed — тел без учёта, которые не удалось прочитать, чтобы завести // запись. Тело остаётся в архиве. AdoptFailed int // Orphans — строк учёта, у которых тела в архиве нет. Станет штатным, когда // появится ретеншен архива. Orphans int // Outcome — счётчики свёртки, те же самые, что считает фоновый воркер // приёма. Встроены, а не продублированы именами: два набора имён для одних // исходов разошлись бы при первой же правке классификации. Outcome Buckets int64 // Workouts и Records — остальные единицы хранения витрины. Считаются рядом // с объектами потому, что отпечаток отвечает «да/нет» за витрину целиком, а // решение о подмене базы необратимо и требует направления расхождения. Workouts int64 Records int64 // Categories — строки реестра категориальных значений: четвёртая единица // хранения витрины. Без счётчика расхождение по ней безадресно — объекты, // тренировки и записи при этом не меняются вовсе. Categories int64 Fingerprint string // Canceled — проигрывание прервано отменой, а не дошло до конца. Canceled bool } // Run проигрывает журнал в базу назначения. // // Порядок строго `(received_at, id)`: слой доставки без плотных метрик // наследуется от ПРЕДШЕСТВУЮЩЕЙ доставки той же автоматизации, то есть является // функцией префикса журнала. Обход каталога совпадает с хронологией только по // датам каталогов и внутри суток не упорядочивает ничего. // // Проигрывание последовательное. Распараллеливать его нельзя по той же причине: // слой зависит от префикса, а запись часового объекта — это чтение, слияние и // запись обратно. func Run(ctx context.Context, o Options) (Report, error) { var rep Report log := o.Log if log == nil { log = slog.New(slog.DiscardHandler) } log = log.With("capability", "replay") // База назначения обязана быть пустой: проигрывание поверх накопленного // оставило бы в витрине результат прежнего разбора — точки из объекта не // удаляются никогда, — то есть не сделало бы того, ради чего пересборка и // существует. n, err := o.Target.CountDeliveries(ctx) if err != nil { return stopOr(rep, err) } if n > 0 { return rep, fmt.Errorf("база назначения не пуста: %d доставок", n) } journal, orphans, rep, err := o.collect(ctx, rep, log) if err != nil { return stopOr(rep, err) } // Строка старта: прогон идёт минутами, и убитый на середине не оставлял бы // о себе в логе ни следа — только чужие с виду записи свёртки. log.InfoContext(ctx, "journal replay started", "archive_dir", o.Archive.Root(), "bodies", len(journal), "orphans", len(orphans)) // Учётные записи, тела которых в архиве нет, переносятся ПЕРВЫМИ и не // сворачиваются. Не перенести их значило бы потерять факты журнала — // заголовки, хеш, автоматизацию — при первой же подмене базы: тела уже нет, // и восстановить их будет нечем. В наследовании слоя они не участвуют: // выведенного слоя у них нет. for _, d := range orphans { if err := o.Target.CreateDelivery(ctx, d); err != nil { return stopOr(rep, err) } } player := Player{Fold: o.Fold} for i, d := range journal { if ctx.Err() != nil { rep.Canceled = true return rep, nil } if err := o.Target.CreateDelivery(ctx, d); err != nil { return stopOr(rep, err) } out, _ := player.Play(ctx, d.ID) // Отмена, застигшая свёртку, — не отказ доставки: считать её отказом // значило бы обвинить разбор в том, чего он не делал, и отправить // человека искать дефект по логу. Решение принимает цикл по СВОЕМУ // контексту — тому же, на котором шла свёртка; классификатор о нём не // знает намеренно, у воркера тот же признак означает другое. if ctx.Err() != nil { rep.Canceled = true return rep, nil } rep.Add(out) if o.Progress != nil { o.Progress(i+1, len(journal)) } } // Отмена могла прийти на последней доставке: без этой проверки запросы ниже // вернули бы context.Canceled как обычную ошибку, и команда завершилась бы // одной строкой «fatal», не напечатав частичного отчёта. if ctx.Err() != nil { rep.Canceled = true return rep, nil } rep.Buckets, err = o.Target.CountBuckets(ctx) if err != nil { return stopOr(rep, err) } rep.Workouts, err = o.Target.CountWorkouts(ctx) if err != nil { return stopOr(rep, err) } rep.Records, err = o.Target.CountRecords(ctx) if err != nil { return stopOr(rep, err) } rep.Categories, err = o.Target.CountCategoryValues(ctx) if err != nil { return stopOr(rep, err) } rep.Fingerprint, err = o.Target.Fingerprint(ctx) if err != nil { return stopOr(rep, err) } log.InfoContext(ctx, "journal replayed", "bodies", rep.Bodies, "skipped_files", rep.SkippedFiles, "adopted", rep.Adopted, "adopt_failed", rep.AdoptFailed, "orphans", rep.Orphans, "duplicates", rep.Duplicates, "folded", rep.Folded, "failed_layer", rep.FailedLayer, "failed_malformed", rep.FailedMalformed, "deferred", rep.Deferred, "failed_other", rep.FailedOther, "partial", rep.Partial, "incomparable", rep.Incomparable, "entities_held", rep.EntitiesHeld, "entities_diverging", rep.EntitiesDiverging, "buckets", rep.Buckets, "workouts", rep.Workouts, "records", rep.Records, "category_values", rep.Categories) return rep, nil } // collect собирает состав журнала и упорядочивает его. // // Второй возврат — учётные записи, тел которых в архиве нет. Они переносятся в // базу назначения, но не сворачиваются. func (o Options) collect(ctx context.Context, rep Report, log *slog.Logger) ([]store.Delivery, []store.Delivery, Report, error) { listing, err := o.Archive.List() if err != nil { // Нечитаемый каталог означает «неизвестно, есть ли тела», а не «тел // нет»: пустой журнал дал бы пустую витрину, чей отпечаток совпадает с // отпечатком любой другой пустой витрины. return nil, nil, rep, err } rep.SkippedFiles = len(listing.Skipped) for _, rel := range listing.Skipped { // Путь в архиве — дата и ULID, содержимого тела в нём нет. Без этой // строки счётчик пропусков в отчёте не на что раскрыть: имя файла не // узнать иначе как обходом архива руками. log.DebugContext(ctx, "archive file is not a body", "file", rel) } var known []store.Delivery if o.Source != nil { known, err = o.Source.ListDeliveries(ctx) if err != nil { return nil, nil, rep, err } } byID := make(map[string]store.Delivery, len(known)) for _, d := range known { byID[d.ID] = d } journal := make([]store.Delivery, 0, len(listing.Bodies)) seen := make(map[string]struct{}, len(listing.Bodies)) for _, rel := range listing.Bodies { // Имя обязано быть КАНОНИЧЕСКИМ идентификатором, а не приводиться к // нему: ident.Parse — функция входной границы, она обрезает пробелы и // поднимает регистр, и `\u00a0.json.gz` дал бы ту же координату // журнала, что настоящее тело. Подложенный файл сортируется раньше и // вытеснил бы настоящее — а невидимые пробелы в этих данных уже // встречались. id, err := ident.Parse(archive.BodyID(rel)) if err != nil || archive.BodyID(rel) != id { rep.SkippedFiles++ log.DebugContext(ctx, "archive file is not a body", "file", rel) continue } if _, dup := seen[id]; dup { // Одно имя в двух каталогах: копия, восстановленная руками, или // тело, переложенное не туда. Вторая запись журнала с тем же // идентификатором сорвала бы весь прогон отказом по первичному // ключу — то есть один посторонний файл лишал бы пересборки всё // остальное. rep.Duplicates++ log.WarnContext(ctx, "archive body id repeats", "file", rel, "delivery_id", id) continue } rep.Bodies++ seen[id] = struct{}{} if d, ok := byID[id]; ok { journal = append(journal, prepare(d)) continue } d, err := o.adopt(id, rel) if err != nil { // Тело есть, а завести по нему запись не вышло: без этой строки // счётчик в отчёте не на что раскрыть, а отчёт при этом отсылает // человека «разобраться по логу». rep.AdoptFailed++ log.WarnContext(ctx, "archive body not adopted", "error", err, "file", rel, "delivery_id", id) continue } rep.Adopted++ journal = append(journal, d) } // Учётная запись, тела которой в архиве нет. Переносится, но не // сворачивается: не перенести значило бы стереть первой же подменой базы // единственное свидетельство, что доставка была, — тела-то уже нет. orphans := make([]store.Delivery, 0) for _, d := range known { if _, ok := seen[d.ID]; !ok { rep.Orphans++ // Не `pending`: тела нет и не будет, а «этим разбором ещё не // смотрели» обещало бы данные, которых не появится, и ретеншен, // который pending не трогает никогда, берёг бы такие строки вечно. o := prepare(d) o.ParseStatus = store.ParseFailed orphans = append(orphans, o) log.DebugContext(ctx, "delivery body missing", "delivery_id", d.ID, "file", d.RawPath) } } // Тотальный ключ: `received_at` хранится с секундной точностью, и доставки // одной секунды без второго ключа шли бы в неопределённом порядке — два // прогона одного журнала могли бы разойтись. sort.Slice(journal, func(i, j int) bool { a, b := journal[i], journal[j] if !a.ReceivedAt.Equal(b.ReceivedAt) { return a.ReceivedAt.Before(b.ReceivedAt) } return a.ID < b.ID }) return journal, orphans, rep, nil } // stopOr отличает отмену от настоящего отказа: первая не ошибка операции, а // требование прекратить работу, и застать она может любой шаг. // // Единая точка на весь пакет: разбросанные по шагам проверки означали бы, что // каждый новый шаг обязан вспомнить правило руками, а забытый превращает // Ctrl-C в «ошибку окружения» без частичного отчёта. func stopOr(rep Report, err error) (Report, error) { if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { rep.Canceled = true return rep, nil } return rep, err } // prepare оставляет от учётной записи ФАКТЫ ЖУРНАЛА и сбрасывает производные от // разбора поля. // // `parse_status`, `points`, `derived_layer`, `uncovered_sections` и // `skipped_entities` — результат ПРЕДЫДУЩЕЙ свёртки, а не то, что приехало // вместе с доставкой. У последнего пустота означает «не измерялось», так что // перенос выдал бы измерение прежнего разбора за измерение текущего — а по нему // решают, можно ли удалить тело. Перенести их // значило бы сделать пересобранную витрину функцией прошлого прогона: доставка, // чей повторный разбор отказал (штатный исход, когда слой не выводится), // сохранила бы слой прежнего разбора — свёртка не затирает его намеренно, — и // следующая доставка той же автоматизации унаследовала бы его молча. Оба прогона // при этом самосогласованы, поэтому проверка «повторная пересборка ничего не // меняет» такого не ловит. // // Незаполненные здесь колонки получают значения по умолчанию схемы: пустой слой // и пустой список непокрытых секций. func prepare(d store.Delivery) store.Delivery { return store.Delivery{ ID: d.ID, ReceivedAt: d.ReceivedAt, AutomationName: d.AutomationName, AutomationID: d.AutomationID, Aggregation: d.Aggregation, Period: d.Period, SessionID: d.SessionID, Bytes: d.Bytes, SHA256: d.SHA256, RawPath: d.RawPath, Headers: d.Headers, ParseStatus: store.ParsePending, } } // adopt заводит учётную запись для тела, лежащего в архиве без неё. // // Такое тело — не экзотика: приём кладёт тело на диск раньше строки в базе // (обратный порядок дал бы учтённую доставку без данных), и отказ на вставке // оставляет тело без учёта. Обещание подобрать его записано в пакете приёма. // // Восстанавливается ровно то, что выводится из самого тела и его имени. // Заголовков в архиве нет вовсе, поэтому вывод слоя у такой доставки честно // деградирует: наследовать не от чего и подтверждать нечем. func (o Options) adopt(id, rel string) (store.Delivery, error) { at, err := ident.TimeOf(id) if err != nil { return store.Delivery{}, err } body, err := o.Fold.ReadBody(rel) if err != nil { return store.Delivery{}, err } sum := sha256.Sum256(body) return store.Delivery{ ID: id, ReceivedAt: at, // Размер и хеш — по РАСПАКОВАННОМУ телу, как их считает приём. Иначе в // тех же колонках появились бы значения другой природы, и индекс по // хешу начал бы врать на границе подобранных тел. Bytes: int64(len(body)), SHA256: hex.EncodeToString(sum[:]), RawPath: rel, ParseStatus: store.ParsePending, }, nil }