package store import ( "bytes" "context" "database/sql" "encoding/json" "errors" "fmt" "time" "git.vakhrushev.me/av/healthlog/internal/canon" ) // Роды таблиц сущностей. Тренировка живёт в своей таблице: у неё есть // заголовок, по которому идёт выборка, а у записи его нет. const ( workoutTable = "workout" recordTable = "record" ) // IncomingEntity — сущность с собственным идентификатором, пришедшая на запись. // // Содержимое хранится исходными байтами: сущность, пересобранная повторной // сериализацией, теряет литерал ровно так же, как точка. type IncomingEntity struct { ID string Kind string Name string Start time.Time End time.Time // OffsetSeconds — смещение зоны начала. OffsetSeconds int // Duration — длительность тренировки в секундах. nil означает «источник не // прислал» и отличим от нуля: ноль — законная длительность. Duration *float64 Raw json.RawMessage } // Incoming — всё, что дала одна доставка. Единицей записи является доставка, а // не точка и не сущность: частичное состояние ломает инвариант «состояние // пересобираемо». type Incoming struct { Points []IncomingPoint Workouts []IncomingEntity Records []IncomingEntity } // DeliveryRef — место доставки в журнале. Пара, а не идентификатор: по ней // разрешается тай-брейк между версиями сущности равной полноты, а порядок // журнала задан парой `(received_at, id)`. type DeliveryRef struct { ID string ReceivedAt time.Time } func (d DeliveryRef) before(other DeliveryRef) bool { if !d.ReceivedAt.Equal(other.ReceivedAt) { return d.ReceivedAt.Before(other.ReceivedAt) } return d.ID < other.ID } // EntityRef — координаты сущности для записи в лог. Содержимого не несёт: // маршрут тренировки — это геотрек до дома, а метки состояния разума — // измерение душевного состояния. type EntityRef struct { Kind string ID string } // maxEntityRefsReported — сколько координат сущностей попадает в лог. const maxEntityRefsReported = 5 // prepareEntities считает хеш каждой сущности. // // Вынесено из транзакции намеренно: хеширование материализует значение целиком // (маршрут — до мегабайта), а транзакция повторяется до пяти раз при занятости // базы. func prepareEntities(in []IncomingEntity, from DeliveryRef) ([]entityVersion, error) { if len(in) == 0 { return nil, nil } out := make([]entityVersion, 0, len(in)) for _, e := range in { v, err := newEntityVersion(e, from) if err != nil { return nil, err } out = append(out, v) } return out, nil } // clipRefs держит список координат в потолке: он зацепка для разбора, а не // отчёт; масштаб события считает счётчик. func clipRefs(refs []EntityRef) []EntityRef { if len(refs) <= maxEntityRefsReported { return refs } return refs[:maxEntityRefsReported] } // entityVersion — версия сущности вместе с тем, что нужно знать при выборе // победителя. // // Разбор и канонизация ОТЛОЖЕНЫ: они нужны только когда хеш разошёлся с // сохранённым, то есть на одной доставке из сорока четырёх. Считать их сразу // значило бы разворачивать маршрут (95% веса тренировки, до мегабайта) в дерево // значений на каждой копии — ровно та форма, от которой разбор тела отказался // замером (197 МиБ кучи против 54 МиБ на теле 42 МиБ). Хеш при этом считается // сразу и один раз на доставку: он и есть быстрый путь. type entityVersion struct { raw json.RawMessage hash string from DeliveryRef // key и fields заполняются лениво, методом analyze(). key []byte fields canon.Fields parsed bool // head — заголовок, который пишется колонками. У сохранённой версии он не // нужен: она либо побеждает и остаётся как есть, либо замещается целиком. head IncomingEntity } func newEntityVersion(e IncomingEntity, from DeliveryRef) (entityVersion, error) { h, err := canon.Hash(e.Raw) if err != nil { return entityVersion{}, fmt.Errorf("хеш сущности: %w", err) } return entityVersion{raw: e.Raw, hash: h, from: from, head: e}, nil } // analyze разбирает версию, если этого ещё не делали. func (v *entityVersion) analyze() { if v.parsed { return } v.key = canon.SortKey(v.raw) v.fields = canon.Analyze(v.raw) v.parsed = true } // pickEntity выбирает между сохранённой и приехавшей версией. // // Второй возврат — потеряла ли бы витрина содержание, приняв приехавшую. Это и // есть плата за отказ объединять поля: событие не предотвращается молча, а // считается и уходит в WARN. // // 1. приехавшая несёт всё содержание сохранённой и сверх того → приехавшая // 2. сохранённая несёт всё содержание приехавшей и сверх того → сохранённая // 3. содержание равно → версия из более поздней доставки ЖУРНАЛА // 4. наборы несравнимы → сохранённая // // Пункт 3 — не «побеждает приехавшая». Приехавшая есть функция порядка // СВЁРТКИ, а он порядку журнала не равен: воркер сворачивает в порядке журнала // только среди видимых ему доставок. Доставка с более ранней меткой, свёрнутая // позже, вернула бы витрину к недосчитанной версии, и пересборка разошлась бы // с живым приёмом молча — в содержимом тренировки, где это не видно ничем, // кроме отпечатка. // // Равные позиции означают две версии одного ключа внутри ОДНОЙ доставки; там // решает минимум канонической формы, потому что порядок элементов в // JSON-массиве нестабилен. func pickEntity(stored, incoming *entityVersion) (takeIncoming, lost bool) { switch v := compareEntities(stored, incoming); v { case entityIncomingRicher: return true, false case entityStoredRicher: return false, true case entityIncomparable: // Несравнимы: у каждой версии есть содержание, которого нет у другой. // Объединение полей отвергнуто там же и по той же причине, что для // точек, — на живом потоке событие не наступало ни разу, — а из двух // версий остаётся сохранённая: правило называется «не теряет // содержания», и приехавшая его теряет. Исход при этом остаётся // функцией журнала: доставки проигрываются в его порядке. return false, true default: return laterInJournal(stored, incoming), false } } // pickWithinDelivery выбирает между двумя версиями одного ключа ВНУТРИ одной // доставки. Второй возврат — различается ли их содержание вообще. // // Отдельно от pickEntity, и не ради симметрии: «сохранённой» версии здесь нет, // есть только порядок элементов в JSON-массиве, а он нестабилен. Правило // «остаётся первая встреченная» сделало бы исход функцией порядка на проводе, // поэтому при равном и при несравнимом содержании решает тотальный порядок // канонических форм. func pickWithinDelivery(a, b *entityVersion) (takeB, differs bool) { switch v := compareEntities(a, b); v { case entityIncomingRicher: return true, true case entityStoredRicher: return false, true case entityIncomparable: return laterInJournal(a, b), true default: return laterInJournal(a, b), false } } // entityVerdict — как соотносится СОДЕРЖАНИЕ двух версий одной сущности. // Нумерация с единицы: нулевое значение не должно выглядеть как «равны». type entityVerdict int const ( // entityEqualContent — множества содержательных ключей совпадают, длины // верхнеуровневых массивов тоже. Значения при этом могут расходиться: их // сравнение здесь неприменимо (см. canon.Fields.Covers). entityEqualContent entityVerdict = iota + 1 entityIncomingRicher entityStoredRicher entityIncomparable ) func compareEntities(stored, incoming *entityVersion) entityVerdict { stored.analyze() incoming.analyze() storedCovers := stored.fields.Covers(incoming.fields) incomingCovers := incoming.fields.Covers(stored.fields) switch { case incomingCovers && storedCovers: return entityEqualContent case incomingCovers: return entityIncomingRicher case storedCovers: return entityStoredRicher default: return entityIncomparable } } // laterInJournal говорит, стоит ли приехавшая версия позже сохранённой в // журнале. Позиции равны у двух версий одного ключа внутри одной доставки; // там решает минимум канонической формы — порядок тотальный и от порядка // элементов в массиве не зависит. func laterInJournal(stored, incoming *entityVersion) bool { if stored.from.before(incoming.from) { return true } if incoming.from.before(stored.from) { return false } return bytes.Compare(incoming.key, stored.key) < 0 } // dedupeEntities сворачивает версии одного ключа ВНУТРИ доставки тем же // правилом — до сравнения с сохранённой. // // Без этого исход зависел бы от того, как написан цикл: карта по ключу дала бы // победу последнему элементу массива мимо правила полноты, а порядок элементов // в JSON-массиве нестабилен. // // Счётчик здесь считает СИММЕТРИЧНО — «в одном теле приехали две версии одного // ключа с разным содержанием», — а не «приехавшая обеднена». Внутри доставки // «сохранённой» версии не существует, есть только порядок элементов массива, и // счётчик, зависящий от него, наблюдал бы событие через раз. func dedupeEntities(versions []entityVersion) ([]entityVersion, int, []EntityRef) { type slot struct { v entityVersion pos int } byKey := make(map[EntityRef]slot, len(versions)) order := make([]EntityRef, 0, len(versions)) held := 0 var heldAt []EntityRef for _, v := range versions { ref := EntityRef{Kind: v.head.Kind, ID: v.head.ID} prev, seen := byKey[ref] if !seen { byKey[ref] = slot{v: v, pos: len(order)} order = append(order, ref) continue } takeB, differs := pickWithinDelivery(&prev.v, &v) if differs { held++ if len(heldAt) < maxEntityRefsReported { heldAt = append(heldAt, ref) } } winner := prev.v if takeB { winner = v } byKey[ref] = slot{v: winner, pos: prev.pos} } out := make([]entityVersion, 0, len(order)) for _, ref := range order { out = append(out, byKey[ref].v) } return out, held, heldAt } // mergeEntities сливает сущности одной секции с сохранёнными. func mergeEntities(ctx context.Context, tx *sql.Tx, table string, versions []entityVersion, now time.Time) (written, held int, heldAt []EntityRef, err error) { for _, v := range versions { stored, found, err := readEntityHead(ctx, tx, table, v.head.Kind, v.head.ID) if err != nil { return 0, 0, nil, err } if !found { if err := writeEntity(ctx, tx, table, v, now); err != nil { return 0, 0, nil, err } written++ continue } // Хеш — детектор изменений: совпал, значит писать нечего, и содержимое // сохранённой сущности читать не приходится вовсе. Тренировка // переприсылается каждой доставкой, пока не доедет маршрут, — на живом // архиве 44 копии дают три различных содержимых. if stored.hash == v.hash { continue } storedRaw, err := readEntityPayload(ctx, tx, table, v.head.Kind, v.head.ID) if err != nil { return 0, 0, nil, err } prev := entityVersion{raw: storedRaw, hash: stored.hash, from: stored.from} takeIncoming, lost := pickEntity(&prev, &v) if lost { held++ if len(heldAt) < maxEntityRefsReported { heldAt = append(heldAt, EntityRef{Kind: v.head.Kind, ID: v.head.ID}) } } if !takeIncoming { continue } if err := writeEntity(ctx, tx, table, v, now); err != nil { return 0, 0, nil, err } written++ } return written, held, heldAt, nil } type storedEntityHead struct { hash string from DeliveryRef } func readEntityHead(ctx context.Context, tx *sql.Tx, table, kind, id string) (storedEntityHead, bool, error) { q := `SELECT content_hash, delivery_id, delivery_received_at FROM ` + table + entityWhere(table) var ( head storedEntityHead receivedAt string ) row := queryEntity(ctx, tx, q, table, kind, id) err := row.Scan(&head.hash, &head.from.ID, &receivedAt) if errors.Is(err, sql.ErrNoRows) { return storedEntityHead{}, false, nil } if err != nil { return storedEntityHead{}, false, fmt.Errorf("select %s: %w", table, err) } // Пустую метку не терпим: колонка NOT NULL без умолчания, и пустота здесь // означала бы дефект писателя. Молчаливый нулевой момент сделал бы // сохранённую версию «самой ранней в журнале», и её затирала бы любая // приехавшая — то есть дефект проявился бы потерей данных, а не отказом. head.from.ReceivedAt, err = ParseTime(receivedAt) if err != nil { return storedEntityHead{}, false, err } return head, true, nil } func readEntityPayload(ctx context.Context, tx *sql.Tx, table, kind, id string) (json.RawMessage, error) { q := `SELECT payload FROM ` + table + entityWhere(table) var payload []byte if err := queryEntity(ctx, tx, q, table, kind, id).Scan(&payload); err != nil { return nil, fmt.Errorf("select %s payload: %w", table, err) } raw, err := gunzipBytes(payload) if err != nil { return nil, err } return raw, nil } // entityWhere и queryEntity держат разницу между таблицами в одном месте: // у тренировки ключ — `id`, у записи — пара `kind + id`. func entityWhere(table string) string { if table == recordTable { return ` WHERE kind = ? AND id = ?` } return ` WHERE id = ?` } func queryEntity(ctx context.Context, tx *sql.Tx, q, table, kind, id string) *sql.Row { if table == recordTable { return tx.QueryRowContext(ctx, q, kind, id) } return tx.QueryRowContext(ctx, q, id) } func writeEntity(ctx context.Context, tx *sql.Tx, table string, v entityVersion, now time.Time) error { payload, err := gzipBytes(v.raw) if err != nil { return err } stamp := FormatTime(now) received := FormatTime(v.from.ReceivedAt) if table == recordTable { const q = ` INSERT INTO record (kind, id, ts_utc, tz_offset, payload, content_hash, delivery_id, delivery_received_at, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (kind, id) DO UPDATE SET ts_utc = excluded.ts_utc, tz_offset = excluded.tz_offset, payload = excluded.payload, content_hash = excluded.content_hash, delivery_id = excluded.delivery_id, delivery_received_at = excluded.delivery_received_at, updated_at = excluded.updated_at` if _, err := tx.ExecContext(ctx, q, v.head.Kind, v.head.ID, FormatTime(v.head.Start), v.head.OffsetSeconds, payload, v.hash, v.from.ID, received, stamp, stamp); err != nil { return fmt.Errorf("upsert record: %w", err) } return nil } const q = ` INSERT INTO workout (id, name, start_utc, end_utc, tz_offset, duration_sec, payload, content_hash, delivery_id, delivery_received_at, created_at, updated_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT (id) DO UPDATE SET name = excluded.name, start_utc = excluded.start_utc, end_utc = excluded.end_utc, tz_offset = excluded.tz_offset, duration_sec = excluded.duration_sec, payload = excluded.payload, content_hash = excluded.content_hash, delivery_id = excluded.delivery_id, delivery_received_at = excluded.delivery_received_at, updated_at = excluded.updated_at` var duration any if v.head.Duration != nil { duration = *v.head.Duration } if _, err := tx.ExecContext(ctx, q, v.head.ID, v.head.Name, FormatTime(v.head.Start), FormatTime(v.head.End), v.head.OffsetSeconds, duration, payload, v.hash, v.from.ID, received, stamp, stamp); err != nil { return fmt.Errorf("upsert workout: %w", err) } return nil } // Workout — тренировка, прочитанная из витрины. Нужна тестам и будущему // Read API: содержимое отдаётся целиком, заголовок — из колонок. type Workout struct { ID string Name string Start time.Time End time.Time OffsetSeconds int Duration *float64 Raw json.RawMessage Delivery string } // Workout читает тренировку по идентификатору. func (s *Store) Workout(ctx context.Context, id string) (Workout, error) { const q = ` SELECT id, name, start_utc, end_utc, tz_offset, duration_sec, payload, delivery_id FROM workout WHERE id = ?` var ( w Workout start, end string duration sql.NullFloat64 payload []byte deliveryFrom string ) err := s.db.QueryRowxContext(ctx, q, id). Scan(&w.ID, &w.Name, &start, &end, &w.OffsetSeconds, &duration, &payload, &deliveryFrom) if errors.Is(err, sql.ErrNoRows) { return Workout{}, ErrNotFound } if err != nil { return Workout{}, fmt.Errorf("select workout: %w", err) } if w.Start, err = ParseTime(start); err != nil { return Workout{}, err } if w.End, err = ParseTime(end); err != nil { return Workout{}, err } if duration.Valid { v := duration.Float64 w.Duration = &v } if w.Raw, err = gunzipBytes(payload); err != nil { return Workout{}, err } w.Delivery = deliveryFrom return w, nil } // Record — запись секции с собственным идентификатором. type Record struct { Kind string ID string TS time.Time OffsetSeconds int Raw json.RawMessage Delivery string } // Record читает запись по роду и идентификатору. func (s *Store) Record(ctx context.Context, kind, id string) (Record, error) { const q = ` SELECT kind, id, ts_utc, tz_offset, payload, delivery_id FROM record WHERE kind = ? AND id = ?` var ( r Record ts string payload []byte deliveryFrom string ) err := s.db.QueryRowxContext(ctx, q, kind, id). Scan(&r.Kind, &r.ID, &ts, &r.OffsetSeconds, &payload, &deliveryFrom) if errors.Is(err, sql.ErrNoRows) { return Record{}, ErrNotFound } if err != nil { return Record{}, fmt.Errorf("select record: %w", err) } if r.TS, err = ParseTime(ts); err != nil { return Record{}, err } if r.Raw, err = gunzipBytes(payload); err != nil { return Record{}, err } r.Delivery = deliveryFrom return r, nil } // CountWorkouts и CountRecords нужны отчёту пересборки: отпечаток отвечает // «да/нет», а по «да/нет» нельзя судить о направлении расхождения. func (s *Store) CountWorkouts(ctx context.Context) (int64, error) { var n int64 if err := s.db.GetContext(ctx, &n, `SELECT count(*) FROM workout`); err != nil { return 0, fmt.Errorf("count workouts: %w", err) } return n, nil } // CountRecords возвращает число записей секций с собственным `id`. func (s *Store) CountRecords(ctx context.Context) (int64, error) { var n int64 if err := s.db.GetContext(ctx, &n, `SELECT count(*) FROM record`); err != nil { return 0, fmt.Errorf("count records: %w", err) } return n, nil }