разбор метрик HAE: фикстуры, канонизация, парсер
- tmp/research/fixtures.py собирает фикстуры из архива, вычищая измерения и сохраняя порядок ключей, форму литералов, выравнивание меток и невидимые символы; шесть фикстур в internal/hae/testdata - internal/canon — общий дом канонической формы, полноты и хеша: числа читаются литералом через json.Number, округление до 12 значащих цифр - internal/hae — разбор секции metrics, вывод слоя по метрике, разделение схем сна, координаты интервалом; recover внутри Parse, фаззинг - миграция bucket и docs/database.md
This commit is contained in:
@@ -0,0 +1,429 @@
|
||||
// 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"
|
||||
"time"
|
||||
)
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
// Result — итог разбора доставки. Частичные исходы живут в счётчиках, а не в
|
||||
// ошибке: пакет, у которого не разобралась одна точка из тысячи, — обычное
|
||||
// дело, и терять из-за неё остальное нельзя.
|
||||
type Result struct {
|
||||
Points []Point
|
||||
|
||||
// Metrics — сколько метрик встретилось в секции.
|
||||
Metrics int
|
||||
// SkippedNoTime — точки без разбираемой метки времени.
|
||||
SkippedNoTime int
|
||||
// SkippedMalformed — точки, не разобравшиеся как объект JSON.
|
||||
SkippedMalformed 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, err := decodeMetrics(body)
|
||||
if err != nil {
|
||||
return Result{}, err
|
||||
}
|
||||
|
||||
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 g.spillover != nil {
|
||||
groups = append(groups, *g.spillover)
|
||||
g.spillover = 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 {
|
||||
return Result{Metrics: res.Metrics}, 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
|
||||
|
||||
// spillover — вторая группа, отделившаяся от этой при разделении схем под
|
||||
// одним именем метрики.
|
||||
spillover *group
|
||||
}
|
||||
|
||||
// envelope — форма тела, ровно настолько подробная, насколько нужно разбору.
|
||||
//
|
||||
// Точки держатся сырыми сообщениями и декодируются по одной: разбор тела в
|
||||
// 42 МиБ через map[string]any удерживает 197 МиБ кучи против 54 МиБ у этой
|
||||
// формы. Вместе с самим телом пик доходил бы до ~300 МиБ на доставку — это
|
||||
// OOM ровно на пике потока, когда терять доставки дороже всего.
|
||||
type envelope struct {
|
||||
Data struct {
|
||||
Metrics []metricEnvelope `json:"metrics"`
|
||||
} `json:"data"`
|
||||
}
|
||||
|
||||
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"`
|
||||
}
|
||||
|
||||
func decodeMetrics(body []byte) ([]metricEnvelope, error) {
|
||||
var env envelope
|
||||
dec := json.NewDecoder(bytes.NewReader(body))
|
||||
if err := dec.Decode(&env); err != nil {
|
||||
return nil, fmt.Errorf("%w: %v", ErrMalformed, err) //nolint:errorlint // причина уходит в лог, наружу не раскрывается
|
||||
}
|
||||
return env.Data.Metrics, nil
|
||||
}
|
||||
|
||||
// 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 != "" {
|
||||
if e, err := time.Parse(timeLayout, head.End); err == nil {
|
||||
end = e
|
||||
}
|
||||
}
|
||||
|
||||
_, offset := start.Zone()
|
||||
g.points = append(g.points, Point{
|
||||
Metric: m.Name,
|
||||
Units: m.Units,
|
||||
Start: start.UTC(),
|
||||
End: end.UTC(),
|
||||
OffsetSeconds: offset,
|
||||
Raw: raw,
|
||||
})
|
||||
}
|
||||
|
||||
if len(g.points) == 0 {
|
||||
return g
|
||||
}
|
||||
|
||||
// Разделение схем выполняется ДО вывода слоя: у суточной сводки слой
|
||||
// назначен, и её полуночные метки не должны участвовать в голосовании.
|
||||
if m.Name == "sleep_analysis" {
|
||||
splitSleep(&g, m.Data)
|
||||
}
|
||||
|
||||
g.dense = len(g.points) >= denseThreshold
|
||||
g.alignment = finestAlignment(g.points)
|
||||
return g
|
||||
}
|
||||
|
||||
// splitSleep разводит суточную сводку сна на собственное имя метрики.
|
||||
//
|
||||
// Под именем sleep_analysis HAE шлёт две несовместимые схемы: поэпизодную
|
||||
// (start/end/value/qty) и суточную сводку (totalSleep/core/rem/deep/awake с
|
||||
// меткой на местной полуночи). Имя sleep_analysis_summary — наше; инвариант
|
||||
// «форма Apple не транслируется» это не нарушает, потому что поля внутри точки
|
||||
// не переименовываются, разделяются только имена метрик, под которыми HAE
|
||||
// смешал две схемы.
|
||||
//
|
||||
// Если в группе оказались обе схемы, сводки переезжают в отдельную группу,
|
||||
// возвращаемую через g.spillover. Живьём смеси не наблюдалось, но полагаться
|
||||
// на это нельзя: HAE меняется между версиями приложения.
|
||||
func splitSleep(g *group, raws []json.RawMessage) {
|
||||
summaries := make([]Point, 0)
|
||||
episodes := make([]Point, 0, len(g.points))
|
||||
|
||||
for i, p := range g.points {
|
||||
if i < len(raws) && isSleepSummary(raws[i]) {
|
||||
p.Metric = sleepSummaryMetric
|
||||
summaries = append(summaries, p)
|
||||
continue
|
||||
}
|
||||
episodes = append(episodes, p)
|
||||
}
|
||||
|
||||
if len(summaries) == 0 {
|
||||
return
|
||||
}
|
||||
g.spillover = &group{
|
||||
metric: sleepSummaryMetric,
|
||||
units: g.units,
|
||||
points: summaries,
|
||||
layer: LayerDay,
|
||||
fixed: true,
|
||||
}
|
||||
g.points = episodes
|
||||
}
|
||||
|
||||
const sleepSummaryMetric = "sleep_analysis_summary"
|
||||
|
||||
func isSleepSummary(raw json.RawMessage) bool {
|
||||
var head pointHead
|
||||
if err := json.Unmarshal(raw, &head); err != nil {
|
||||
return false
|
||||
}
|
||||
return head.TotalSleep != nil
|
||||
}
|
||||
|
||||
// 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 {
|
||||
switch {
|
||||
case p.Start.Second() != 0 || p.Start.Nanosecond() != 0:
|
||||
return LayerRaw
|
||||
case p.Start.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
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user