// Package hae — разбор тела доставки Health Auto Export в точки. // // Отдельно от хранения потому, что у разбора будет второй потребитель // (пересборка витрины из архива) и второй источник (родной экспорт Apple — // свой формат поверх того же хранилища). Пакет ничего не знает ни про SQLite, // ни про архив, и ничего не пишет: `Parse` — чистая функция от тела и // заголовков. // // Правила разбора выведены измерением живого потока, а не спроектированы: // docs/local-research.md, находки 2, 30, 33, 35, 36, 38, 39, 47. Документация // HAE местами расходится с тем, что приложение шлёт на самом деле, поэтому // источник истины по формату — пакеты в testdata. package hae import ( "bytes" "encoding/json" "errors" "fmt" "sort" "time" "unicode/utf8" ) // Layer — подробность, в которой метрика приехала. Выводится из выравнивания // меток времени, а не из заголовка доставки: заголовок `Default` наблюдался // одновременно у посекундного, минутного и часового режимов (находка 33). type Layer string // Слои хранения. `sample` появится с импортом родного экспорта Apple, `day` // назначается схемам с фиксированной гранулярностью и не выводится. const ( LayerSample Layer = "sample" LayerRaw Layer = "raw" LayerMinute Layer = "minute" LayerHour Layer = "hour" LayerDay Layer = "day" ) // Ошибки разбора. var ( // ErrMalformed — тело не разбирается как JSON ожидаемой формы. ErrMalformed = errors.New("тело не разбирается") // ErrLayerUnknown — в доставке есть метрики, но определить их слой нечем: // плотных метрик нет, у автоматизации нет предыдущего надёжного слоя, а // заголовок ненадёжен. Точки не сохраняются, тело остаётся в архиве — // доставку подберёт пересборка, когда слой станет известен. // // Молчаливый выбор `raw` здесь недопустим: призрачный разрез поедет в // каталог и в правило Read API «самый мелкий слой, покрывающий диапазон». ErrLayerUnknown = errors.New("слой доставки не определяется") ) // Порог плотности: метрика, у которой не меньше стольких точек с метками, // классифицируется по собственному выравниванию. Редкая наследует слой // доставки — её собственное выравнивание ничего не значит, потому что одна // метка на часе бывает и у минутного ряда. const denseThreshold = 10 // timeLayout — формат метки точки в секции metrics. Другого там не // встречается: RFC 3339 живёт только в stateOfMind, а Unix-эпоха — внутри // heartbeatSeries, и меткой точки не является. const timeLayout = "2006-01-02 15:04:05 -0700" // Point — одна разобранная точка. type Point struct { // Metric — имя метрики как прислал HAE. Исключение — sleep_analysis: под // одним именем приезжают две несовместимые схемы, и они разводятся. Metric string Units string Layer Layer // Start и End — координаты точки в UTC. У точки-измерения End равен Start: // ключ одной формы для всех точек, потому что интервальная и точечная // формы не встречаются вперемешку внутри метрики одной доставки // (находка 47). Start time.Time End time.Time // OffsetSeconds — смещение исходной зоны. Нормализовать время без него // значит потерять, в каком часовом поясе человек находился. OffsetSeconds int // Raw — содержимое точки исходными байтами, как пришло в теле. Пересборка // повторной сериализацией теряет литерал (`1.0` → `1`, целые больше 2^53 // сдвигаются, невалидный UTF-8 → U+FFFD), и потеря не видна тестам на // фикстурах: они сравнивают разобранное с разобранным. Raw json.RawMessage // local — метка в исходной зоне. Наружу не отдаётся: хранится всегда UTC // плюс офсет. Нужна выводу слоя — HAE строит сетку по МЕСТНОМУ времени, и // выравнивание, посчитанное по UTC, объявляет часовую выгрузку минутной в // зонах с получасовым смещением (+0530, +0545, +0930). local time.Time } // Result — итог разбора доставки. Частичные исходы живут в счётчиках, а не в // ошибке: пакет, у которого не разобралась одна точка из тысячи, — обычное // дело, и терять из-за неё остальное нельзя. type Result struct { Points []Point // Uncovered — верхнеуровневые ключи `data`, которых разбор не покрывает, // отсортированные и без повторов. Половина живого потока состоит из таких // доставок целиком (48 из 99: workouts и stateOfMind), и без этого списка // они неотличимы от разобранной доставки с пустой секцией метрик. // // Список канонизирован потому, что уезжает в базу и сравнивается между // доставками, а порядок ключей в JSON от HAE нестабилен. Uncovered []string // UncoveredDropped — сколько имён отброшено границей списка. Молчаливое // усечение сделало бы список уверенным, но неполным ответом на вопрос «что // останется потерянным, если тело удалить». UncoveredDropped int // Metrics — сколько метрик встретилось в секции. Metrics int // SkippedNoTime — точки без разбираемой метки времени. SkippedNoTime int // SkippedMalformed — точки, не разобравшиеся как объект JSON. SkippedMalformed int // SkippedBadEnd — точки, у которых есть `end`, но он не разбирается. // Вырождать такую точку в мгновенную нельзя: она схлопнулась бы с соседней // по координате. SkippedBadEnd int // Layer — слой, выведенный для доставки в целом (тот, что наследуют редкие // метрики). Пустой, если плотных метрик не было и наследовать было нечего. // Его сохраняет вызывающий, чтобы следующая доставка той же автоматизации // могла его унаследовать. Layer Layer // LayerMismatch — выведенный слой разошёлся с НАДЁЖНЫМ заголовком. // Заголовок `Default` в сравнении не участвует: он не означает режима, и // сравнение с ним давало бы WARN на каждой доставке потока в пять минут. LayerMismatch bool // HeaderLayer — слой по заголовку, если заголовок надёжен. HeaderLayer Layer } // Meta — что доставка рассказала о себе заголовками, плюс память о прошлых // доставках той же автоматизации. type Meta struct { // Aggregation — заголовок `automation-aggregation`. Надёжен только в // значениях `Minutes` и `Hours`. Aggregation string // FallbackLayer — последний надёжно выведенный слой этой же автоматизации // (`automation-id`). Нужен доставкам без плотных метрик: измерено 2 такие // из 89, обе с заголовком `Default`. Ищет и передаёт его вызывающий — // разбор остаётся чистой функцией. FallbackLayer Layer } // Parse разбирает секцию metrics тела доставки в точки. // // Ошибка возвращается только когда точек не будет вовсе: тело не JSON // (ErrMalformed) или слой не определяется (ErrLayerUnknown). Всё остальное — // счётчики в Result. Отсутствие секции metrics ошибкой не является: доставки // с одними тренировками или состоянием разума — норма. func Parse(body []byte, meta Meta) (res Result, err error) { // Разбор чужого формата обязан отвечать ошибкой, а не паникой: приём не // имеет права упасть из-за того, что HAE прислал невиданное. Перехват // стоит здесь, внутри разбора, а не выше: паника из хранилища — настоящий // дефект, и глушить её нельзя. defer func() { if r := recover(); r != nil { res = Result{} err = fmt.Errorf("%w: паника разбора: %v", ErrMalformed, r) } }() metrics, uncovered, dropped, err := decodeEnvelope(body) if err != nil { return Result{}, err } res.Uncovered = uncovered res.UncoveredDropped = dropped res.Metrics = len(metrics) if len(metrics) == 0 { return res, nil } groups := make([]group, 0, len(metrics)) for _, m := range metrics { g := decodeGroup(m, &res) if len(g.summaries) > 0 { groups = append(groups, group{ metric: sleepSummaryMetric, units: g.units, points: g.summaries, layer: LayerDay, fixed: true, }) g.summaries = nil } if len(g.points) > 0 { groups = append(groups, g) } } if len(groups) == 0 { return res, nil } res.HeaderLayer = headerLayer(meta.Aggregation) if err := assignLayers(groups, meta, &res); err != nil { // Список непокрытых секций переживает отказ: доставка, у которой не // определился слой, обязана остаться записью о том, что в теле есть // невосстановимая секция. Иначе ретеншен увидит failed без списка и // решит, что терять нечего. return Result{ Metrics: res.Metrics, Uncovered: res.Uncovered, UncoveredDropped: res.UncoveredDropped, }, err } total := 0 for _, g := range groups { total += len(g.points) } res.Points = make([]Point, 0, total) for _, g := range groups { for _, p := range g.points { p.Layer = g.layer res.Points = append(res.Points, p) } } return res, nil } // group — точки одной метрики одной доставки: единица, для которой выводится // слой. Классификация именно по метрике, а не по доставке: при перенастройке // автоматизации приезжают смешанные доставки, и отнесение такой доставки к // одному слою складывает минутные точки с посекундными (находка 33). type group struct { metric string units string points []Point // layer — итоговый слой группы. layer Layer // fixed — слой назначен схемой, а не выведен (суточная сводка сна). Такие // группы не участвуют в определении слоя доставки: сводок бывает больше // порога плотности, и их полуночные метки назначили бы всей доставке hour. fixed bool // alignment — самое мелкое выравнивание среди меток группы. alignment Layer dense bool // summaries — точки, схема которых опознана прямо при разборе и слой // которым назначен, а не выведен (суточная сводка сна). Держатся отдельно, // потому что в определении слоя доставки не участвуют: сводок бывает больше // порога плотности, и их полуночные метки назначили бы всей доставке `hour`. summaries []Point } // envelope — форма тела, ровно настолько подробная, насколько нужно разбору. // // Точки держатся сырыми сообщениями и декодируются по одной: разбор тела в // 42 МиБ через map[string]any удерживает 197 МиБ кучи против 54 МиБ у этой // формы. Вместе с самим телом пик доходил бы до ~300 МиБ на доставку — это // OOM ровно на пике потока, когда терять доставки дороже всего. // metricsSection — единственная секция, которую разбор покрывает сегодня. const metricsSection = "metrics" // covered говорит, покрывает ли разбор секцию с таким именем. // // Функция, а не изменяемая карта: разбор и перечисление непокрытых ходят по // одному источнику, поэтому состояние «секция разбирается, но числится // непокрытой» невыразимо. func covered(section string) bool { return section == metricsSection } // Границы на список непокрытых ключей. Тело контролирует отправитель целиком: // без границ сто тысяч однобуквенных ключей превращаются в одну строку в базе // и одну строку в логе того же порядка. Секций у HAE восемь, самое длинное имя // — heartRateNotifications (22 байта), так что запас велик. const ( maxUncovered = 32 maxUncoveredLen = 64 ) type metricEnvelope struct { Name string `json:"name"` Units string `json:"units"` Data []json.RawMessage `json:"data"` } // pointHead — поля точки, нужные разбору. Всё остальное остаётся в Raw и // хранится дословно: «служебных» полей у точки нет, отбрасывать нечего. type pointHead struct { Date string `json:"date"` Start string `json:"start"` End string `json:"end"` // TotalSleep различает две схемы под именем sleep_analysis: поэпизодную и // суточную сводку. Общих полей, кроме date и source, у них нет. TotalSleep *json.RawMessage `json:"totalSleep"` } // decodeEnvelope разбирает конверт: отдаёт секцию metrics и имена секций, // которых разбор не покрывает. // // Идёт по верхнему уровню одним декодером: Token() читает рамку объекта и имена // членов, Decode() — значения. Значение покрытого ключа декодируется на месте, // значение непокрытого ПРОГЛАТЫВАЕТСЯ декодированием в выбрасываемый // RawMessage. Это форма из ExampleDecoder_Decode_stream стандартной библиотеки; // в encoding/json/v2 та же операция названа прямо — SkipValue. // // Пропуск ручным счётом глубины по Token() выглядит дешевле и измеримо хуже: // делимитеры идут мимо сканера, поэтому ограничитель вложенности encoding/json // не работает, а стек токенов растёт как O(глубины). Тело 40 МиБ из вложенных // скобок даёт пик 488 МиБ вместо контрактных четырёх тел — Decode отвергает его // мгновенно. Разбор `data` в map[string]json.RawMessage дешевле по коду, но // копирует байты ВСЕХ секций и держит их до конца разбора; у проглатывания // копия одна и живёт до следующего члена. func decodeEnvelope(body []byte) (metrics []metricEnvelope, uncovered []string, dropped int, err error) { fail := func(e error) ([]metricEnvelope, []string, int, error) { return nil, nil, 0, fmt.Errorf("%w: %v", ErrMalformed, e) //nolint:errorlint // причина уходит в лог, наружу не раскрывается } dec := json.NewDecoder(bytes.NewReader(body)) // Верхний уровень тела: интересует только data. Прочие ключи конверта в // список не идут — иначе в одном списке смешались бы имена секций и мусор // конверта, а форму `{"data": …}` проверяет приём. tok, err := dec.Token() if err != nil { return fail(err) } // Голый null телом ошибкой не был и не становится: прежний разбор // раскладывал его в пустую структуру. Границы поведения этой задачей не // двигаются — она добавляет список, а не строгость. if tok == nil { return nil, nil, 0, nil } if d, ok := tok.(json.Delim); !ok || d != '{' { return fail(fmt.Errorf("ожидался объект, встречено %v", tok)) } seen := make(map[string]struct{}) for dec.More() { name, err := memberName(dec) if err != nil { return fail(err) } if name != "data" { if err := swallow(dec); err != nil { return fail(err) } continue } metrics, uncovered, dropped, err = decodeData(dec, seen) if err != nil { return fail(err) } } if err := expectDelim(dec, '}'); err != nil { return fail(err) } // Список канонизируется: порядок ключей в JSON от HAE нестабилен, а // значение уезжает в базу и сравнивается между доставками. sort.Strings(uncovered) return metrics, uncovered, dropped, nil } // decodeData разбирает объект data, собирая metrics и имена непокрытых секций. func decodeData(dec *json.Decoder, seen map[string]struct{}) ([]metricEnvelope, []string, int, error) { tok, err := dec.Token() if err != nil { return nil, nil, 0, err } // data не объект — прежнее поведение: ошибка ровно там, где была. if d, ok := tok.(json.Delim); !ok || d != '{' { return nil, nil, 0, fmt.Errorf("data: ожидался объект, встречено %v", tok) } var ( metrics []metricEnvelope uncovered []string dropped int ) for dec.More() { name, err := memberName(dec) if err != nil { return nil, nil, 0, err } if covered(name) { // Повтор ключа metrics JSON допускает; секции ОБЪЕДИНЯЮТСЯ, а не // побеждает последняя: терять точки молча нельзя. var part []metricEnvelope if err := dec.Decode(&part); err != nil { return nil, nil, 0, err } metrics = append(metrics, part...) continue } if err := swallow(dec); err != nil { return nil, nil, 0, err } if _, dup := seen[name]; dup { continue } seen[name] = struct{}{} if len(uncovered) >= maxUncovered { dropped++ continue } uncovered = append(uncovered, clipSection(name)) } if _, err := dec.Token(); err != nil { // закрывающая скобка data return nil, nil, 0, err } return metrics, uncovered, dropped, nil } // memberName читает имя члена объекта. Token() отдаёт имя уже после разбора // escape-последовательностей, поэтому границы считаются по декодированному. func memberName(dec *json.Decoder) (string, error) { tok, err := dec.Token() if err != nil { return "", err } name, ok := tok.(string) if !ok { return "", fmt.Errorf("ожидалось имя члена, встречено %v", tok) } return name, nil } // swallow проглатывает значение целиком, ничего не удерживая. func swallow(dec *json.Decoder) error { var skip json.RawMessage return dec.Decode(&skip) } func expectDelim(dec *json.Decoder, want json.Delim) error { tok, err := dec.Token() if err != nil { return err } if d, ok := tok.(json.Delim); !ok || d != want { return fmt.Errorf("ожидалось %q, встречено %v", want, tok) } return nil } // clipSection обрезает слишком длинное имя по границе рун и помечает обрезку. // Маркер приписывается СВЕРХ предела: обрезка не инъективна, и обрезанное имя // сравнению со словарём известных секций не подлежит. func clipSection(name string) string { if len(name) <= maxUncoveredLen { return name } cut := maxUncoveredLen for cut > 0 && !utf8.RuneStart(name[cut]) { cut-- } return name[:cut] + "…" } // decodeGroup разбирает точки одной метрики. Точка без разбираемой метки // пропускается со счётчиком — ронять из-за неё остальную доставку незачем. // // Схема точки определяется ЗДЕСЬ ЖЕ, в том же проходе, где точка разобрана. // Второй проход по исходному массиву был бы неверен: список разобранных точек // уже отфильтрован пропусками, и любой пропуск сдвигал бы соответствие — эпизод // сна уезжал бы под имя суточной сводки, а сводка под имя эпизода. Индексной // корреляции между двумя списками здесь не существует по построению. func decodeGroup(m metricEnvelope, res *Result) group { g := group{metric: m.Name, units: m.Units, points: make([]Point, 0, len(m.Data))} for _, raw := range m.Data { var head pointHead if err := json.Unmarshal(raw, &head); err != nil { res.SkippedMalformed++ continue } start, ok := parseTime(head.Date, head.Start) if !ok { res.SkippedNoTime++ continue } end := start if head.End != "" { e, err := time.Parse(timeLayout, head.End) if err != nil { // Интервал, конец которого не читается, вырождать в точку // нельзя: две записи с общим началом получили бы одну // координату и одна из них исчезла бы. Пропускаем со // счётчиком — тело остаётся в архиве. res.SkippedBadEnd++ continue } end = e } _, offset := start.Zone() p := Point{ Metric: m.Name, Units: m.Units, Start: start.UTC(), End: end.UTC(), OffsetSeconds: offset, Raw: raw, local: start, } // Под именем sleep_analysis HAE шлёт две несовместимые схемы: // поэпизодную (start/end/value/qty) и суточную сводку // (totalSleep/core/rem/deep/awake с меткой на местной полуночи). Имя // sleep_analysis_summary — наше; инвариант «форма Apple не // транслируется» это не нарушает: поля внутри точки не // переименовываются, разделяются только имена метрик, под которыми // HAE смешал две схемы. if m.Name == sleepMetric && head.TotalSleep != nil { p.Metric = sleepSummaryMetric g.summaries = append(g.summaries, p) continue } g.points = append(g.points, p) } g.dense = len(g.points) >= denseThreshold g.alignment = finestAlignment(g.points) return g } // Имена метрик сна: пришедшее от HAE и наше для суточной сводки. const ( sleepMetric = "sleep_analysis" sleepSummaryMetric = "sleep_analysis_summary" ) // parseTime разбирает метку точки. Начало берётся из start, а при его // отсутствии — из date. Измерено: start, когда он есть, всегда совпадает с // date, поэтому правило не вводит второго источника метки — оно закрывает // случай, когда HAE перестанет их дублировать. func parseTime(date, start string) (time.Time, bool) { s := start if s == "" { s = date } if s == "" { return time.Time{}, false } t, err := time.Parse(timeLayout, s) if err != nil { return time.Time{}, false } return t, true } // finestAlignment возвращает самое мелкое выравнивание среди меток. // // Именно самое мелкое, а не преобладающее: у плотных метрик выравнивания // перемешаны (active_energy — 1320 минутных меток и 21 часовая), потому что // метка ровно на часе одновременно является и минутной. Метрика, у которой // хоть одна метка стоит на середине часа, часовой не является. func finestAlignment(points []Point) Layer { finest := LayerHour for _, p := range points { // По МЕСТНОЙ метке, а не по UTC: HAE строит сетку в зоне телефона. // В зонах с получасовым смещением (+0530 Индия, +0545 Непал, +0930 // Аделаида) ровный местный час превращается в UTC-метку на половине, // и вся часовая выгрузка уехала бы в слой `minute` — а там столкнулась // бы с настоящей минутной автоматизацией на координате hh:30 и завысила // сумму минутного слоя вдвое. switch { case p.local.Second() != 0 || p.local.Nanosecond() != 0: return LayerRaw case p.local.Minute() != 0: finest = LayerMinute } } return finest } // headerLayer переводит заголовок в слой, но только когда заголовок надёжен. // `Default` соответствует трём разным режимам выгрузки и не означает ничего. func headerLayer(aggregation string) Layer { switch aggregation { case "Minutes": return LayerMinute case "Hours": return LayerHour default: return "" } } // finer возвращает более мелкий из двух слоёв. func finer(a, b Layer) Layer { if rank(a) < rank(b) { return a } return b } func rank(l Layer) int { switch l { case LayerSample: return 0 case LayerRaw: return 1 case LayerMinute: return 2 case LayerHour: return 3 case LayerDay: return 4 default: return 5 } }