непокрытые секции доставки видны в статусе разбора
- половина потока (50 доставок из 104) не несёт metrics вовсе и до сих пор числилась parsed: ретеншен, поверив статусу, срезал бы тела stateOfMind, которых в экспорте Apple нет - разбор перечисляет верхнеуровневые ключи data, непокрытые проглатываются декодированием: тело 40 МиБ из непокрытой секции удерживает 0 МиБ - статус partial и колонка delivery.uncovered_sections; миграция переводит прежние parsed в pending — им верить нельзя - витрина не изменилась: отпечаток совпал с прогоном до изменения
This commit is contained in:
+46
-8
@@ -82,6 +82,12 @@ type Stats struct {
|
||||
LayerMismatch bool
|
||||
Collisions []store.Collision
|
||||
IncomparableAt []store.Collision
|
||||
// Uncovered — верхнеуровневые ключи `data`, которых разбор не покрывает.
|
||||
// Это ответ на вопрос «что останется потерянным, если тело удалить»:
|
||||
// для stateOfMind он необратим — в экспорте Apple этой секции нет.
|
||||
Uncovered []string
|
||||
// UncoveredDropped — сколько имён отброшено границей списка.
|
||||
UncoveredDropped int
|
||||
}
|
||||
|
||||
// Fold разбирает тело доставки и раскладывает точки по часовым объектам.
|
||||
@@ -103,7 +109,7 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) {
|
||||
|
||||
body, err := s.readBody(d.RawPath)
|
||||
if err != nil {
|
||||
s.fail(ctx, deliveryID, err)
|
||||
s.fail(ctx, deliveryID, err, nil)
|
||||
return stats, err
|
||||
}
|
||||
|
||||
@@ -118,10 +124,15 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) {
|
||||
FallbackLayer: hae.Layer(fallback),
|
||||
})
|
||||
if err != nil {
|
||||
s.fail(ctx, deliveryID, err)
|
||||
// Список непокрытых секций переживает отказ: доставка, у которой не
|
||||
// определился слой, обязана остаться записью о том, что в теле есть
|
||||
// невосстановимая секция.
|
||||
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
|
||||
@@ -132,7 +143,7 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) {
|
||||
|
||||
merge, err := s.store.MergePoints(ctx, toIncoming(parsed.Points), deliveryID)
|
||||
if err != nil {
|
||||
s.fail(ctx, deliveryID, err)
|
||||
s.fail(ctx, deliveryID, err, parsed.Uncovered)
|
||||
return stats, err
|
||||
}
|
||||
|
||||
@@ -146,7 +157,20 @@ func (s *Service) Fold(ctx context.Context, deliveryID string) (Stats, error) {
|
||||
stats.Collisions = merge.Collisions
|
||||
stats.IncomparableAt = merge.IncomparableAt
|
||||
|
||||
if err := s.finish(ctx, deliveryID, store.ParseDone, int64(stats.Points), stats.Layer); err != nil {
|
||||
// Источник истины — список; статус производен от него и от факта отказа.
|
||||
// Приоритет назван явно, иначе два будущих читателя (ретеншен и /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
|
||||
}
|
||||
@@ -186,6 +210,11 @@ func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) {
|
||||
"skipped_bad_end", st.SkippedBadEnd,
|
||||
"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))
|
||||
@@ -200,6 +229,10 @@ func (s *Service) logResult(ctx context.Context, deliveryID string, st Stats) {
|
||||
allSkipped := st.Points == 0 && skipped > 0
|
||||
|
||||
switch {
|
||||
case st.UncoveredDropped > 0:
|
||||
// Не частичный разбор, а тело, не похожее на HAE: секций у HAE восемь,
|
||||
// а границу выбило больше тридцати двух.
|
||||
s.log.WarnContext(ctx, "delivery folded, uncovered section list truncated", attrs...)
|
||||
case st.Incomparable > 0:
|
||||
// Выше перезаписей намеренно: несравнимый набор полей — событие реже и
|
||||
// информативнее, на живом потоке не случавшееся ни разу. Признаки при
|
||||
@@ -239,11 +272,11 @@ func formatCollisions(cs []store.Collision) string {
|
||||
const keepLayer = ""
|
||||
|
||||
// finish записывает исход разбора на контексте, переживающем отмену исходного.
|
||||
func (s *Service) finish(ctx context.Context, deliveryID, status string, points int64, layer string) error {
|
||||
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, status, points, layer); err != nil {
|
||||
if err := s.store.FinishParse(ctx, deliveryID, out); err != nil {
|
||||
return fmt.Errorf("запись исхода разбора: %w", err)
|
||||
}
|
||||
return nil
|
||||
@@ -251,7 +284,7 @@ func (s *Service) finish(ctx context.Context, deliveryID, status string, points
|
||||
|
||||
// fail отмечает доставку неразобранной. Тело остаётся в архиве, и её подберёт
|
||||
// пересборка — приём при этом не затрагивается: сохранили значит приняли.
|
||||
func (s *Service) fail(ctx context.Context, deliveryID string, cause error) {
|
||||
func (s *Service) fail(ctx context.Context, deliveryID string, cause error, uncovered []string) {
|
||||
level := slog.LevelError
|
||||
switch {
|
||||
case errors.Is(cause, hae.ErrLayerUnknown):
|
||||
@@ -269,7 +302,12 @@ func (s *Service) fail(ctx context.Context, deliveryID string, cause error) {
|
||||
// Слой НЕ затирается: доставка могла свернуться успешно раньше, и пустая
|
||||
// строка здесь оборвала бы цепочку наследования, то есть изменила бы
|
||||
// результат пересборки журнала.
|
||||
if err := s.finish(ctx, deliveryID, store.ParseFailed, 0, keepLayer); err != nil {
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
"log/slog"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -228,3 +229,98 @@ func mustHour(t *testing.T, st *store.Store, metric, layer string) time.Time {
|
||||
}
|
||||
return hours[0]
|
||||
}
|
||||
|
||||
// Доставка с непокрытой секцией обязана быть ОТЛИЧИМА от разобранной целиком.
|
||||
// Без этого ретеншен, ориентируясь на статус, срежет тело — а для stateOfMind
|
||||
// это необратимо: в экспорте Apple секции нет, доставки HAE единственный
|
||||
// источник.
|
||||
func TestFoldЧастичныйРазборВиденВУчёте(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
f, arch, st := newFold(t)
|
||||
ctx := context.Background()
|
||||
|
||||
deliver(t, arch, st, "d1", "Minutes", "a1", fixture(t, "uncovered_sections.json"))
|
||||
stats, err := f.Fold(ctx, "d1")
|
||||
if err != nil {
|
||||
t.Fatalf("свёртка: %v", err)
|
||||
}
|
||||
|
||||
if len(stats.Uncovered) == 0 {
|
||||
t.Error("список непокрытых секций пуст")
|
||||
}
|
||||
|
||||
d, err := st.LastDelivery(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("чтение доставки: %v", err)
|
||||
}
|
||||
if d.ParseStatus != store.ParsePartial {
|
||||
t.Errorf("статус %q, ожидался %q", d.ParseStatus, store.ParsePartial)
|
||||
}
|
||||
if !strings.Contains(d.UncoveredSections, "stateOfMind") {
|
||||
t.Errorf("список в базе %q не содержит stateOfMind", d.UncoveredSections)
|
||||
}
|
||||
// Точки метрик обязаны сохраниться: частичность не отменяет разобранного.
|
||||
if stats.Points == 0 {
|
||||
t.Error("точек нет — покрытая секция потерялась вместе с непокрытой")
|
||||
}
|
||||
}
|
||||
|
||||
// Доставка из одних метрик списка не получает и остаётся parsed: частичность
|
||||
// производна от списка, а не назначена.
|
||||
func TestFoldПолныйРазборОстаётсяParsed(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
f, arch, st := newFold(t)
|
||||
ctx := context.Background()
|
||||
|
||||
deliver(t, arch, st, "d1", "Minutes", "a1", fixture(t, "minute.json"))
|
||||
if _, err := f.Fold(ctx, "d1"); err != nil {
|
||||
t.Fatalf("свёртка: %v", err)
|
||||
}
|
||||
|
||||
d, err := st.LastDelivery(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("чтение доставки: %v", err)
|
||||
}
|
||||
if d.ParseStatus != store.ParseDone {
|
||||
t.Errorf("статус %q, ожидался %q", d.ParseStatus, store.ParseDone)
|
||||
}
|
||||
if d.UncoveredSections != "[]" {
|
||||
t.Errorf("список %q, ожидался `[]` — ровно одно представление пустоты", d.UncoveredSections)
|
||||
}
|
||||
}
|
||||
|
||||
// Список замещает прежнее значение ЦЕЛИКОМ, включая замещение пустым. Иначе
|
||||
// доставка, чья секция стала покрытой, осталась бы partial навсегда, и
|
||||
// ретеншен вечно щадил бы тело, которое уже не нужно.
|
||||
func TestFoldПересвёрткаОчищаетСписок(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
f, arch, st := newFold(t)
|
||||
ctx := context.Background()
|
||||
|
||||
// Сперва доставка с непокрытой секцией.
|
||||
deliver(t, arch, st, "d1", "Minutes", "a1", fixture(t, "uncovered_sections.json"))
|
||||
if _, err := f.Fold(ctx, "d1"); err != nil {
|
||||
t.Fatalf("свёртка: %v", err)
|
||||
}
|
||||
|
||||
// Затем — та же доставка, но тело уже без непокрытых секций: так выглядит
|
||||
// пересвёртка после того, как секцию научились разбирать.
|
||||
deliver(t, arch, st, "d2", "Minutes", "a1", fixture(t, "minute.json"))
|
||||
if _, err := f.Fold(ctx, "d2"); err != nil {
|
||||
t.Fatalf("свёртка: %v", err)
|
||||
}
|
||||
|
||||
d, err := st.LastDelivery(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("чтение доставки: %v", err)
|
||||
}
|
||||
if d.UncoveredSections != "[]" {
|
||||
t.Errorf("список %q не очистился — пустой список обязан замещать прежний", d.UncoveredSections)
|
||||
}
|
||||
if d.ParseStatus != store.ParseDone {
|
||||
t.Errorf("статус %q, ожидался %q", d.ParseStatus, store.ParseDone)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -182,3 +182,75 @@ func TestFoldНесравнимыеНаборыДаютWarn(t *testing.T) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Имена непокрытых секций в логе нужны — по ним видно, что поток принёс новое.
|
||||
// Содержимого секций там быть не может: это данные о здоровье. И уровень от
|
||||
// самой частичности не растёт — `partial` установившееся состояние половины
|
||||
// потока, а постоянный WARN каждые пять минут обесценивает уровень.
|
||||
func TestFoldЧастичныйРазборВЛоге(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var buf bytes.Buffer
|
||||
log := slog.New(slog.NewJSONHandler(&buf, &slog.HandlerOptions{Level: slog.LevelInfo}))
|
||||
|
||||
dir := t.TempDir()
|
||||
arch, err := archive.New(filepath.Join(dir, "raw"))
|
||||
if err != nil {
|
||||
t.Fatalf("архив: %v", err)
|
||||
}
|
||||
st, err := store.Open(filepath.Join(dir, "healthlog.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("база: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
|
||||
f := fold.New(arch, st, 0, log)
|
||||
ctx := context.Background()
|
||||
|
||||
// Содержимое непокрытой секции помечено так, чтобы его нельзя было спутать
|
||||
// ни с чем: если оно окажется в логе, тест это увидит.
|
||||
const body = `{"data":{
|
||||
"metrics":[{"name":"step_count","units":"count","data":[
|
||||
{"date":"2025-06-05 10:00:00 +0300","qty":1},
|
||||
{"date":"2025-06-05 10:01:00 +0300","qty":2},
|
||||
{"date":"2025-06-05 10:02:00 +0300","qty":3},
|
||||
{"date":"2025-06-05 10:03:00 +0300","qty":4},
|
||||
{"date":"2025-06-05 10:04:00 +0300","qty":5},
|
||||
{"date":"2025-06-05 10:05:00 +0300","qty":6},
|
||||
{"date":"2025-06-05 10:06:00 +0300","qty":7},
|
||||
{"date":"2025-06-05 10:07:00 +0300","qty":8},
|
||||
{"date":"2025-06-05 10:08:00 +0300","qty":9},
|
||||
{"date":"2025-06-05 10:09:00 +0300","qty":10}]}],
|
||||
"stateOfMind":[{"valence":"СЕКРЕТНОЕ-НАСТРОЕНИЕ"}]}}`
|
||||
|
||||
deliver(t, arch, st, "d1", "Minutes", "a1", []byte(body))
|
||||
if _, err := f.Fold(ctx, "d1"); err != nil {
|
||||
t.Fatalf("свёртка: %v", err)
|
||||
}
|
||||
|
||||
out := buf.String()
|
||||
if !strings.Contains(out, "stateOfMind") {
|
||||
t.Error("имени непокрытой секции нет в логе — момент появления новой секции незаметен")
|
||||
}
|
||||
if strings.Contains(out, "СЕКРЕТНОЕ-НАСТРОЕНИЕ") {
|
||||
t.Error("содержимое непокрытой секции утекло в лог")
|
||||
}
|
||||
|
||||
var rec struct {
|
||||
Level string `json:"level"`
|
||||
Uncovered []string `json:"uncovered"`
|
||||
}
|
||||
line := strings.TrimSpace(out)
|
||||
if i := strings.LastIndex(line, "\n"); i >= 0 {
|
||||
line = line[i+1:]
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &rec); err != nil {
|
||||
t.Fatalf("запись лога не разбирается: %v", err)
|
||||
}
|
||||
if rec.Level != "INFO" {
|
||||
t.Errorf("уровень %q, ожидался INFO: частичность — не отклонение", rec.Level)
|
||||
}
|
||||
if len(rec.Uncovered) != 1 || rec.Uncovered[0] != "stateOfMind" {
|
||||
t.Errorf("атрибут uncovered = %v, ожидался структурный список из stateOfMind", rec.Uncovered)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -39,6 +39,8 @@ func TestReplayЖивогоАрхива(t *testing.T) {
|
||||
|
||||
f, arch, st := newFold(t)
|
||||
ctx := context.Background()
|
||||
partial := 0
|
||||
sections := map[string]int{}
|
||||
|
||||
var folded, failed, incomparable int
|
||||
for _, path := range bodies {
|
||||
@@ -62,6 +64,12 @@ func TestReplayЖивогоАрхива(t *testing.T) {
|
||||
}
|
||||
folded++
|
||||
incomparable += res.Incomparable
|
||||
if len(res.Uncovered) > 0 {
|
||||
partial++
|
||||
for _, s := range res.Uncovered {
|
||||
sections[s]++
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Несравнимые наборы полей — посылка, на которой стоит отказ от объединения
|
||||
@@ -70,11 +78,19 @@ func TestReplayЖивогоАрхива(t *testing.T) {
|
||||
// сходимости.
|
||||
t.Logf("доставок %d: свёрнуто %d, не свёрнуто %d, несравнимых наборов %d",
|
||||
len(bodies), folded, failed, incomparable)
|
||||
t.Logf("частично разобрано %d, непокрытые секции: %v", partial, sections)
|
||||
|
||||
if folded == 0 {
|
||||
t.Fatal("ни одна доставка не свернулась")
|
||||
}
|
||||
|
||||
// Половина живого потока не несёт metrics вовсе (находка 50): такие
|
||||
// доставки обязаны быть отличимы от разобранных целиком, иначе ретеншен
|
||||
// срежет тела, которые для stateOfMind единственный источник.
|
||||
if partial == 0 {
|
||||
t.Error("ни одной частично разобранной доставки — перечисление непокрытых секций не работает")
|
||||
}
|
||||
|
||||
// Повторный прогон того же журнала не меняет состояния: свёртка
|
||||
// детерминирована, и пересборка даёт то же, что живой приём.
|
||||
//
|
||||
|
||||
+206
-12
@@ -17,7 +17,9 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
)
|
||||
|
||||
// Layer — подробность, в которой метрика приехала. Выводится из выравнивания
|
||||
@@ -99,6 +101,19 @@ type Point struct {
|
||||
type Result struct {
|
||||
Points []Point
|
||||
|
||||
// Uncovered — верхнеуровневые ключи `data`, которых разбор не покрывает,
|
||||
// отсортированные и без повторов. Половина живого потока состоит из таких
|
||||
// доставок целиком (48 из 99: workouts и stateOfMind), и без этого списка
|
||||
// они неотличимы от разобранной доставки с пустой секцией метрик.
|
||||
//
|
||||
// Список канонизирован потому, что уезжает в базу и сравнивается между
|
||||
// доставками, а порядок ключей в JSON от HAE нестабилен.
|
||||
Uncovered []string
|
||||
// UncoveredDropped — сколько имён отброшено границей списка. Молчаливое
|
||||
// усечение сделало бы список уверенным, но неполным ответом на вопрос «что
|
||||
// останется потерянным, если тело удалить».
|
||||
UncoveredDropped int
|
||||
|
||||
// Metrics — сколько метрик встретилось в секции.
|
||||
Metrics int
|
||||
// SkippedNoTime — точки без разбираемой метки времени.
|
||||
@@ -155,10 +170,12 @@ func Parse(body []byte, meta Meta) (res Result, err error) {
|
||||
}
|
||||
}()
|
||||
|
||||
metrics, err := decodeMetrics(body)
|
||||
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 {
|
||||
@@ -188,7 +205,15 @@ func Parse(body []byte, meta Meta) (res Result, err error) {
|
||||
|
||||
res.HeaderLayer = headerLayer(meta.Aggregation)
|
||||
if err := assignLayers(groups, meta, &res); err != nil {
|
||||
return Result{Metrics: res.Metrics}, err
|
||||
// Список непокрытых секций переживает отказ: доставка, у которой не
|
||||
// определился слой, обязана остаться записью о том, что в теле есть
|
||||
// невосстановимая секция. Иначе ретеншен увидит failed без списка и
|
||||
// решит, что терять нечего.
|
||||
return Result{
|
||||
Metrics: res.Metrics,
|
||||
Uncovered: res.Uncovered,
|
||||
UncoveredDropped: res.UncoveredDropped,
|
||||
}, err
|
||||
}
|
||||
|
||||
total := 0
|
||||
@@ -237,12 +262,27 @@ type group struct {
|
||||
// 42 МиБ через map[string]any удерживает 197 МиБ кучи против 54 МиБ у этой
|
||||
// формы. Вместе с самим телом пик доходил бы до ~300 МиБ на доставку — это
|
||||
// OOM ровно на пике потока, когда терять доставки дороже всего.
|
||||
type envelope struct {
|
||||
Data struct {
|
||||
Metrics []metricEnvelope `json:"metrics"`
|
||||
} `json:"data"`
|
||||
// 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"`
|
||||
@@ -261,13 +301,167 @@ type pointHead struct {
|
||||
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 // причина уходит в лог, наружу не раскрывается
|
||||
// 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 // причина уходит в лог, наружу не раскрывается
|
||||
}
|
||||
return env.Data.Metrics, nil
|
||||
|
||||
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 разбирает точки одной метрики. Точка без разбираемой метки
|
||||
|
||||
@@ -3,11 +3,14 @@ package hae_test
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"slices"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
"unicode/utf8"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/hae"
|
||||
)
|
||||
@@ -544,3 +547,157 @@ func FuzzParse(f *testing.F) {
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// Половина живого потока состоит из непокрытых секций целиком (48 доставок из
|
||||
// 99: workouts и stateOfMind). Без списка они неотличимы от разобранной
|
||||
// доставки с пустой секцией метрик, и ретеншен, ориентируясь на статус, срезал
|
||||
// бы тела, которые для stateOfMind единственный источник.
|
||||
func TestParseПеречисляетНепокрытыеСекции(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
res, err := hae.Parse(load(t, "uncovered_sections.json"), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
|
||||
want := []string{"stateOfMind", "workouts"}
|
||||
if !slices.Equal(res.Uncovered, want) {
|
||||
t.Errorf("непокрытые %v, ожидались %v", res.Uncovered, want)
|
||||
}
|
||||
if len(res.Points) == 0 {
|
||||
t.Error("точки метрик обязаны сохраниться: непокрытая секция рядом им не мешает")
|
||||
}
|
||||
if res.UncoveredDropped != 0 {
|
||||
t.Errorf("отброшено %d имён, ожидалось 0", res.UncoveredDropped)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseНепокрытыеСекцииГраницыИДетерминизм(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
t.Run("доставка из одной непокрытой секции", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
res, err := hae.Parse([]byte(`{"data":{"stateOfMind":[{"x":1}]}}`), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
if !slices.Equal(res.Uncovered, []string{"stateOfMind"}) {
|
||||
t.Errorf("непокрытые %v", res.Uncovered)
|
||||
}
|
||||
if len(res.Points) != 0 {
|
||||
t.Errorf("точек %d, ожидалось 0", len(res.Points))
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("доставка из одних метрик списка не даёт", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
res, err := hae.Parse([]byte(`{"data":{"metrics":[]}}`), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
if len(res.Uncovered) != 0 {
|
||||
t.Errorf("непокрытые %v, ожидался пустой список", res.Uncovered)
|
||||
}
|
||||
})
|
||||
|
||||
// Порядок ключей в JSON от HAE нестабилен, а список уезжает в базу и
|
||||
// сравнивается между доставками: зависящий от порядка на проводе список
|
||||
// сравнивать нельзя.
|
||||
t.Run("порядок и повторы не влияют", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
bodies := []string{
|
||||
`{"data":{"workouts":[],"stateOfMind":[],"ecg":[]}}`,
|
||||
`{"data":{"ecg":[],"workouts":[],"stateOfMind":[]}}`,
|
||||
`{"data":{"stateOfMind":[],"ecg":[],"workouts":[],"ecg":[]}}`,
|
||||
}
|
||||
want := []string{"ecg", "stateOfMind", "workouts"}
|
||||
for _, b := range bodies {
|
||||
res, err := hae.Parse([]byte(b), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор %s: %v", b, err)
|
||||
}
|
||||
if !slices.Equal(res.Uncovered, want) {
|
||||
t.Errorf("тело %s дало %v, ожидалось %v", b, res.Uncovered, want)
|
||||
}
|
||||
}
|
||||
})
|
||||
|
||||
// Повтор ключа metrics JSON допускает; секции объединяются, а не побеждает
|
||||
// последняя — терять точки молча нельзя.
|
||||
t.Run("повтор metrics объединяет секции", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const body = `{"data":{"metrics":[{"name":"a","data":[]}],"metrics":[{"name":"b","data":[]}]}}`
|
||||
res, err := hae.Parse([]byte(body), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
if res.Metrics != 2 {
|
||||
t.Errorf("метрик %d, ожидалось 2 — секции обязаны объединиться", res.Metrics)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("границы списка видны", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var sb strings.Builder
|
||||
sb.WriteString(`{"data":{`)
|
||||
for i := range 40 {
|
||||
if i > 0 {
|
||||
sb.WriteByte(',')
|
||||
}
|
||||
fmt.Fprintf(&sb, `"s%02d":[]`, i)
|
||||
}
|
||||
sb.WriteString(`}}`)
|
||||
|
||||
res, err := hae.Parse([]byte(sb.String()), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
if len(res.Uncovered) != 32 {
|
||||
t.Errorf("имён %d, ожидалось 32", len(res.Uncovered))
|
||||
}
|
||||
if res.UncoveredDropped != 8 {
|
||||
t.Errorf("отброшено %d, ожидалось 8 — иначе «ровно 32» неотличимо от «пришло пятьсот»", res.UncoveredDropped)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("длинное имя обрезано и помечено", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
long := strings.Repeat("щ", 60) // 120 байт
|
||||
res, err := hae.Parse([]byte(`{"data":{"`+long+`":[]}}`), hae.Meta{})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
if len(res.Uncovered) != 1 {
|
||||
t.Fatalf("непокрытые %v", res.Uncovered)
|
||||
}
|
||||
got := res.Uncovered[0]
|
||||
if !strings.HasSuffix(got, "…") {
|
||||
t.Errorf("имя %q без маркера обрезки: обрезка не инъективна и обязана быть видимой", got)
|
||||
}
|
||||
if !utf8.ValidString(got) {
|
||||
t.Errorf("имя %q обрезано посреди руны", got)
|
||||
}
|
||||
})
|
||||
|
||||
// Ошибка после разобранной секции metrics: правило прежнее — при ошибке
|
||||
// точек нет вовсе, иначе часть точек оказалась бы в витрине под статусом,
|
||||
// по которому доставку никто не подберёт.
|
||||
t.Run("обрыв после метрик точек не отдаёт", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const body = `{"data":{"metrics":[{"name":"a","data":[{"date":"2025-06-05 10:00:00 +0300","qty":1}]}],"work`
|
||||
res, err := hae.Parse([]byte(body), hae.Meta{})
|
||||
if !errors.Is(err, hae.ErrMalformed) {
|
||||
t.Fatalf("ошибка %v, ожидалась ErrMalformed", err)
|
||||
}
|
||||
if len(res.Points) != 0 {
|
||||
t.Errorf("точек %d, ожидалось 0", len(res.Points))
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -82,3 +82,79 @@ func largeBody(metrics, points int) []byte {
|
||||
b.WriteString(`]}}`)
|
||||
return []byte(b.String())
|
||||
}
|
||||
|
||||
// Перечисление непокрытых секций не имеет права стать вторым способом удержать
|
||||
// тело. Непокрытая секция ПРОГЛАТЫВАЕТСЯ декодированием в выбрасываемый
|
||||
// RawMessage: копия одна, живёт до следующего члена, удерживается ноль. Разбор
|
||||
// `data` в map[string]json.RawMessage дешевле по коду, но копирует байты всех
|
||||
// секций и держит их до конца разбора — этот тест сторожит разницу.
|
||||
func TestParseУдержаниеКучиНепокрытойСекции(t *testing.T) {
|
||||
if testing.Short() {
|
||||
t.Skip("измерение кучи: не для -short")
|
||||
}
|
||||
|
||||
body := bodyWithUncovered(40 << 20)
|
||||
t.Logf("тело %d МиБ, из них непокрытая секция почти всё", len(body)>>20)
|
||||
|
||||
var before, after runtime.MemStats
|
||||
runtime.GC()
|
||||
runtime.ReadMemStats(&before)
|
||||
|
||||
// Слой единственной точке взяться неоткуда — наследуем, иначе разбор
|
||||
// откажет раньше, чем дойдёт до замера.
|
||||
res, err := hae.Parse(body, hae.Meta{FallbackLayer: hae.LayerMinute})
|
||||
if err != nil {
|
||||
t.Fatalf("разбор: %v", err)
|
||||
}
|
||||
|
||||
runtime.GC()
|
||||
runtime.ReadMemStats(&after)
|
||||
|
||||
retained := int64(after.HeapAlloc) - int64(before.HeapAlloc)
|
||||
limit := int64(len(body)) * 4
|
||||
|
||||
t.Logf("непокрытых %v, удержано %d МиБ при теле %d МиБ",
|
||||
res.Uncovered, retained>>20, len(body)>>20)
|
||||
|
||||
if retained > limit {
|
||||
t.Errorf("удержано %d МиБ при теле %d МиБ — больше четырёх тел; "+
|
||||
"похоже, секции удерживаются, а не проглатываются",
|
||||
retained>>20, len(body)>>20)
|
||||
}
|
||||
|
||||
runtime.KeepAlive(res)
|
||||
runtime.KeepAlive(body)
|
||||
}
|
||||
|
||||
// Тело из вложенных скобок обязано быть ОТВЕРГНУТО, а не съедено. Пропуск
|
||||
// ручным счётом глубины по Token() этого не давал: делимитеры идут мимо
|
||||
// сканера, ограничитель вложенности encoding/json не работает, и стек токенов
|
||||
// растёт как O(глубины) — измерено 488 МиБ пика на теле 40 МиБ.
|
||||
func TestParseВложенныеСкобкиОтвергаются(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const depth = 100_000
|
||||
var b strings.Builder
|
||||
b.WriteString(`{"data":{"deep":`)
|
||||
b.WriteString(strings.Repeat("[", depth))
|
||||
b.WriteString(strings.Repeat("]", depth))
|
||||
b.WriteString(`}}`)
|
||||
|
||||
if _, err := hae.Parse([]byte(b.String()), hae.Meta{}); err == nil {
|
||||
t.Error("тело из ста тысяч уровней вложенности принято — ограничитель вложенности не работает")
|
||||
}
|
||||
}
|
||||
|
||||
func bodyWithUncovered(size int) []byte {
|
||||
var b strings.Builder
|
||||
b.WriteString(`{"data":{"metrics":[{"name":"m","units":"count","data":[`)
|
||||
b.WriteString(`{"date":"2025-06-05 10:00:00 +0300","qty":1}`)
|
||||
b.WriteString(`]}],"stateOfMind":[`)
|
||||
|
||||
const entry = `{"id":"00000000-0000-0000-0000-000000000000","kind":"momentary_emotion","valence":0.5},`
|
||||
for b.Len() < size {
|
||||
b.WriteString(entry)
|
||||
}
|
||||
b.WriteString(`{"id":"tail"}]}}`)
|
||||
return []byte(b.String())
|
||||
}
|
||||
|
||||
+49
@@ -0,0 +1,49 @@
|
||||
{
|
||||
"data": {
|
||||
"metrics": [
|
||||
{
|
||||
"name": "step_count",
|
||||
"units": "count",
|
||||
"data": [
|
||||
{"date": "2025-06-05 09:00:00 +0300", "qty": 41.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:01:00 +0300", "qty": 17.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:02:00 +0300", "qty": 82.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:03:00 +0300", "qty": 5.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:04:00 +0300", "qty": 63.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:05:00 +0300", "qty": 28.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:06:00 +0300", "qty": 94.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:07:00 +0300", "qty": 12.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:08:00 +0300", "qty": 71.0, "source": "Device A"},
|
||||
{"date": "2025-06-05 09:09:00 +0300", "qty": 36.0, "source": "Device A"}
|
||||
]
|
||||
}
|
||||
],
|
||||
"workouts": [
|
||||
{
|
||||
"id": "00000000-0000-0000-0000-000000000001",
|
||||
"name": "В помещении Ходьба",
|
||||
"start": "2025-06-05 08:00:00 +0300",
|
||||
"end": "2025-06-05 08:30:00 +0300",
|
||||
"duration": 1800,
|
||||
"heartRateData": [
|
||||
{"date": "2025-06-05 08:00:07 +0300", "Min": 91.0, "Avg": 94.5, "Max": 98.0, "units": "count/min"},
|
||||
{"date": "2025-06-05 08:00:21 +0300", "Min": 93.0, "Avg": 95.5, "Max": 99.0, "units": "count/min"}
|
||||
],
|
||||
"route": [
|
||||
{"lat": 10.0, "lon": 20.0, "altitude": 30.0, "timestamp": "2025-06-05 08:00:07 +0300"}
|
||||
]
|
||||
}
|
||||
],
|
||||
"stateOfMind": [
|
||||
{
|
||||
"id": "00000000-0000-0000-0000-000000000002",
|
||||
"start": "2025-06-05T05:12:33Z",
|
||||
"end": "2025-06-05T05:12:33Z",
|
||||
"kind": "momentary_emotion",
|
||||
"valence": 0.25,
|
||||
"labels": ["slightly_pleasant"],
|
||||
"associations": []
|
||||
}
|
||||
]
|
||||
}
|
||||
}
|
||||
@@ -3,6 +3,7 @@ package store
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"time"
|
||||
@@ -11,10 +12,18 @@ import (
|
||||
// Статусы разбора доставки. Код ответа приёма от них не зависит: сохранили —
|
||||
// значит приняли.
|
||||
const (
|
||||
// ParsePending — тело сохранено, разбора ещё не было.
|
||||
// ParsePending — этим разбором тело ещё не смотрели.
|
||||
//
|
||||
// Смысл именно такой, а не «тела ещё не касались»: миграция 00005 перевела
|
||||
// сюда доставки, разобранные кодом, который не различал частичный разбор.
|
||||
// Статус консервативный — ретеншен не трогает pending никогда.
|
||||
ParsePending = "pending"
|
||||
// ParseDone — тело разобрано, точки разложены по объектам.
|
||||
// ParseDone — разобрано всё, что в теле было.
|
||||
ParseDone = "parsed"
|
||||
// ParsePartial — разобрано покрытое, но в теле остались секции, которых
|
||||
// разбор не покрывает. Не отклонение, а установившееся состояние половины
|
||||
// потока: 48 доставок из 99 несут только workouts или stateOfMind.
|
||||
ParsePartial = "partial"
|
||||
// ParseFailed — разобрать не удалось. Тело лежит в архиве, доставку
|
||||
// подберёт пересборка.
|
||||
ParseFailed = "failed"
|
||||
@@ -38,6 +47,10 @@ type Delivery struct {
|
||||
// значений), уже без секретов. Именованные поля выше дублируют часть из
|
||||
// них: по ним ходят запросы, а Headers хранит всё остальное на будущее.
|
||||
Headers string
|
||||
// UncoveredSections — секции тела, которых разбор не покрыл, JSON-массивом
|
||||
// имён. Ответ на вопрос «что останется потерянным, если тело удалить»:
|
||||
// ретеншен обязан спрашивать его прежде, чем срезать тело.
|
||||
UncoveredSections string
|
||||
}
|
||||
|
||||
// CreateDelivery записывает факт приёма пакета.
|
||||
@@ -69,7 +82,7 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
|
||||
const q = `
|
||||
SELECT id, received_at, automation_name, automation_id, aggregation,
|
||||
period, session_id, bytes, sha256, raw_path, parse_status,
|
||||
points, headers
|
||||
points, headers, uncovered_sections
|
||||
FROM delivery ORDER BY received_at DESC, id DESC LIMIT 1`
|
||||
|
||||
var d Delivery
|
||||
@@ -77,7 +90,7 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
|
||||
err := s.db.QueryRowxContext(ctx, q).Scan(
|
||||
&d.ID, &receivedAt, &d.AutomationName, &d.AutomationID, &d.Aggregation,
|
||||
&d.Period, &d.SessionID, &d.Bytes, &d.SHA256, &d.RawPath,
|
||||
&d.ParseStatus, &d.Points, &d.Headers)
|
||||
&d.ParseStatus, &d.Points, &d.Headers, &d.UncoveredSections)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return Delivery{}, ErrNotFound
|
||||
}
|
||||
@@ -101,14 +114,40 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
|
||||
// Пустой layer означает «не трогать»: у неудачной свёртки слоя нет, а
|
||||
// затирание оборвало бы цепочку наследования — в том числе у доставки, которая
|
||||
// раньше свернулась успешно, — и результат пересборки журнала изменился бы.
|
||||
func (s *Store) FinishParse(ctx context.Context, id, status string, points int64, layer string) error {
|
||||
// ParseOutcome — исход разбора доставки. Структурой, а не растущим списком
|
||||
// позиционных параметров: у FinishParse их было уже четыре, и пятый неизбежно
|
||||
// перепутали бы местами с четвёртым.
|
||||
type ParseOutcome struct {
|
||||
Status string
|
||||
Points int64
|
||||
// Layer пустой означает «не трогать» — см. FinishParse.
|
||||
Layer string
|
||||
// Uncovered замещает прежнее значение ЦЕЛИКОМ, включая замещение пустым:
|
||||
// у слоя пустота — отсутствие знания, у списка — знание об отсутствии.
|
||||
// Пересвёртка доставки, чья секция стала покрытой, обязана список очистить.
|
||||
Uncovered []string
|
||||
}
|
||||
|
||||
func (s *Store) FinishParse(ctx context.Context, id string, out ParseOutcome) error {
|
||||
const q = `
|
||||
UPDATE delivery
|
||||
SET parse_status = ?, points = ?,
|
||||
derived_layer = CASE WHEN ? = '' THEN derived_layer ELSE ? END
|
||||
derived_layer = CASE WHEN ? = '' THEN derived_layer ELSE ? END,
|
||||
uncovered_sections = ?
|
||||
WHERE id = ?`
|
||||
|
||||
res, err := s.db.ExecContext(ctx, q, status, points, layer, layer, id)
|
||||
// Ровно одно представление пустоты — `[]`: nil-срез Go сериализуется как
|
||||
// null, и в колонке появилось бы второе значение с тем же смыслом.
|
||||
sections := out.Uncovered
|
||||
if sections == nil {
|
||||
sections = []string{}
|
||||
}
|
||||
encoded, err := json.Marshal(sections)
|
||||
if err != nil {
|
||||
return fmt.Errorf("encode uncovered sections: %w", err)
|
||||
}
|
||||
|
||||
res, err := s.db.ExecContext(ctx, q, out.Status, out.Points, out.Layer, out.Layer, string(encoded), id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update parse status: %w", err)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"io/fs"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
|
||||
"github.com/jmoiron/sqlx"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
// Миграция 00005 переводит существующие `parsed` в `pending`. Статус, который
|
||||
// поставил код, не различавший частичного разбора, ничего не доказывает: под
|
||||
// ним лежат и полностью разобранные доставки, и доставки из одних тренировок с
|
||||
// нулём точек (48 из 99 на живом архиве). Ретеншен, ради которого признак и
|
||||
// заводится, поверил бы им и срезал тела — а для stateOfMind это необратимо.
|
||||
func TestMigrationПрежниеParsedСтановятсяPending(t *testing.T) {
|
||||
db, err := sqlx.Connect("sqlite", dsn(filepath.Join(t.TempDir(), "healthlog.db")))
|
||||
if err != nil {
|
||||
t.Fatalf("открытие базы: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = db.Close() })
|
||||
|
||||
sub, err := fs.Sub(migrationsFS, "migrations")
|
||||
if err != nil {
|
||||
t.Fatalf("миграции: %v", err)
|
||||
}
|
||||
p, err := goose.NewProvider(goose.DialectSQLite3, db.DB, sub)
|
||||
if err != nil {
|
||||
t.Fatalf("провайдер: %v", err)
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
// Состояние ДО этой миграции: схема четвёртой версии.
|
||||
if _, err := p.UpTo(ctx, 4); err != nil {
|
||||
t.Fatalf("миграция до 4: %v", err)
|
||||
}
|
||||
|
||||
const insert = `
|
||||
INSERT INTO delivery (id, received_at, automation_name, automation_id,
|
||||
aggregation, period, session_id, bytes, sha256, raw_path,
|
||||
parse_status, points)
|
||||
VALUES (?, '2025-06-05T10:00:00Z', '', '', '', '', '', 0, '-', '-', ?, ?)`
|
||||
for _, c := range []struct {
|
||||
id string
|
||||
status string
|
||||
points int
|
||||
}{
|
||||
{"d-parsed-points", "parsed", 12},
|
||||
{"d-parsed-empty", "parsed", 0},
|
||||
{"d-failed", "failed", 0},
|
||||
{"d-pending", "pending", 0},
|
||||
} {
|
||||
if _, err := db.ExecContext(ctx, insert, c.id, c.status, c.points); err != nil {
|
||||
t.Fatalf("вставка %s: %v", c.id, err)
|
||||
}
|
||||
}
|
||||
|
||||
if _, err := p.Up(ctx); err != nil {
|
||||
t.Fatalf("миграция до последней: %v", err)
|
||||
}
|
||||
|
||||
want := map[string]string{
|
||||
// Переводятся ОБА варианта parsed, а не только пустой: наблюдение
|
||||
// «секции не смешиваются» собрано за двое суток потока, и ставить на
|
||||
// него необратимое удаление тел значит повторять ошибку, ради которой
|
||||
// задача и заведена.
|
||||
"d-parsed-points": "pending",
|
||||
"d-parsed-empty": "pending",
|
||||
"d-failed": "failed",
|
||||
"d-pending": "pending",
|
||||
}
|
||||
for id, wantStatus := range want {
|
||||
var status, sections string
|
||||
err := db.QueryRowContext(ctx,
|
||||
`SELECT parse_status, uncovered_sections FROM delivery WHERE id = ?`, id).
|
||||
Scan(&status, §ions)
|
||||
if err != nil {
|
||||
t.Fatalf("чтение %s: %v", id, err)
|
||||
}
|
||||
if status != wantStatus {
|
||||
t.Errorf("%s: статус %q, ожидался %q", id, status, wantStatus)
|
||||
}
|
||||
if sections != "[]" {
|
||||
t.Errorf("%s: список %q, ожидался `[]`", id, sections)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
-- +goose Up
|
||||
-- Секции тела, которых разбор не покрыл. JSON-массив имён (`["stateOfMind"]`),
|
||||
-- пустой список — `[]`.
|
||||
--
|
||||
-- Массивом, а не строкой с разделителем: имя ключа приходит из чужого тела и
|
||||
-- может содержать что угодно, включая пробел и запятую. JSON снимает вопрос
|
||||
-- разделителя, согласуется с колонкой headers и читается из SQLite через
|
||||
-- json_each, если ретеншену это понадобится.
|
||||
--
|
||||
-- Колонка — ответ на вопрос «что останется потерянным, если тело удалить».
|
||||
-- Для stateOfMind ответ необратим: в экспорте Apple его нет (находка 46),
|
||||
-- доставки HAE для него единственный источник.
|
||||
ALTER TABLE delivery ADD COLUMN uncovered_sections TEXT NOT NULL DEFAULT '[]';
|
||||
|
||||
-- Статус parsed, поставленный кодом, который частичного разбора не различал,
|
||||
-- ничего не доказывает: под ним лежат и полностью разобранные доставки, и
|
||||
-- доставки из одних тренировок с нулём точек (48 из 99 на живом архиве).
|
||||
-- Ретеншен, ради которого признак и заводится, поверил бы им и срезал тела.
|
||||
--
|
||||
-- pending — консервативный статус: ретеншен не трогает его никогда, а подбор
|
||||
-- pending пересвернёт доставки из архива. Свёртка идемпотентна, повторный
|
||||
-- прогон журнала состояния не меняет.
|
||||
--
|
||||
-- Целевой перевод только строк с points = 0 рассматривался и отвергнут: он
|
||||
-- опирается на наблюдение «секции не смешиваются», собранное за двое суток
|
||||
-- потока, а ставить на такое наблюдение необратимое удаление тел значит
|
||||
-- повторять ошибку, ради которой задача и заведена.
|
||||
UPDATE delivery SET parse_status = 'pending' WHERE parse_status = 'parsed';
|
||||
|
||||
-- +goose Down
|
||||
-- ВНИМАНИЕ: миграция односторонняя ПО ДАННЫМ. Down снимает колонку, но какие
|
||||
-- доставки были parsed, восстановить неоткуда — состояние пересобирается из
|
||||
-- архива, а не откатом. До появления подбора pending строки останутся в этом
|
||||
-- статусе.
|
||||
ALTER TABLE delivery DROP COLUMN uncovered_sections;
|
||||
Reference in New Issue
Block a user