Дозакрыты находки ревью по слиянию сущностей
- Правило покрытия получило второй разряд (условный, как у точек), запрет вырождения формы и счёт содержательных элементов ряда: скелет из скаляров и ряд из null больше не затирают маршрут. Победитель внутри доставки стал функцией множества версий — общим помощником с точками, — а провенанс поднимается и при совпавшем хеше, иначе отложенная доставка возвращала витрину к прежнему содержимому. - Одно поле не того типа больше не уносит сущность, а пропуски видны в учётной записи доставки (миграция 00008, NULL = «не измерялось»); каноническая форма считается один раз и вне транзакции; откат бинаря поверх новой схемы отказывает на старте; текст ошибки разбора не несёт значений из тела. - Ревью кода профилем deep (девять проходов) нашло две регрессии и обе закрыты: безусловный второй разряд запирал законный досчёт навсегда, а выбор победителя был квадратичен по числу присланных версий одного ключа.
This commit is contained in:
+80
-54
@@ -84,9 +84,23 @@ type MergeStats struct {
|
||||
// EntitiesHeld — приехавшие версии, отклонённые как теряющие содержание
|
||||
// сохранённой (включая несравнимые наборы). Это и есть плата за отказ
|
||||
// объединять поля: событие считается, а не предотвращается молча.
|
||||
//
|
||||
// На него опирается ЕДИНСТВЕННЫЙ контроль того, что правило покрытия не
|
||||
// стало слишком строгим: сходимость отпечатка этого не проверяет по
|
||||
// построению — живой приём и пересборка пользуются одним правилом и
|
||||
// одинаково сойдутся на одинаково удержанной версии. Поэтому счётчик
|
||||
// обязан считать ровно удержания и ничего сверх.
|
||||
EntitiesHeld int
|
||||
// HeldAt — координаты первых таких сущностей, для записи в лог.
|
||||
HeldAt []EntityRef
|
||||
// EntitiesDiverging — версии одного ключа, приехавшие в ОДНОМ теле с разным
|
||||
// содержанием. Событие другого рода: победитель ложится в витрину целиком,
|
||||
// терять нечего, лечится оно не тем же. Считается отдельно от удержаний,
|
||||
// иначе одно число отвечало бы на два вопроса — и число удержаний, по
|
||||
// которому судят о строгости правила, стало бы неотличимо от шума.
|
||||
EntitiesDiverging int
|
||||
// DivergingAt — координаты первых таких сущностей.
|
||||
DivergingAt []EntityRef
|
||||
}
|
||||
|
||||
// Collision — координаты объекта, где столкновение разрешилось перезаписью
|
||||
@@ -140,10 +154,22 @@ func (s *Store) Merge(ctx context.Context, in Incoming, from DeliveryRef) (Merge
|
||||
groups := groupByHour(in.Points)
|
||||
keys := sortedKeys(groups)
|
||||
|
||||
// Хеш и каноническая форма сущности считаются ОДИН раз на доставку, до
|
||||
// входа в транзакцию: канонизация материализует значение целиком, а
|
||||
// транзакция повторяется до пяти раз при занятости базы — внутри неё пик
|
||||
// кучи умножился бы на число попыток.
|
||||
// Каноническая форма сущности и её хеш считаются ОДИН раз на версию и здесь
|
||||
// — до входа в транзакцию. Транзакция открыта `immediate`, то есть блокирует
|
||||
// запись, и повторяется до пяти раз при занятости базы: канонизация внутри
|
||||
// неё умножала бы и пик кучи, и время удержания блокировки. Измерено: тело
|
||||
// 40 МиБ даёт 768 МиБ пика, 63 МиБ удерживают блокировку 5.019 с при
|
||||
// busy_timeout 5000, после чего конкурентный CreateDelivery исчерпывает
|
||||
// повторы.
|
||||
//
|
||||
// Пределов остаётся два, и оба названы вслух. Первый: разбор СОХРАНЁННОЙ
|
||||
// версии остаётся внутри транзакции — её содержимое читается оттуда же и
|
||||
// только когда хеш разошёлся; удержание блокировки пропорционально её
|
||||
// размеру. Второй: форма и множества ключей всех версий доставки
|
||||
// УДЕРЖИВАЮТСЯ в памяти до конца транзакции, то есть расход пропорционален
|
||||
// размеру доставки, а не самой большой её сущности. Про процессор здесь
|
||||
// стало лучше, про память — хуже, и закрыть оба может лишь предел на размер
|
||||
// сущности вместе с потоковым расчётом.
|
||||
workouts, err := prepareEntities(in.Workouts, from)
|
||||
if err != nil {
|
||||
return MergeStats{}, err
|
||||
@@ -152,18 +178,25 @@ func (s *Store) Merge(ctx context.Context, in Incoming, from DeliveryRef) (Merge
|
||||
if err != nil {
|
||||
return MergeStats{}, err
|
||||
}
|
||||
workouts, workoutsHeld, workoutsHeldAt := dedupeEntities(workouts)
|
||||
records, recordsHeld, recordsHeldAt := dedupeEntities(records)
|
||||
workouts, workoutsDiverging, workoutsDivergingAt, err := dedupeEntities(ctx, workouts)
|
||||
if err != nil {
|
||||
return MergeStats{}, err
|
||||
}
|
||||
records, recordsDiverging, recordsDivergingAt, err := dedupeEntities(ctx, records)
|
||||
if err != nil {
|
||||
return MergeStats{}, err
|
||||
}
|
||||
|
||||
var stats MergeStats
|
||||
err = s.inTx(ctx, func(tx *sql.Tx) error {
|
||||
// Счётчики обнуляются на каждой попытке: повтор транзакции начинает
|
||||
// слияние заново, и накопленное от прошлой попытки посчиталось бы дважды.
|
||||
stats = MergeStats{
|
||||
Workouts: len(in.Workouts),
|
||||
Records: len(in.Records),
|
||||
EntitiesHeld: workoutsHeld + recordsHeld,
|
||||
HeldAt: clipRefs(append(append([]EntityRef{}, workoutsHeldAt...), recordsHeldAt...)),
|
||||
Workouts: len(in.Workouts),
|
||||
Records: len(in.Records),
|
||||
EntitiesDiverging: workoutsDiverging + recordsDiverging,
|
||||
DivergingAt: clipRefs(append(append([]EntityRef{},
|
||||
workoutsDivergingAt...), recordsDivergingAt...)),
|
||||
}
|
||||
|
||||
for _, key := range keys {
|
||||
@@ -300,7 +333,10 @@ func mergeBucket(ctx context.Context, tx *sql.Tx, key bucketKey, group *pointGro
|
||||
return res, err
|
||||
}
|
||||
|
||||
merged, overwrites, incomparable := mergePoints(stored.Points, group.points)
|
||||
merged, overwrites, incomparable, err := mergePoints(ctx, stored.Points, group.points)
|
||||
if err != nil {
|
||||
return res, err
|
||||
}
|
||||
res.overwrites = overwrites
|
||||
res.incomparable = incomparable
|
||||
// Считаем сохранённые точки, а не присланные: точные повторы внутри
|
||||
@@ -356,7 +392,7 @@ func mergeBucket(ctx context.Context, tx *sql.Tx, key bucketKey, group *pointGro
|
||||
// не заговорит.
|
||||
//
|
||||
// Точки из объекта не удаляются никогда.
|
||||
func mergePoints(stored, incoming []Point) (merged []Point, overwrites, incomparable int) {
|
||||
func mergePoints(ctx context.Context, stored, incoming []Point) (merged []Point, overwrites, incomparable int, err error) {
|
||||
type coord struct {
|
||||
start int64
|
||||
end int64
|
||||
@@ -404,7 +440,10 @@ func mergePoints(stored, incoming []Point) (merged []Point, overwrites, incompar
|
||||
out := make([]Point, 0, len(order))
|
||||
for _, c := range order {
|
||||
cands := byCoord[c]
|
||||
winner, unrelated := resolve(cands)
|
||||
winner, unrelated, err := resolve(ctx, cands)
|
||||
if err != nil {
|
||||
return nil, 0, 0, err
|
||||
}
|
||||
// Перезаписей столько, сколько точек уступило: при двух кандидатах
|
||||
// одна, при трёх две. Так счёт остаётся сравнимым с прежним, где
|
||||
// столкновение считалось на каждую приехавшую точку.
|
||||
@@ -425,7 +464,7 @@ func mergePoints(stored, incoming []Point) (merged []Point, overwrites, incompar
|
||||
}
|
||||
return out[i].End.Before(out[j].End)
|
||||
})
|
||||
return out, overwrites, incomparable
|
||||
return out, overwrites, incomparable, nil
|
||||
}
|
||||
|
||||
// candidate — точка вместе с тем, что о ней нужно знать при выборе
|
||||
@@ -444,14 +483,11 @@ func newCandidate(p Point) candidate {
|
||||
|
||||
// resolve выбирает победителя среди кандидатов одной координаты.
|
||||
//
|
||||
// Победитель — функция МНОЖЕСТВА кандидатов, а не порядка их поступления.
|
||||
// Сперва отбрасываются те, кого превосходит по полноте кто-то другой
|
||||
// (полнота — частичный порядок, поэтому «непревзойдённые» определены
|
||||
// однозначно), затем среди оставшихся берётся минимум по каноническому
|
||||
// порядку — он тотальный, поэтому минимум единственен. Обе операции зависят
|
||||
// только от состава множества, поэтому пересборка журнала даёт то же
|
||||
// состояние, что живой приём, а повторная свёртка той же доставки не меняет
|
||||
// ничего.
|
||||
// Механизм общий с выбором версии сущности — pickBest: отбрасываем
|
||||
// превзойдённых по частичному порядку, среди оставшихся берём минимум по
|
||||
// тотальному. Отношения разные (полнота у точек, покрытие у сущностей), а
|
||||
// рассуждение одно, и второй его экземпляр однажды уже разошёлся со стандартом
|
||||
// нетранзитивностью.
|
||||
//
|
||||
// Победителем остаётся одна из пришедших точек ДОСЛОВНО: правило выбирает, а
|
||||
// не конструирует. Каноническая форма существует только в момент сравнения, и
|
||||
@@ -461,46 +497,36 @@ func newCandidate(p Point) candidate {
|
||||
// несравнимыми наборами содержательных полей. На живом потоке этого не
|
||||
// случилось ни разу (0 из 2 897 столкновений), поэтому объединение полей не
|
||||
// реализовано: вместо него счётчик, который скажет, если событие наступит.
|
||||
func resolve(cands []candidate) (Point, bool) {
|
||||
if len(cands) == 1 {
|
||||
return cands[0].pt, false
|
||||
}
|
||||
|
||||
maximal := make([]candidate, 0, len(cands))
|
||||
for i, a := range cands {
|
||||
dominated := false
|
||||
for j, b := range cands {
|
||||
if i == j {
|
||||
continue
|
||||
}
|
||||
if b.fields.Relate(a.fields) == canon.FullnessSuperset {
|
||||
dominated = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !dominated {
|
||||
maximal = append(maximal, a)
|
||||
}
|
||||
}
|
||||
|
||||
best := maximal[0]
|
||||
for _, c := range maximal[1:] {
|
||||
if bytes.Compare(c.key, best.key) < 0 {
|
||||
best = c
|
||||
}
|
||||
func resolve(ctx context.Context, cands []candidate) (Point, bool, error) {
|
||||
winner, maximal, err := pickBest(ctx, cands, pointDominates, pointLess)
|
||||
if err != nil {
|
||||
return Point{}, false, err
|
||||
}
|
||||
|
||||
// Несравнимость — не «осталось больше одного»: точки с одинаковыми
|
||||
// наборами полей и разными значениями тоже остаются обе, и это рядовой
|
||||
// тай-брейк. Считается только то, ради чего отложено объединение полей:
|
||||
// у каждой из двух есть содержательный ключ, которого нет у другой.
|
||||
return best.pt, hasIncomparablePair(maximal)
|
||||
return cands[winner].pt, hasIncomparablePair(cands, maximal), nil
|
||||
}
|
||||
|
||||
func hasIncomparablePair(cands []candidate) bool {
|
||||
for i := range cands {
|
||||
for j := i + 1; j < len(cands); j++ {
|
||||
if cands[i].fields.Relate(cands[j].fields) == canon.FullnessIncomparable {
|
||||
// pointDominates — строгое превосходство по полноте. Relate возвращает
|
||||
// FullnessSuperset только когда a несёт всё, что b, и сверх того, поэтому
|
||||
// отношение уже строгое.
|
||||
func pointDominates(a, b candidate) bool {
|
||||
return a.fields.Relate(b.fields) == canon.FullnessSuperset
|
||||
}
|
||||
|
||||
// pointLess — тотальный порядок по канонической форме. Минимум единствен:
|
||||
// кандидаты с равной формой схлопываются ещё при сборе множества.
|
||||
func pointLess(a, b candidate) bool {
|
||||
return bytes.Compare(a.key, b.key) < 0
|
||||
}
|
||||
|
||||
func hasIncomparablePair(cands []candidate, maximal []int) bool {
|
||||
for i := range maximal {
|
||||
for j := i + 1; j < len(maximal); j++ {
|
||||
if cands[maximal[i]].fields.Relate(cands[maximal[j]].fields) == canon.FullnessIncomparable {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
@@ -51,6 +51,20 @@ type Delivery struct {
|
||||
// имён. Ответ на вопрос «что останется потерянным, если тело удалить»:
|
||||
// ретеншен обязан спрашивать его прежде, чем срезать тело.
|
||||
UncoveredSections string
|
||||
// SkippedEntities — сколько сущностей с собственным `id` разбор пропустил.
|
||||
// Второй половина ответа на тот же вопрос: сущность, которую разбор не
|
||||
// понял, в витрину не попала, а список непокрытых секций про неё молчит.
|
||||
//
|
||||
// Отсутствие значения означает «не измерялось» и НЕ равно нулю: так
|
||||
// выглядят доставки, свёрнутые разбором, который пропусков не считал, и те,
|
||||
// чей разбор не досчитал. Читатель, принимающий по счётчику необратимое
|
||||
// решение, обязан трактовать отсутствие как «не удалять».
|
||||
//
|
||||
// Указателем, а не sql.NullInt64: поле уедет в JSON `/stats` и в MCP, а
|
||||
// NullInt64 сериализуется формой драйвера (`{"Int64":0,"Valid":false}`) —
|
||||
// первый, кто про это забудет, опубликует её наружу, и она станет
|
||||
// контрактом. Указатель даёт `null` бесплатно и означает ровно то же.
|
||||
SkippedEntities *int64
|
||||
}
|
||||
|
||||
// CreateDelivery записывает факт приёма пакета.
|
||||
@@ -92,21 +106,26 @@ 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, uncovered_sections
|
||||
points, headers, uncovered_sections, skipped_entities
|
||||
FROM delivery ORDER BY received_at DESC, id DESC LIMIT 1`
|
||||
|
||||
var d Delivery
|
||||
var receivedAt string
|
||||
// sql.NullInt64 живёт ровно на границе сканирования и наружу не выходит.
|
||||
var skipped sql.NullInt64
|
||||
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.UncoveredSections)
|
||||
&d.ParseStatus, &d.Points, &d.Headers, &d.UncoveredSections, &skipped)
|
||||
if errors.Is(err, sql.ErrNoRows) {
|
||||
return Delivery{}, ErrNotFound
|
||||
}
|
||||
if err != nil {
|
||||
return Delivery{}, fmt.Errorf("select last delivery: %w", err)
|
||||
}
|
||||
if skipped.Valid {
|
||||
d.SkippedEntities = &skipped.Int64
|
||||
}
|
||||
|
||||
d.ReceivedAt, err = ParseTime(receivedAt)
|
||||
if err != nil {
|
||||
@@ -120,11 +139,19 @@ func (s *Store) LastDelivery(ctx context.Context) (Delivery, error) {
|
||||
//
|
||||
// Отдаются только **факты журнала**: то, что пришло вместе с доставкой.
|
||||
// Производные от разбора поля (`parse_status`, `points`, `derived_layer`,
|
||||
// `uncovered_sections`) сюда не попадают намеренно — перенос их в пересобранную
|
||||
// базу сделал бы витрину функцией предыдущего прогона. Особенно `derived_layer`:
|
||||
// доставка, чей повторный разбор отказал, отдала бы в наследование слой
|
||||
// прежнего разбора, и следующая доставка той же автоматизации унаследовала бы
|
||||
// его молча.
|
||||
// `uncovered_sections`, `skipped_entities`) сюда не попадают намеренно —
|
||||
// перенос их в пересобранную базу сделал бы витрину функцией предыдущего
|
||||
// прогона. Особенно `derived_layer`: доставка, чей повторный разбор отказал,
|
||||
// отдала бы в наследование слой прежнего разбора, и следующая доставка той же
|
||||
// автоматизации унаследовала бы его молча. У `skipped_entities` цена та же и
|
||||
// хуже: пустота у него значит «не измерялось», и перенесённое число выдавало бы
|
||||
// измерение прежнего разбора за измерение текущего — а по нему принимается
|
||||
// необратимое решение об удалении тела.
|
||||
//
|
||||
// Перечень пополняется ТЕМ ЖЕ изменением, которое заводит новое поле: он
|
||||
// единственное место, где сказано, чему нельзя пережить пересборку, и следующий
|
||||
// автор решает по нему. Поле, не внесённое сюда, однажды перенесут «для полноты
|
||||
// учёта».
|
||||
func (s *Store) ListDeliveries(ctx context.Context) ([]Delivery, error) {
|
||||
const q = `
|
||||
SELECT id, received_at, automation_name, automation_id, aggregation,
|
||||
@@ -255,6 +282,17 @@ type ParseOutcome struct {
|
||||
// у слоя пустота — отсутствие знания, у списка — знание об отсутствии.
|
||||
// Пересвёртка доставки, чья секция стала покрытой, обязана список очистить.
|
||||
Uncovered []string
|
||||
// SkippedEntities — сколько сущностей с собственным `id` разбор пропустил.
|
||||
// Пишется, когда разбор ДОСЧИТАЛ, включая ноль: доставка, пропуски которой
|
||||
// исчезли вместе с поумневшим разбором, не должна остаться помеченной
|
||||
// навсегда.
|
||||
//
|
||||
// nil означает «не измерялось» и колонку НЕ ТРОГАЕТ — та же идиома, что у
|
||||
// пустого Layer. Без неё отказ на чтении тела или паника разбора писали бы
|
||||
// ноль, то есть «проверено, терять нечего», в доставку, содержимое которой
|
||||
// никто не смотрел: ровно та подстановка, ради отказа от которой колонка
|
||||
// заведена без DEFAULT.
|
||||
SkippedEntities *int64
|
||||
}
|
||||
|
||||
func (s *Store) FinishParse(ctx context.Context, id string, out ParseOutcome) error {
|
||||
@@ -262,7 +300,8 @@ func (s *Store) FinishParse(ctx context.Context, id string, out ParseOutcome) er
|
||||
UPDATE delivery
|
||||
SET parse_status = ?, points = ?,
|
||||
derived_layer = CASE WHEN ? = '' THEN derived_layer ELSE ? END,
|
||||
uncovered_sections = ?
|
||||
uncovered_sections = ?,
|
||||
skipped_entities = CASE WHEN ? THEN ? ELSE skipped_entities END
|
||||
WHERE id = ?`
|
||||
|
||||
// Ровно одно представление пустоты — `[]`: nil-срез Go сериализуется как
|
||||
@@ -276,7 +315,17 @@ func (s *Store) FinishParse(ctx context.Context, id string, out ParseOutcome) er
|
||||
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)
|
||||
// Отсутствие числа не пишется нулём: ноль означает «измерено, пропусков не
|
||||
// было», а нам нужно «не измерялось». Колонка остаётся какой была — та же
|
||||
// форма, что у слоя строкой выше.
|
||||
measured := out.SkippedEntities != nil
|
||||
var skipped int64
|
||||
if measured {
|
||||
skipped = *out.SkippedEntities
|
||||
}
|
||||
|
||||
res, err := s.db.ExecContext(ctx, q, out.Status, out.Points, out.Layer, out.Layer,
|
||||
string(encoded), measured, skipped, id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update parse status: %w", err)
|
||||
}
|
||||
|
||||
+259
-93
@@ -106,21 +106,23 @@ func clipRefs(refs []EntityRef) []EntityRef {
|
||||
// entityVersion — версия сущности вместе с тем, что нужно знать при выборе
|
||||
// победителя.
|
||||
//
|
||||
// Разбор и канонизация ОТЛОЖЕНЫ: они нужны только когда хеш разошёлся с
|
||||
// сохранённым, то есть на одной доставке из сорока четырёх. Считать их сразу
|
||||
// значило бы разворачивать маршрут (95% веса тренировки, до мегабайта) в дерево
|
||||
// значений на каждой копии — ровно та форма, от которой разбор тела отказался
|
||||
// замером (197 МиБ кучи против 54 МиБ на теле 42 МиБ). Хеш при этом считается
|
||||
// сразу и один раз на доставку: он и есть быстрый путь.
|
||||
// Всё считается СРАЗУ и один раз на версию, до входа в транзакцию. Ленивость
|
||||
// здесь была мнимой: хеш всё равно требует полной канонической формы, то есть
|
||||
// самая дорогая работа платилась на каждой копии и так, а отложенный разбор
|
||||
// считал ту же форму ВТОРОЙ раз — и делал это внутри транзакции, которая
|
||||
// открыта `immediate` и повторяется до пяти раз при занятости базы.
|
||||
//
|
||||
// Баланс назван честно: на пути разошедшегося хеша (одна доставка из сорока
|
||||
// четырёх) стало на одну полную канонизацию меньше; на пути совпавшего хеша
|
||||
// добавился мелкий разбор в map[string]json.RawMessage — проход по телу без
|
||||
// разворачивания значений. Внутри транзакции для приехавших версий не остаётся
|
||||
// ничего.
|
||||
type entityVersion struct {
|
||||
raw json.RawMessage
|
||||
hash string
|
||||
from DeliveryRef
|
||||
|
||||
// key и fields заполняются лениво, методом analyze().
|
||||
raw json.RawMessage
|
||||
hash string
|
||||
from DeliveryRef
|
||||
key []byte
|
||||
fields canon.Fields
|
||||
parsed bool
|
||||
|
||||
// head — заголовок, который пишется колонками. У сохранённой версии он не
|
||||
// нужен: она либо побеждает и остаётся как есть, либо замещается целиком.
|
||||
@@ -128,21 +130,56 @@ type entityVersion struct {
|
||||
}
|
||||
|
||||
func newEntityVersion(e IncomingEntity, from DeliveryRef) (entityVersion, error) {
|
||||
h, err := canon.Hash(e.Raw)
|
||||
v, err := analyzeVersion(e.Raw, from)
|
||||
if err != nil {
|
||||
return entityVersion{}, fmt.Errorf("хеш сущности: %w", err)
|
||||
return entityVersion{}, err
|
||||
}
|
||||
return entityVersion{raw: e.Raw, hash: h, from: from, head: e}, nil
|
||||
v.head = e
|
||||
return v, nil
|
||||
}
|
||||
|
||||
// analyze разбирает версию, если этого ещё не делали.
|
||||
func (v *entityVersion) analyze() {
|
||||
if v.parsed {
|
||||
return
|
||||
// newStoredVersion собирает версию, прочитанную из витрины.
|
||||
//
|
||||
// Каноническая форма здесь НЕ считается, и это существенно: разбор сохранённой
|
||||
// версии — единственная работа, которая осталась внутри транзакции, открытой
|
||||
// `immediate`. Замер на тренировке в 168 КБ: полная канонизация с хешем — 4.5 мс
|
||||
// и 2.3 МБ на 38 тысячах аллокаций, множества ключей — 1.3 мс и 174 КБ на
|
||||
// тридцати. Хеш сохранённой уже лежит колонкой, а форма нужна ровно в одной
|
||||
// ветке тай-брейка (равные позиции журнала — та же доставка, свёрнутая
|
||||
// повторно) и считается там лениво.
|
||||
func newStoredVersion(raw json.RawMessage, hash string, from DeliveryRef) entityVersion {
|
||||
return entityVersion{
|
||||
raw: raw,
|
||||
hash: hash,
|
||||
from: from,
|
||||
fields: canon.Analyze(raw),
|
||||
}
|
||||
v.key = canon.SortKey(v.raw)
|
||||
v.fields = canon.Analyze(v.raw)
|
||||
v.parsed = true
|
||||
}
|
||||
|
||||
func analyzeVersion(raw json.RawMessage, from DeliveryRef) (entityVersion, error) {
|
||||
form, hash, err := canon.FormAndHash(raw)
|
||||
if err != nil {
|
||||
return entityVersion{}, fmt.Errorf("канонизация сущности: %w", err)
|
||||
}
|
||||
return entityVersion{
|
||||
raw: raw,
|
||||
hash: hash,
|
||||
from: from,
|
||||
key: form,
|
||||
fields: canon.Analyze(raw),
|
||||
}, nil
|
||||
}
|
||||
|
||||
// sortKey отдаёт каноническую форму версии, считая её при необходимости.
|
||||
//
|
||||
// Ленивость здесь одна на весь файл и нужна ровно сохранённой версии: у неё
|
||||
// форма требуется только в тай-брейке равных позиций журнала, а стоит она
|
||||
// втрое дороже разбора и платится под блокировкой записи.
|
||||
func (v *entityVersion) sortKey() []byte {
|
||||
if v.key == nil {
|
||||
v.key = canon.SortKey(v.raw)
|
||||
}
|
||||
return v.key
|
||||
}
|
||||
|
||||
// pickEntity выбирает между сохранённой и приехавшей версией.
|
||||
@@ -177,35 +214,23 @@ func pickEntity(stored, incoming *entityVersion) (takeIncoming, lost bool) {
|
||||
// Объединение полей отвергнуто там же и по той же причине, что для
|
||||
// точек, — на живом потоке событие не наступало ни разу, — а из двух
|
||||
// версий остаётся сохранённая: правило называется «не теряет
|
||||
// содержания», и приехавшая его теряет. Исход при этом остаётся
|
||||
// функцией журнала: доставки проигрываются в его порядке.
|
||||
// содержания», и приехавшая его теряет.
|
||||
//
|
||||
// ЗДЕСЬ И ТОЛЬКО ЗДЕСЬ исход зависит от порядка свёртки, а не от
|
||||
// журнала: в витрине лежит победитель прошлых слияний, а не все
|
||||
// кандидаты истории, и «сохранённая выигрывает» означает разный итог
|
||||
// при разном порядке. Порядок свёртки журналу не равен — доставка,
|
||||
// получившая ErrBusy, остаётся `pending` и сворачивается следующим
|
||||
// проходом, — так что живой приём и пересборка на несравнимых версиях
|
||||
// законно расходятся. Это единственная точка, где витрина не является
|
||||
// функцией множества доставок; она названа вслух в architecture.md, и
|
||||
// счётчик удержаний ниже — единственное, что о ней сообщает.
|
||||
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
|
||||
@@ -221,9 +246,6 @@ const (
|
||||
)
|
||||
|
||||
func compareEntities(stored, incoming *entityVersion) entityVerdict {
|
||||
stored.analyze()
|
||||
incoming.analyze()
|
||||
|
||||
storedCovers := stored.fields.Covers(incoming.fields)
|
||||
incomingCovers := incoming.fields.Covers(stored.fields)
|
||||
|
||||
@@ -250,58 +272,135 @@ func laterInJournal(stored, incoming *entityVersion) bool {
|
||||
if incoming.from.before(stored.from) {
|
||||
return false
|
||||
}
|
||||
return bytes.Compare(incoming.key, stored.key) < 0
|
||||
return bytes.Compare(incoming.sortKey(), stored.sortKey()) < 0
|
||||
}
|
||||
|
||||
// dedupeEntities сворачивает версии одного ключа ВНУТРИ доставки тем же
|
||||
// правилом — до сравнения с сохранённой.
|
||||
// entityDominates говорит, СТРОГО ли a превосходит b по содержанию: покрывает и
|
||||
// не покрывается в ответ.
|
||||
//
|
||||
// Без этого исход зависел бы от того, как написан цикл: карта по ключу дала бы
|
||||
// победу последнему элементу массива мимо правила полноты, а порядок элементов
|
||||
// в JSON-массиве нестабилен.
|
||||
//
|
||||
// Счётчик здесь считает СИММЕТРИЧНО — «в одном теле приехали две версии одного
|
||||
// ключа с разным содержанием», — а не «приехавшая обеднена». Внутри доставки
|
||||
// «сохранённой» версии не существует, есть только порядок элементов массива, и
|
||||
// счётчик, зависящий от него, наблюдал бы событие через раз.
|
||||
func dedupeEntities(versions []entityVersion) ([]entityVersion, int, []EntityRef) {
|
||||
type slot struct {
|
||||
v entityVersion
|
||||
pos int
|
||||
}
|
||||
// Строгость обязательна. Covers — предпорядок, а не строгий порядок: две версии
|
||||
// могут покрывать друг друга взаимно (тот же набор ключей, другие значения), и
|
||||
// отбрасывание «всего, что кем-то покрыто» опустошило бы множество, потеряв обе.
|
||||
func entityDominates(a, b entityVersion) bool {
|
||||
return a.fields.Covers(b.fields) && !b.fields.Covers(a.fields)
|
||||
}
|
||||
|
||||
byKey := make(map[EntityRef]slot, len(versions))
|
||||
// entityLess — тотальный порядок на версиях равного содержания.
|
||||
//
|
||||
// Сперва каноническая форма, потом ИСХОДНЫЕ БАЙТЫ. Второй разряд не украшение:
|
||||
// у сущностей версии с равной формой не схлопываются (в отличие от точек, где
|
||||
// это делает дедупликация по ключу), а у HAE порядок ключей в JSON и запись
|
||||
// числа нестабильны — то есть без него минимум неединствен, и в витрину лёг бы
|
||||
// тот элемент, что стоял в массиве раньше. Порядок элементов на проводе не
|
||||
// имеет права решать, какие байты хранятся.
|
||||
func entityLess(a, b entityVersion) bool {
|
||||
if c := bytes.Compare(a.key, b.key); c != 0 {
|
||||
return c < 0
|
||||
}
|
||||
return bytes.Compare(a.raw, b.raw) < 0
|
||||
}
|
||||
|
||||
// dedupeEntities сворачивает версии одного ключа ВНУТРИ доставки — до сравнения
|
||||
// с сохранённой.
|
||||
//
|
||||
// Победитель здесь — функция МНОЖЕСТВА версий, а не порядка элементов массива:
|
||||
// сперва отбрасываются строго покрытые, среди оставшихся берётся минимум
|
||||
// тотального порядка. Попарная свёртка была неверна ровно так же, как она была
|
||||
// неверна для точек: покрытие — частичный порядок, тай-брейк — тотальный, и
|
||||
// вместе они дают нетранзитивную победу, при которой [A,B,C] и [B,C,A] дают
|
||||
// разных победителей.
|
||||
//
|
||||
// Версии с СОВПАВШЕЙ канонической формой схлопываются ДО выбора победителя, и
|
||||
// это не оптимизация ради красоты: выбор квадратичен по числу кандидатов, а их
|
||||
// число приходит из чужого тела. Точки схлопываются так же и в том же месте
|
||||
// (см. mergePoints). Внутри схлопнутой группы остаются минимальные байты —
|
||||
// тот же второй разряд тотального порядка, что и между группами.
|
||||
//
|
||||
// Второй возврат — счётчик «в одном теле приехали версии одного ключа с РАЗНЫМ
|
||||
// содержанием», симметричный по построению: считаются кандидаты, чья форма
|
||||
// отличается от формы победителя. По форме, а не по байтам: порядок ключей у
|
||||
// HAE нестабилен и дребезг последнего разряда тоже, так что побайтовый счётчик
|
||||
// срабатывал бы на норме потока и стал бы неотличим от шума ровно тогда, когда
|
||||
// понадобился бы.
|
||||
//
|
||||
// Счётчик отдельный от «удержаний», а не общий с ними. Две версии в одном теле —
|
||||
// это НЕ потеря содержания: победитель ложится в витрину целиком, и удерживать
|
||||
// нечего. Смешивать их значило бы отвечать одним числом на два вопроса, которые
|
||||
// лечатся по-разному, — а на число удержаний опирается единственный контроль
|
||||
// того, что правило покрытия не стало слишком строгим.
|
||||
func dedupeEntities(ctx context.Context, versions []entityVersion) ([]entityVersion, int, []EntityRef, error) {
|
||||
byKey := make(map[EntityRef][]entityVersion, 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)}
|
||||
if _, seen := byKey[ref]; !seen {
|
||||
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}
|
||||
byKey[ref] = append(byKey[ref], v)
|
||||
}
|
||||
|
||||
out := make([]entityVersion, 0, len(order))
|
||||
diverging := 0
|
||||
var divergingAt []EntityRef
|
||||
|
||||
for _, ref := range order {
|
||||
out = append(out, byKey[ref].v)
|
||||
cands, dropped := collapseEqualForms(byKey[ref])
|
||||
winner, _, err := pickBest(ctx, cands, entityDominates, entityLess)
|
||||
if err != nil {
|
||||
return nil, 0, nil, err
|
||||
}
|
||||
out = append(out, cands[winner])
|
||||
|
||||
// Схлопнутые копии победителя различием не считаются: их форма ему
|
||||
// равна. Считаются все прочие — и оставшиеся кандидаты, и те, что
|
||||
// схлопнулись в них.
|
||||
differing := 0
|
||||
for i := range cands {
|
||||
if i != winner {
|
||||
differing += 1 + dropped[i]
|
||||
}
|
||||
}
|
||||
if differing > 0 {
|
||||
diverging += differing
|
||||
if len(divergingAt) < maxEntityRefsReported {
|
||||
divergingAt = append(divergingAt, ref)
|
||||
}
|
||||
}
|
||||
}
|
||||
return out, held, heldAt
|
||||
return out, diverging, divergingAt, nil
|
||||
}
|
||||
|
||||
// collapseEqualForms схлопывает версии с одинаковой канонической формой в одну,
|
||||
// оставляя минимальные исходные байты. Второй возврат — сколько копий сложилось
|
||||
// в каждого оставшегося кандидата (нужно счётчику различий).
|
||||
//
|
||||
// Схлопывание обязательно, а не желательно: без него тело с двадцатью тысячами
|
||||
// повторов одного `id` даёт четыреста миллионов сравнений покрытия, каждое с
|
||||
// обходом массивов. Тело в пределах приёма такое вмещает.
|
||||
func collapseEqualForms(versions []entityVersion) ([]entityVersion, []int) {
|
||||
byForm := make(map[string]int, len(versions))
|
||||
out := make([]entityVersion, 0, len(versions))
|
||||
dropped := make([]int, 0, len(versions))
|
||||
|
||||
for _, v := range versions {
|
||||
form := string(v.key)
|
||||
i, seen := byForm[form]
|
||||
if !seen {
|
||||
byForm[form] = len(out)
|
||||
out = append(out, v)
|
||||
dropped = append(dropped, 0)
|
||||
continue
|
||||
}
|
||||
dropped[i]++
|
||||
if bytes.Compare(v.raw, out[i].raw) < 0 {
|
||||
// Байты решают внутри группы ровно так же, как между группами:
|
||||
// порядок элементов на проводе не имеет права выбирать содержимое.
|
||||
// Заголовок едет вместе с байтами — он от них производен.
|
||||
out[i] = v
|
||||
}
|
||||
}
|
||||
return out, dropped
|
||||
}
|
||||
|
||||
// mergeEntities сливает сущности одной секции с сохранёнными.
|
||||
@@ -320,11 +419,31 @@ func mergeEntities(ctx context.Context, tx *sql.Tx, table string, versions []ent
|
||||
continue
|
||||
}
|
||||
|
||||
// Хеш — детектор изменений: совпал, значит писать нечего, и содержимое
|
||||
// сохранённой сущности читать не приходится вовсе. Тренировка
|
||||
// переприсылается каждой доставкой, пока не доедет маршрут, — на живом
|
||||
// архиве 44 копии дают три различных содержимых.
|
||||
// Хеш — детектор изменений: совпал, значит содержимое то же, и читать
|
||||
// его не приходится вовсе. Тренировка переприсылается каждой доставкой,
|
||||
// пока не доедет маршрут, — на живом архиве 44 копии дают три различных
|
||||
// содержимых.
|
||||
//
|
||||
// Но провенанс при этом обновить НАДО. Сохранённая позиция журнала
|
||||
// участвует в тай-брейке «содержание равно», и если в ней осталась
|
||||
// первая свёрнутая копия вместо победителя журнала, отложенная доставка
|
||||
// вернёт витрину к прежнему содержимому — то есть живая витрина
|
||||
// разойдётся с пересборкой, молча и в содержимом тренировки.
|
||||
//
|
||||
// Предел назван вслух: обновляется провенанс, но НЕ байты. При
|
||||
// совпавшей канонической форме в витрине остаются байты той доставки,
|
||||
// что свернулась первой, — а порядок ключей у HAE нестабилен, значит у
|
||||
// живого приёма и пересборки они могут различаться. Отпечаток этого не
|
||||
// различает (он считает по канонической форме), содержания не теряется
|
||||
// ничего, а переписывать мегабайтный маршрут на каждой из двадцати
|
||||
// шести присылок ради выбора между эквивалентными литералами — цена
|
||||
// несоразмерная.
|
||||
if stored.hash == v.hash {
|
||||
if stored.from.before(v.from) {
|
||||
if err := touchEntityProvenance(ctx, tx, table, v); err != nil {
|
||||
return 0, 0, nil, err
|
||||
}
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
@@ -332,7 +451,7 @@ func mergeEntities(ctx context.Context, tx *sql.Tx, table string, versions []ent
|
||||
if err != nil {
|
||||
return 0, 0, nil, err
|
||||
}
|
||||
prev := entityVersion{raw: storedRaw, hash: stored.hash, from: stored.from}
|
||||
prev := newStoredVersion(storedRaw, stored.hash, stored.from)
|
||||
|
||||
takeIncoming, lost := pickEntity(&prev, &v)
|
||||
if lost {
|
||||
@@ -352,6 +471,44 @@ func mergeEntities(ctx context.Context, tx *sql.Tx, table string, versions []ent
|
||||
return written, held, heldAt, nil
|
||||
}
|
||||
|
||||
// touchEntityProvenance поднимает провенанс сущности до более поздней доставки
|
||||
// журнала, не трогая содержимое.
|
||||
//
|
||||
// `updated_at` НЕ двигается, и это отдельное решение, а не экономия. Тренировка
|
||||
// приезжает до двадцати шести раз; бамп метки на каждой сделал бы её меткой
|
||||
// касания строки, а не изменения содержимого, и потребитель запроса «что
|
||||
// изменилось с момента X» получил бы двадцать шесть ложных изменений,
|
||||
// неотличимых от настоящего досчёта. Провенанс несёт собственную метку —
|
||||
// времени приёма своей доставки, — и для тай-брейка её достаточно.
|
||||
//
|
||||
// Счётчик записанных сущностей такое обновление тоже не увеличивает: он считает
|
||||
// СОДЕРЖИМОЕ витрины, и сравнимость его с прежними замерами важнее учёта
|
||||
// обновлённой ссылки.
|
||||
func touchEntityProvenance(ctx context.Context, tx *sql.Tx, table string, v entityVersion) error {
|
||||
q := `UPDATE ` + table + ` SET delivery_id = ?, delivery_received_at = ?` + entityWhere(table)
|
||||
|
||||
args := append([]any{v.from.ID, FormatTime(v.from.ReceivedAt)},
|
||||
entityKeyArgs(table, v.head.Kind, v.head.ID)...)
|
||||
|
||||
res, err := tx.ExecContext(ctx, q, args...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("update %s provenance: %w", table, err)
|
||||
}
|
||||
// Строка гарантированно существует: её заголовок прочитан этой же
|
||||
// транзакцией десятью строками выше. Ноль означал бы, что ключ собран не
|
||||
// теми колонками, — а провенанс в отпечаток витрины не входит, значит
|
||||
// молчаливый промах не поймает ни один оракул сходимости. Соседи по файлу
|
||||
// (FinishParse, MarkSealed) проверяют по той же причине.
|
||||
n, err := res.RowsAffected()
|
||||
if err != nil {
|
||||
return fmt.Errorf("update %s provenance: %w", table, err)
|
||||
}
|
||||
if n == 0 {
|
||||
return fmt.Errorf("update %s provenance: %w", table, ErrNotFound)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
type storedEntityHead struct {
|
||||
hash string
|
||||
from DeliveryRef
|
||||
@@ -406,11 +563,20 @@ func entityWhere(table string) string {
|
||||
return ` WHERE id = ?`
|
||||
}
|
||||
|
||||
func queryEntity(ctx context.Context, tx *sql.Tx, q, table, kind, id string) *sql.Row {
|
||||
// entityKeyArgs — аргументы к entityWhere. Живут рядом с ним намеренно: число
|
||||
// `?` в тексте и длина этого среза обязаны меняться вместе, а компилятор их
|
||||
// соответствия не видит. Промах даст ошибку SQLite внутри транзакции слияния,
|
||||
// то есть на пути, который повторяется до пяти раз и оканчивается `failed` у
|
||||
// доставки, а не отказом сборки.
|
||||
func entityKeyArgs(table, kind, id string) []any {
|
||||
if table == recordTable {
|
||||
return tx.QueryRowContext(ctx, q, kind, id)
|
||||
return []any{kind, id}
|
||||
}
|
||||
return tx.QueryRowContext(ctx, q, id)
|
||||
return []any{id}
|
||||
}
|
||||
|
||||
func queryEntity(ctx context.Context, tx *sql.Tx, q, table, kind, id string) *sql.Row {
|
||||
return tx.QueryRowContext(ctx, q, entityKeyArgs(table, kind, id)...)
|
||||
}
|
||||
|
||||
func writeEntity(ctx context.Context, tx *sql.Tx, table string, v entityVersion, now time.Time) error {
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
package store
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Внутренний тест, потому что проверяемое наружу не отдаётся: `updated_at` —
|
||||
// колонка, а не поле модели. Обещание «метка означает изменение содержимого, а
|
||||
// не касание строки» держится только этой проверкой, и внешний тест для неё
|
||||
// потребовал бы публичного метода ради теста.
|
||||
func TestMergeПовторНеДвигаетМеткуИзменения(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
st, err := Open(filepath.Join(t.TempDir(), "healthlog.db"))
|
||||
if err != nil {
|
||||
t.Fatalf("открытие базы: %v", err)
|
||||
}
|
||||
t.Cleanup(func() { _ = st.Close() })
|
||||
|
||||
raw := json.RawMessage(`{"id":"w1","route":[{"lat":1},{"lat":2}],"totalEnergy":{"qty":20}}`)
|
||||
wo := IncomingEntity{
|
||||
ID: "w1",
|
||||
Kind: "workouts",
|
||||
Start: time.Date(2025, 6, 5, 7, 0, 0, 0, time.UTC),
|
||||
End: time.Date(2025, 6, 5, 7, 10, 0, 0, time.UTC),
|
||||
Raw: raw,
|
||||
}
|
||||
merge := func(id string, at time.Time) MergeStats {
|
||||
t.Helper()
|
||||
s, err := st.Merge(ctx, Incoming{Workouts: []IncomingEntity{wo}},
|
||||
DeliveryRef{ID: id, ReceivedAt: at})
|
||||
if err != nil {
|
||||
t.Fatalf("слияние: %v", err)
|
||||
}
|
||||
return s
|
||||
}
|
||||
updatedAt := func() string {
|
||||
t.Helper()
|
||||
var v string
|
||||
if err := st.db.QueryRowContext(ctx,
|
||||
`SELECT updated_at FROM workout WHERE id = 'w1'`).Scan(&v); err != nil {
|
||||
t.Fatalf("чтение updated_at: %v", err)
|
||||
}
|
||||
return v
|
||||
}
|
||||
|
||||
base := time.Date(2025, 6, 5, 8, 0, 0, 0, time.UTC)
|
||||
merge("d1", base)
|
||||
before := updatedAt()
|
||||
|
||||
// Тренировка приезжает до 26 раз, пока источник её досчитывает. Провенанс
|
||||
// при этом обязан подняться, а метка изменения — нет: иначе она становится
|
||||
// меткой касания строки, и запрос «что изменилось с момента X» получает 26
|
||||
// ложных изменений, неотличимых от настоящего досчёта.
|
||||
merge("d2", base.Add(5*time.Minute))
|
||||
if after := updatedAt(); after != before {
|
||||
t.Errorf("метка изменения двинулась без изменения содержимого: %s → %s", before, after)
|
||||
}
|
||||
|
||||
w, err := st.Workout(ctx, "w1")
|
||||
if err != nil {
|
||||
t.Fatalf("чтение тренировки: %v", err)
|
||||
}
|
||||
if w.Delivery != "d2" {
|
||||
t.Fatal("провенанс не обновился — тест проверяет не то")
|
||||
}
|
||||
}
|
||||
@@ -3,11 +3,25 @@ package store_test
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.vakhrushev.me/av/healthlog/internal/store"
|
||||
)
|
||||
|
||||
func fingerprint(t *testing.T, st *store.Store) string {
|
||||
t.Helper()
|
||||
|
||||
fp, err := st.Fingerprint(context.Background())
|
||||
if err != nil {
|
||||
t.Fatalf("отпечаток: %v", err)
|
||||
}
|
||||
return fp
|
||||
}
|
||||
|
||||
// workout собирает тренировку с заданным содержимым. Заголовок в этих тестах
|
||||
// вторичен: правило замены смотрит на содержание, а не на колонки.
|
||||
func workout(t *testing.T, id, raw string) store.IncomingEntity {
|
||||
@@ -206,22 +220,32 @@ func TestMergeПовторТойЖеТренировкиНеПишет(t *testin
|
||||
}
|
||||
}
|
||||
|
||||
// Три версии в разных порядках подачи: пункты правила, не зависящие от порядка
|
||||
// Версии в разных порядках подачи: пункты правила, не зависящие от порядка
|
||||
// свёртки, обязаны давать одно состояние. Конвенция требует перестановки трёх,
|
||||
// а не пары: попарная свёртка уже давала нетранзитивную победу на точках.
|
||||
func TestMergeПерестановкаТрёхВерсийДаётОдноСостояние(t *testing.T) {
|
||||
// а не пары (попарная свёртка уже давала нетранзитивную победу на точках) и
|
||||
// версии с содержимым, равным одной из присланных, — иначе ветка «содержание
|
||||
// равно» не посещается ни разу.
|
||||
func TestMergeПерестановкаВерсийДаётОдноСостояние(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
type step struct {
|
||||
d store.DeliveryRef
|
||||
raw string
|
||||
}
|
||||
// Четвёртая версия несёт содержимое, РАВНОЕ одной из уже присланных. Без неё
|
||||
// перебор троек с разными хешами ветку «содержание равно» не посещает ни
|
||||
// разу — а именно на ней провенанс устаревал, и живая витрина расходилась с
|
||||
// пересборкой молча.
|
||||
steps := []step{
|
||||
{from(t, "d1", "2025-06-05T08:00:00Z"), woNoRouteNewValues},
|
||||
{from(t, "d2", "2025-06-05T08:05:00Z"), woWithRoute},
|
||||
{from(t, "d3", "2025-06-05T08:10:00Z"), woRicher},
|
||||
{from(t, "d4", "2025-06-05T08:15:00Z"), woWithRoute},
|
||||
}
|
||||
orders := [][]int{
|
||||
{0, 1, 2, 3}, {3, 2, 1, 0}, {1, 0, 3, 2},
|
||||
{2, 3, 0, 1}, {1, 3, 0, 2}, {3, 0, 2, 1},
|
||||
}
|
||||
orders := [][]int{{0, 1, 2}, {2, 1, 0}, {1, 0, 2}, {1, 2, 0}, {2, 0, 1}, {0, 2, 1}}
|
||||
|
||||
var want string
|
||||
for i, order := range orders {
|
||||
@@ -421,12 +445,18 @@ func TestMergeНесравнимыеВерсииВОдномТелеНеЗави
|
||||
if storedRaw(t, прямой, "w1") != storedRaw(t, обратный, "w1") {
|
||||
t.Error("исход зависит от порядка элементов в массиве секции")
|
||||
}
|
||||
if a.EntitiesHeld != b.EntitiesHeld {
|
||||
t.Errorf("счётчик зависит от порядка: %d против %d", a.EntitiesHeld, b.EntitiesHeld)
|
||||
if a.EntitiesDiverging != b.EntitiesDiverging {
|
||||
t.Errorf("счётчик зависит от порядка: %d против %d", a.EntitiesDiverging, b.EntitiesDiverging)
|
||||
}
|
||||
if a.EntitiesHeld == 0 {
|
||||
if a.EntitiesDiverging == 0 {
|
||||
t.Error("две версии с разным содержанием в одном теле остались незамеченными")
|
||||
}
|
||||
// Удержаний тут нет: победитель лёг в витрину целиком, терять нечего.
|
||||
// Счётчики разведены именно ради этого различия.
|
||||
if a.EntitiesHeld != 0 || b.EntitiesHeld != 0 {
|
||||
t.Errorf("несравнимые версии в одном теле посчитаны удержаниями: %d и %d",
|
||||
a.EntitiesHeld, b.EntitiesHeld)
|
||||
}
|
||||
}
|
||||
|
||||
// Отпечаток обязан различать состояния, а не только содержимое: составной ключ
|
||||
@@ -458,3 +488,375 @@ func TestFingerprintРазличаетСоставнойКлючЗаписи(t *
|
||||
t.Error("два разных состояния витрины дали один отпечаток: составной ключ склеен до взятия длины")
|
||||
}
|
||||
}
|
||||
|
||||
// Оракулы враждебного прохода ревью: тело контролирует отправитель целиком, и
|
||||
// «версия той же формы без содержания» проходила все проверки — а маршрут это
|
||||
// 95% тренировки, которого в экспорте Apple нет вовсе. Восстановить его после
|
||||
// затирания не может даже пересборка: журнал проиграет то же поражение.
|
||||
const (
|
||||
woRealWorkout = `{"id":"w7","name":"В помещении Ходьба","isIndoor":true,
|
||||
"maxHeartRate":{"qty":199,"units":"count/min"},
|
||||
"heartRate":{"max":{"qty":199},"avg":{"qty":47.2},"min":{"qty":41}},
|
||||
"heartRateData":[{"Max":199,"Avg":86.1,"Min":41},{"Max":150,"Avg":80.0,"Min":44}],
|
||||
"activeEnergy":[{"qty":49.4,"units":"kJ"},{"qty":12.1,"units":"kJ"}],
|
||||
"totalEnergy":{"qty":66.4,"units":"kJ"},"duration":11.1}`
|
||||
woSkeleton = `{"id":"w7","name":"x","isIndoor":false,
|
||||
"maxHeartRate":1,"heartRate":1,
|
||||
"heartRateData":[null,null],"activeEnergy":[null,null],
|
||||
"totalEnergy":1,"duration":1}`
|
||||
woNulledRoute = `{"id":"w1","name":"На улице Ходьба","route":[null,null,null],
|
||||
"activeEnergy":[{"qty":10}],"totalEnergy":{"qty":20,"units":"kJ"}}`
|
||||
woEmptyRoute = `{"id":"w1","name":"На улице Ходьба","route":[{},{},{}],
|
||||
"activeEnergy":[{"qty":10}],"totalEnergy":{"qty":20,"units":"kJ"}}`
|
||||
// Тот же смысл, другой порядок ключей и другая запись числа: каноническая
|
||||
// форма совпадает, байты — нет.
|
||||
woSameFormOtherBytes = `{"name":"На улице Ходьба","id":"w1","totalEnergy":{"units":"kJ","qty":20.0},
|
||||
"activeEnergy":[{"qty":10}],"route":[{"lat":1},{"lat":2},{"lat":3}]}`
|
||||
)
|
||||
|
||||
func TestMergeСкелетНеВытесняетТренировку(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st := open(t)
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w7", woRealWorkout))
|
||||
stats := mergeWorkouts(t, st, from(t, "d2", "2025-06-05T08:05:00Z"), workout(t, "w7", woSkeleton))
|
||||
|
||||
if got := storedRaw(t, st, "w7"); !strings.Contains(got, "199") {
|
||||
t.Error("тренировка заменена скелетом из скаляров и null")
|
||||
}
|
||||
if stats.EntitiesHeld != 1 {
|
||||
t.Errorf("удержано %d, ожидалось 1: событие обязано быть видно", stats.EntitiesHeld)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMergeМаршрутНеЗатираетсяПустышкамиТойЖеДлины(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
for name, poor := range map[string]string{"null": woNulledRoute, "пустые объекты": woEmptyRoute} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st := open(t)
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w1", woWithRoute))
|
||||
stats := mergeWorkouts(t, st, from(t, "d2", "2025-06-05T08:05:00Z"), workout(t, "w1", poor))
|
||||
|
||||
if got := storedRaw(t, st, "w1"); !strings.Contains(got, `"lat":1`) {
|
||||
t.Error("маршрут затёрт рядом той же длины без содержания")
|
||||
}
|
||||
if stats.EntitiesHeld != 1 {
|
||||
t.Errorf("удержано %d, ожидалось 1", stats.EntitiesHeld)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// Ключ с пустым значением исчезал бы по жребию тай-брейка. Второй разряд
|
||||
// сравнения записан для точек и здесь применяется к сущностям.
|
||||
func TestMergeКлючСПустымЗначениемНеИсчезает(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const (
|
||||
withEmpty = `{"id":"w1","qty":10,"context":null}`
|
||||
noEmpty = `{"id":"w1","qty":11}`
|
||||
)
|
||||
st := open(t)
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w1", withEmpty))
|
||||
stats := mergeWorkouts(t, st, from(t, "d2", "2025-06-05T08:05:00Z"), workout(t, "w1", noEmpty))
|
||||
|
||||
if got := storedRaw(t, st, "w1"); !strings.Contains(got, "context") {
|
||||
t.Error("ключ с пустым значением исчез по жребию")
|
||||
}
|
||||
if stats.EntitiesHeld != 1 {
|
||||
t.Errorf("удержано %d, ожидалось 1", stats.EntitiesHeld)
|
||||
}
|
||||
}
|
||||
|
||||
// Две версии одного `id` в ОДНОМ теле с разным содержанием: счётчик не имеет
|
||||
// права молчать. До правки он давал ноль — то есть событие проходило бы как
|
||||
// штатное INFO.
|
||||
func TestMergeДвеВерсииВОдномТелеСчитаются(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
d := from(t, "d1", "2025-06-05T08:00:00Z")
|
||||
st := open(t)
|
||||
stats := mergeWorkouts(t, st, d, workout(t, "w1", woWithRoute), workout(t, "w1", woNulledRoute))
|
||||
|
||||
if stats.EntitiesDiverging == 0 {
|
||||
t.Error("две версии одного id в одном теле — счётчик молчит")
|
||||
}
|
||||
if got := storedRaw(t, st, "w1"); !strings.Contains(got, `"lat":1`) {
|
||||
t.Error("в теле победила версия без содержания")
|
||||
}
|
||||
}
|
||||
|
||||
// Побайтовое различие при совпавшей канонической форме событием не является:
|
||||
// порядок ключей у HAE нестабилен и дребезг последнего разряда тоже. Но байты в
|
||||
// витрине обязаны быть одни при любой перестановке — иначе исход зависит от
|
||||
// порядка элементов на проводе.
|
||||
func TestMergeРавнаяФормаРазныеБайтыНеСобытие(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
d := from(t, "d1", "2025-06-05T08:00:00Z")
|
||||
|
||||
прямой := open(t)
|
||||
s1 := mergeWorkouts(t, прямой, d, workout(t, "w1", woWithRoute), workout(t, "w1", woSameFormOtherBytes))
|
||||
обратный := open(t)
|
||||
s2 := mergeWorkouts(t, обратный, d, workout(t, "w1", woSameFormOtherBytes), workout(t, "w1", woWithRoute))
|
||||
|
||||
if s1.EntitiesDiverging != 0 || s2.EntitiesDiverging != 0 {
|
||||
t.Errorf("счётчик сработал на дребезге записи: %d и %d",
|
||||
s1.EntitiesDiverging, s2.EntitiesDiverging)
|
||||
}
|
||||
if storedRaw(t, прямой, "w1") != storedRaw(t, обратный, "w1") {
|
||||
t.Error("в витрине разные байты при одинаковом содержимом: победитель зависит от порядка")
|
||||
}
|
||||
}
|
||||
|
||||
// Победитель внутри доставки — функция МНОЖЕСТВА версий. Попарная свёртка
|
||||
// частичного порядка с тотальным тай-брейком нетранзитивна: [A,B,C] давало C,
|
||||
// [B,C,A] давало A.
|
||||
func TestMergeПерестановкаТрёхВерсийВОдномТеле(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const (
|
||||
vA = `{"a":[9,9],"b":2,"id":"k"}`
|
||||
vB = `{"a":[1,2],"id":"k"}`
|
||||
vC = `{"a":[1],"c":3,"id":"k"}`
|
||||
)
|
||||
orders := [][]string{
|
||||
{vA, vB, vC}, {vA, vC, vB}, {vB, vA, vC},
|
||||
{vB, vC, vA}, {vC, vA, vB}, {vC, vB, vA},
|
||||
}
|
||||
|
||||
var want string
|
||||
for i, order := range orders {
|
||||
st := open(t)
|
||||
ws := make([]store.IncomingEntity, 0, len(order))
|
||||
for _, raw := range order {
|
||||
ws = append(ws, workout(t, "k", raw))
|
||||
}
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), ws...)
|
||||
|
||||
got := storedRaw(t, st, "k")
|
||||
if i == 0 {
|
||||
want = got
|
||||
continue
|
||||
}
|
||||
if got != want {
|
||||
t.Errorf("порядок %d дал другого победителя:\n %s\n %s", i, got, want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Провенанс обязан отражать победителя ЖУРНАЛА, а не первую свёрнутую копию.
|
||||
// Иначе доставка, свёрнутая с опозданием, вернёт витрину к прежнему содержимому,
|
||||
// и живая витрина разойдётся с пересборкой молча — в содержимом тренировки.
|
||||
func TestMergeОтложеннаяДоставкаНеВозвращаетПрежнееСодержимое(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
journal := open(t)
|
||||
mergeWorkouts(t, journal, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w1", woWithRoute))
|
||||
mergeWorkouts(t, journal, from(t, "d3", "2025-06-05T08:05:00Z"), workout(t, "w1", woSameShapeNewValues))
|
||||
mergeWorkouts(t, journal, from(t, "d2", "2025-06-05T08:10:00Z"), workout(t, "w1", woWithRoute))
|
||||
|
||||
deferred := open(t)
|
||||
mergeWorkouts(t, deferred, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w1", woWithRoute))
|
||||
mergeWorkouts(t, deferred, from(t, "d2", "2025-06-05T08:10:00Z"), workout(t, "w1", woWithRoute))
|
||||
mergeWorkouts(t, deferred, from(t, "d3", "2025-06-05T08:05:00Z"), workout(t, "w1", woSameShapeNewValues))
|
||||
|
||||
if a, b := storedRaw(t, journal, "w1"), storedRaw(t, deferred, "w1"); a != b {
|
||||
t.Errorf("состояние зависит от порядка свёртки:\n %s\n %s", a, b)
|
||||
}
|
||||
if a, b := fingerprint(t, journal), fingerprint(t, deferred); a != b {
|
||||
t.Errorf("отпечатки разошлись: %s против %s", a, b)
|
||||
}
|
||||
}
|
||||
|
||||
// Провенанс поднимается до более поздней доставки журнала даже при совпавшем
|
||||
// хеше — иначе тай-брейк «содержание равно» решает по устаревшей позиции.
|
||||
func TestMergeПовторОбновляетПровенанс(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
st := open(t)
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w1", woWithRoute))
|
||||
stats := mergeWorkouts(t, st, from(t, "d2", "2025-06-05T08:05:00Z"), workout(t, "w1", woWithRoute))
|
||||
|
||||
if stats.WorkoutsWritten != 0 {
|
||||
t.Errorf("записано %d, ожидалось 0: содержимое то же", stats.WorkoutsWritten)
|
||||
}
|
||||
w, err := st.Workout(ctx, "w1")
|
||||
if err != nil {
|
||||
t.Fatalf("чтение тренировки: %v", err)
|
||||
}
|
||||
if w.Delivery != "d2" {
|
||||
t.Errorf("провенанс %q, ожидался d2: победитель журнала — более поздняя доставка", w.Delivery)
|
||||
}
|
||||
}
|
||||
|
||||
// Провенанс записи обновляется тем же путём, что и у тренировки, но ключ у неё
|
||||
// СОСТАВНОЙ — а значит текст `WHERE` и хвост аргументов обязаны совпадать. Ветка
|
||||
// не покрывалась ни одним тестом, притом что промах давал бы `UPDATE` в ноль
|
||||
// строк, невидимый ни в отпечатке (провенанса там нет), ни в счётчиках.
|
||||
// Род `record` — единственный, где цена необратима: в экспорте Apple его нет.
|
||||
func TestMergeПровенансЗаписиОбновляетсяПоСоставномуКлючу(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
ctx := context.Background()
|
||||
st := open(t)
|
||||
rec := func() store.IncomingEntity {
|
||||
return store.IncomingEntity{
|
||||
ID: "e1", Kind: "stateOfMind",
|
||||
Start: ts(t, "2025-06-05T18:00:00Z"),
|
||||
End: ts(t, "2025-06-05T18:00:00Z"),
|
||||
Raw: json.RawMessage(`{"id":"e1","kind":"momentary_emotion","valence":0.5}`),
|
||||
}
|
||||
}
|
||||
merge := func(d store.DeliveryRef) store.MergeStats {
|
||||
t.Helper()
|
||||
s, err := st.Merge(ctx, store.Incoming{Records: []store.IncomingEntity{rec()}}, d)
|
||||
if err != nil {
|
||||
t.Fatalf("слияние записи: %v", err)
|
||||
}
|
||||
return s
|
||||
}
|
||||
|
||||
merge(from(t, "d1", "2025-06-05T20:00:00Z"))
|
||||
stats := merge(from(t, "d2", "2025-06-05T20:05:00Z"))
|
||||
if stats.RecordsWritten != 0 {
|
||||
t.Errorf("записано %d, ожидалось 0: содержимое то же", stats.RecordsWritten)
|
||||
}
|
||||
|
||||
got, err := st.Record(ctx, "stateOfMind", "e1")
|
||||
if err != nil {
|
||||
t.Fatalf("чтение записи: %v", err)
|
||||
}
|
||||
if got.Delivery != "d2" {
|
||||
t.Errorf("провенанс %q, ожидался d2", got.Delivery)
|
||||
}
|
||||
}
|
||||
|
||||
// Число версий одного ключа приходит из чужого тела, а выбор победителя по ним
|
||||
// квадратичен. Отмена обязана прерывать отбор: без неё тело с двадцатью
|
||||
// тысячами версий занимает единственного воркера свёртки дольше, чем длится его
|
||||
// собственный дедлайн, и очередь встаёт молча при зелёном `/healthz`.
|
||||
func TestMergeОтменаПрерываетВыборПобедителя(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st := open(t)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
cancel()
|
||||
|
||||
ws := make([]store.IncomingEntity, 0, 3)
|
||||
for i := range 3 {
|
||||
raw := fmt.Sprintf(`{"id":"w1","qty":%d,"route":[{"lat":%d}]}`, i, i)
|
||||
ws = append(ws, workout(t, "w1", raw))
|
||||
}
|
||||
_, err := st.Merge(ctx, store.Incoming{Workouts: ws}, from(t, "d1", "2025-06-05T08:00:00Z"))
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Errorf("слияние дало %v, ожидалась отмена", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Версии с совпавшей канонической формой схлопываются ДО квадратичного отбора:
|
||||
// иначе тело, вмещающее сотни тысяч копий одного `id`, стоит часов работы. При
|
||||
// этом схлопывание не имеет права менять исход — победитель тот же.
|
||||
func TestMergeРавныеФормыСхлопываютсяДоОтбора(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
st := open(t)
|
||||
ws := make([]store.IncomingEntity, 0, 2000)
|
||||
for range 2000 {
|
||||
ws = append(ws, workout(t, "w1", woWithRoute))
|
||||
}
|
||||
ws = append(ws, workout(t, "w1", woRicher))
|
||||
|
||||
done := make(chan store.MergeStats, 1)
|
||||
go func() {
|
||||
done <- mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), ws...)
|
||||
}()
|
||||
select {
|
||||
case stats := <-done:
|
||||
if got := storedRaw(t, st, "w1"); got != woRicher {
|
||||
t.Errorf("схлопывание изменило победителя:\n%s", got)
|
||||
}
|
||||
// Победитель — единственная более полная версия; от неё формой
|
||||
// отличаются две тысячи копий победнее, и схлопывание их не прячет:
|
||||
// счётчик считает версии, а не группы.
|
||||
if stats.EntitiesDiverging != 2000 {
|
||||
t.Errorf("различающихся версий %d, ожидалось 2000", stats.EntitiesDiverging)
|
||||
}
|
||||
case <-time.After(20 * time.Second):
|
||||
t.Fatal("слияние не уложилось в 20 с: схлопывание не работает")
|
||||
}
|
||||
}
|
||||
|
||||
// Регрессия, найденная враждебным проходом ревью: безусловный второй разряд
|
||||
// правила покрытия запирал законный досчёт НАВСЕГДА. Версия с пустым ключом и
|
||||
// без маршрута оказывалась несравнимой с версией, у которой маршрут приехал, а
|
||||
// этого ключа нет, — и маршрут не доезжал ни одной доставкой, причём пересборка
|
||||
// проигрывала то же поражение. Второй разряд разрешает спор равных, а не
|
||||
// отменяет первый.
|
||||
func TestMergeПустойКлючНеЗапираетДосчёт(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const (
|
||||
сПустымКлючом = `{"id":"w2","activeEnergy":[{"qty":10}],"totalEnergy":null}`
|
||||
сМаршрутом = `{"id":"w2","activeEnergy":[{"qty":10}],"route":[{"lat":1},{"lat":2},{"lat":3}]}`
|
||||
)
|
||||
st := open(t)
|
||||
mergeWorkouts(t, st, from(t, "d1", "2025-06-05T08:00:00Z"), workout(t, "w2", сПустымКлючом))
|
||||
stats := mergeWorkouts(t, st, from(t, "d2", "2025-06-05T08:05:00Z"), workout(t, "w2", сМаршрутом))
|
||||
|
||||
if got := storedRaw(t, st, "w2"); !strings.Contains(got, `"lat":1`) {
|
||||
t.Errorf("маршрут не доехал — досчёт заперт пустым ключом:\n%s", got)
|
||||
}
|
||||
if stats.WorkoutsWritten != 1 {
|
||||
t.Errorf("записано %d, ожидалась 1", stats.WorkoutsWritten)
|
||||
}
|
||||
if stats.EntitiesHeld != 0 {
|
||||
t.Errorf("удержано %d: законный досчёт принят за потерю содержания", stats.EntitiesHeld)
|
||||
}
|
||||
}
|
||||
|
||||
// Единственная точка, где витрина НЕ является функцией множества доставок, —
|
||||
// несравнимые версии. Тест не чинит это, а закрепляет: в витрине лежит
|
||||
// победитель прошлых слияний, а не все кандидаты истории, поэтому «сохранённая
|
||||
// выигрывает» даёт разный итог при разном порядке. Порядок свёртки журналу не
|
||||
// равен: доставка, получившая ErrBusy, остаётся `pending` и сворачивается
|
||||
// следующим проходом.
|
||||
//
|
||||
// Предел назван в pickEntity и в architecture.md; наблюдается счётчиком
|
||||
// удержаний. Если он когда-нибудь будет закрыт, красный тест напомнит, что
|
||||
// текст в обоих местах пора переписать.
|
||||
func TestMergeНесравнимыеВерсииЗависятОтПорядкаСвёртки(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
const (
|
||||
база = `{"id":"w3","activeEnergy":[{"qty":10}]}`
|
||||
сТреком = `{"id":"w3","activeEnergy":[{"qty":10}],"route":[{"lat":1},{"lat":2}]}`
|
||||
сЭтажами = `{"id":"w3","activeEnergy":[{"qty":10}],"flightsClimbed":{"qty":3}}`
|
||||
)
|
||||
d1 := from(t, "d1", "2025-06-05T08:00:00Z")
|
||||
d2 := from(t, "d2", "2025-06-05T08:05:00Z")
|
||||
d3 := from(t, "d3", "2025-06-05T08:10:00Z")
|
||||
|
||||
журнальный := open(t)
|
||||
mergeWorkouts(t, журнальный, d1, workout(t, "w3", база))
|
||||
mergeWorkouts(t, журнальный, d2, workout(t, "w3", сТреком))
|
||||
mergeWorkouts(t, журнальный, d3, workout(t, "w3", сЭтажами))
|
||||
|
||||
отложенный := open(t)
|
||||
mergeWorkouts(t, отложенный, d1, workout(t, "w3", база))
|
||||
mergeWorkouts(t, отложенный, d3, workout(t, "w3", сЭтажами))
|
||||
stats := mergeWorkouts(t, отложенный, d2, workout(t, "w3", сТреком))
|
||||
|
||||
a, b := storedRaw(t, журнальный, "w3"), storedRaw(t, отложенный, "w3")
|
||||
if a == b {
|
||||
t.Fatal("предел закрылся — перепиши текст в pickEntity и architecture.md")
|
||||
}
|
||||
// И главное: расхождение не молчит.
|
||||
if stats.EntitiesHeld == 0 {
|
||||
t.Error("расхождение по порядку свёртки не отражено счётчиком удержаний")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -10,6 +10,13 @@ import (
|
||||
// database/sql.
|
||||
var ErrNotFound = errors.New("запись не найдена")
|
||||
|
||||
// errNoCandidates — выбор победителя позван на пустом множестве. Нарушенный
|
||||
// инвариант вызывающего, а не свойство данных: множество собирается из карты и
|
||||
// пустым быть не может. Ошибкой, а не паникой, потому что путь проходит внутри
|
||||
// свёртки принятой доставки — отказ обязан быть диагностируемым, а не «index
|
||||
// out of range» в стеке фоновой горутины.
|
||||
var errNoCandidates = errors.New("выбор победителя на пустом множестве версий")
|
||||
|
||||
// ErrBusy — база занята, и повторы транзакции этого не пересидели.
|
||||
//
|
||||
// Доменная ошибка, а не код драйвера: на неё ветвится свёртка. Отказ по
|
||||
|
||||
@@ -0,0 +1,30 @@
|
||||
-- +goose Up
|
||||
-- Сколько сущностей с собственным `id` разбор этой доставки пропустил: без
|
||||
-- `id`, с непомерно длинным `id`, с неразбираемой меткой времени или не
|
||||
-- разобравшихся как объект.
|
||||
--
|
||||
-- Заводится не ради отчётности. Ретеншен сырого архива решает «что потеряется,
|
||||
-- если тело удалить», ПО БАЗЕ, и до этой колонки получал ответ «терять нечего»
|
||||
-- ровно там, где потеряна тренировка с маршрутом: сущность в витрину не попала,
|
||||
-- список непокрытых секций пуст, статус `parsed`. Лог здесь не годится — он
|
||||
-- ротируется, а решение об удалении тела необратимо.
|
||||
--
|
||||
-- БЕЗ DEFAULT намеренно: NULL означает «этот разбор пропусков не считал», и это
|
||||
-- НЕ то же, что ноль. Подстановка нуля объявила бы весь исторический журнал
|
||||
-- проверенным — то самое ложное «терять нечего», ради которого колонка и
|
||||
-- заводится, только теперь с видом измерения. Читатель, принимающий по
|
||||
-- счётчику необратимое решение, обязан трактовать NULL как «не удалять».
|
||||
--
|
||||
-- Data-миграции при этом нет, и это сказано числом: прогон всех 118 тел живого
|
||||
-- архива через разбор даёт `noID=0 noTime=0 malformed=0`, то есть корпус
|
||||
-- пропусков не производил и пересворачивать нечего. «Ничего не потерял по
|
||||
-- замеру» и «проверено этим разбором» — разные утверждения, поэтому колонка
|
||||
-- всё равно остаётся NULL до первой пересвёртки.
|
||||
--
|
||||
-- Статус разбора от пропуска сущности не зависит: `partial` определён списком
|
||||
-- непокрытых секций, и второй источник истины для него завёл бы расхождение
|
||||
-- читателей, которое учёт частичного разбора запрещает явно.
|
||||
ALTER TABLE delivery ADD COLUMN skipped_entities INTEGER;
|
||||
|
||||
-- +goose Down
|
||||
ALTER TABLE delivery DROP COLUMN skipped_entities;
|
||||
@@ -108,3 +108,88 @@ func seed(t *testing.T, path string) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Откат бинаря поверх новой схемы обязан отказывать, а не стартовать молча:
|
||||
// старый бинарь незнакомые секции игнорирует и доставки за окно отката помечает
|
||||
// разобранными — ничто не намекает, что для этого окна нужна пересборка. Класс
|
||||
// «молчание», и после ретеншена тел окно становится невосстановимым.
|
||||
func TestOpenОтвергаетСхемуИзБудущего(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "healthlog.db")
|
||||
seed(t, path)
|
||||
|
||||
db, err := sql.Open("sqlite", "file:"+path)
|
||||
if err != nil {
|
||||
t.Fatalf("sql.Open: %v", err)
|
||||
}
|
||||
if _, err := db.Exec(
|
||||
`INSERT INTO goose_db_version (version_id, is_applied, tstamp)
|
||||
VALUES (99, 1, datetime('now'))`); err != nil {
|
||||
t.Fatalf("вставка версии из будущего: %v", err)
|
||||
}
|
||||
_ = db.Close()
|
||||
|
||||
st, err := store.Open(path)
|
||||
if err == nil {
|
||||
_ = st.Close()
|
||||
t.Fatal("Open молча открыл базу со схемой, которой бинарь не знает")
|
||||
}
|
||||
if !errors.Is(err, store.ErrSchemaMismatch) {
|
||||
t.Errorf("Open дал %v, ожидался ErrSchemaMismatch", err)
|
||||
}
|
||||
}
|
||||
|
||||
// Версия базы НИЖЕ версии бинаря отказом быть не должна: ради этого случая
|
||||
// миграции и существуют. Асимметрия только у Open — OpenForRead строг.
|
||||
func TestOpenНоваяБазаМигрирует(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "fresh.db")
|
||||
st, err := store.Open(path)
|
||||
if err != nil {
|
||||
t.Fatalf("Open новой базы: %v", err)
|
||||
}
|
||||
if err := st.Close(); err != nil {
|
||||
t.Fatalf("Close: %v", err)
|
||||
}
|
||||
// Повторное открытие уже мигрированной базы тоже проходит: current == target.
|
||||
st2, err := store.Open(path)
|
||||
if err != nil {
|
||||
t.Fatalf("повторный Open: %v", err)
|
||||
}
|
||||
_ = st2.Close()
|
||||
}
|
||||
|
||||
// База без журнала миграций нашей не является, и сказать это надо прямо и
|
||||
// сразу. Через goose такой вопрос стоил бы трёх секунд повторов и ответа
|
||||
// «attempt to write a readonly database» — то есть оператор, спросивший про
|
||||
// версию схемы, получил бы ответ про права на файл.
|
||||
func TestOpenForReadЧужаяБазаОтвергаетсяБыстро(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
path := filepath.Join(t.TempDir(), "alien.db")
|
||||
db, err := sql.Open("sqlite", "file:"+path)
|
||||
if err != nil {
|
||||
t.Fatalf("sql.Open: %v", err)
|
||||
}
|
||||
if _, err := db.Exec(`CREATE TABLE t (a INTEGER)`); err != nil {
|
||||
t.Fatalf("создание чужой таблицы: %v", err)
|
||||
}
|
||||
_ = db.Close()
|
||||
|
||||
start := time.Now()
|
||||
st, err := store.OpenForRead(path)
|
||||
if err == nil {
|
||||
_ = st.Close()
|
||||
t.Fatal("чужая база открылась на чтение")
|
||||
}
|
||||
if !errors.Is(err, store.ErrSchemaMismatch) {
|
||||
t.Errorf("OpenForRead дал %v, ожидался ErrSchemaMismatch", err)
|
||||
}
|
||||
// Проверяется свойство «отказ не идёт через повторы записи», а не
|
||||
// конкретная скорость: повторы у goose — три по секунде.
|
||||
if d := time.Since(start); d > time.Second {
|
||||
t.Errorf("отказ занял %v — путь идёт через попытки записи", d)
|
||||
}
|
||||
}
|
||||
|
||||
+100
-38
@@ -9,8 +9,6 @@ import (
|
||||
"fmt"
|
||||
"io/fs"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/jmoiron/sqlx"
|
||||
@@ -27,7 +25,27 @@ type Store struct {
|
||||
db *sqlx.DB
|
||||
}
|
||||
|
||||
// Open открывает БД по пути и накатывает миграции.
|
||||
// Open открывает БД по пути, сверяет версию схемы и накатывает миграции.
|
||||
//
|
||||
// Версия базы ВЫШЕ версии бинаря — отказ, а не повод мигрировать. Иначе откат
|
||||
// бинаря проходит молча: старый бинарь поверх новой схемы стартует успешно,
|
||||
// незнакомые секции игнорирует и доставки за окно отката помечает
|
||||
// разобранными — ничто не намекает, что для этого окна нужна пересборка. Класс
|
||||
// «молчание», и цена его растёт вместе с ретеншеном: после удаления тел окно
|
||||
// становится невосстановимым.
|
||||
//
|
||||
// Цена самого отказа названа вслух, потому что она реальна: сервис не
|
||||
// поднимется, а телефон шлёт непрерывно и молча — доставка, не попавшая в
|
||||
// архив, в журнал не попадает вовсе. Выбор сделан так потому, что откат бинаря
|
||||
// это действие оператора, который в этот момент рядом и видит отказ сразу
|
||||
// (контейнер уходит в цикл перезапуска), а дыры плотных метрик за время простоя
|
||||
// закроют широкий и глубокий проходы синхронизации. Не закроют `stateOfMind` —
|
||||
// у него доставки HAE единственный источник; это и есть цена. Она меньше цены
|
||||
// молчания, которое портит витрину за всё окно отката незаметно.
|
||||
//
|
||||
// Версия базы НИЖЕ версии бинаря отказом не является: ради этого случая
|
||||
// миграции и существуют. Асимметрия только здесь — OpenForRead остаётся
|
||||
// строгим.
|
||||
func Open(dbPath string) (*Store, error) {
|
||||
db, err := sqlx.Connect("sqlite", dsn(dbPath))
|
||||
if err != nil {
|
||||
@@ -59,50 +77,77 @@ func OpenForRead(dbPath string) (*Store, error) {
|
||||
return nil, fmt.Errorf("open sqlite %q read-only: %w", dbPath, err)
|
||||
}
|
||||
|
||||
want, err := latestMigration()
|
||||
ctx := context.Background()
|
||||
// Журнал миграций спрашивается ДО goose и структурно, а не по тексту ошибки
|
||||
// драйвера. Причина не в стиле: `GetVersions` при отсутствии таблицы идёт
|
||||
// её СОЗДАВАТЬ, на соединении `mode=ro` это три секунды повторов и отказ
|
||||
// «attempt to write a readonly database» — оператор, спросивший про версию
|
||||
// схемы, получал бы ответ про права на файл. База без журнала миграций
|
||||
// нашей не является, и сказать это надо прямо.
|
||||
ok, err := hasMigrationLog(ctx, db)
|
||||
if err != nil {
|
||||
_ = db.Close()
|
||||
return nil, err
|
||||
}
|
||||
var got int64
|
||||
if err := db.Get(&got, `SELECT max(version_id) FROM goose_db_version`); err != nil {
|
||||
if !ok {
|
||||
_ = db.Close()
|
||||
return nil, fmt.Errorf("read schema version: %w", err)
|
||||
return nil, fmt.Errorf("%w: журнала миграций в базе нет", ErrSchemaMismatch)
|
||||
}
|
||||
if got != want {
|
||||
|
||||
inDB, inBinary, err := readSchemaVersion(ctx, db)
|
||||
if err != nil {
|
||||
_ = db.Close()
|
||||
return nil, fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, got, want)
|
||||
return nil, err
|
||||
}
|
||||
// Строгое равенство, в отличие от Open: у чтения нет способа догнать схему,
|
||||
// а база старее бинаря отдала бы колонки, которых в ней ещё нет. Так уже
|
||||
// нормировано пересборкой, и настоящее правило её не ослабляет.
|
||||
if inDB != inBinary {
|
||||
_ = db.Close()
|
||||
return nil, fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, inDB, inBinary)
|
||||
}
|
||||
return &Store{db: db}, nil
|
||||
}
|
||||
|
||||
// latestMigration — номер последней миграции, вшитой в бинарь.
|
||||
func latestMigration() (int64, error) {
|
||||
entries, err := fs.ReadDir(migrationsFS, "migrations")
|
||||
// hasMigrationLog говорит, есть ли в базе журнал миграций goose. Структурный
|
||||
// вопрос к самой базе, а не разбор текста ошибки драйвера: сообщения драйвера
|
||||
// контрактом не являются — правило записано в isBusy и действует здесь.
|
||||
func hasMigrationLog(ctx context.Context, db *sqlx.DB) (bool, error) {
|
||||
const q = `SELECT count(*) FROM sqlite_master WHERE type = 'table' AND name = 'goose_db_version'`
|
||||
|
||||
var n int
|
||||
if err := db.GetContext(ctx, &n, q); err != nil {
|
||||
return false, fmt.Errorf("read migration log presence: %w", err)
|
||||
}
|
||||
return n > 0, nil
|
||||
}
|
||||
|
||||
// readSchemaVersion отвечает, какая версия схемы лежит в базе и какую знает
|
||||
// бинарь. Единственное место, где версия ЧИТАЕТСЯ, — сравнивают её два способа
|
||||
// открытия по-разному, а читают одинаково.
|
||||
//
|
||||
// Спрашиваем сам goose, а не собственный `SELECT max(version_id)`: имя таблицы
|
||||
// учёта, имя колонки и правило «максимум = текущая версия» принадлежат ему.
|
||||
// Рукописная копия его приватной схемы разошлась бы при обновлении зависимости,
|
||||
// причём не отказом, а тем, что страж перестал бы ловить, — то есть ровно тем,
|
||||
// что страж и обязан не допускать. Заодно исчезает собственный разбор имён
|
||||
// `NNNNN_*.sql` и вопрос «как отличить пустую таблицу от отсутствующей, не
|
||||
// читая текст ошибки драйвера»: на новой базе goose отдаёт 0 сам.
|
||||
//
|
||||
// Оговорка, без которой обещание непроверяемо: `GetVersions` при ОТСУТСТВИИ
|
||||
// таблицы учёта идёт её создавать. На соединении только для чтения это отказ, и
|
||||
// вызывающий обязан отсеять такую базу раньше (см. hasMigrationLog); на
|
||||
// соединении с записью создание законно — им и начинается новая база.
|
||||
func readSchemaVersion(ctx context.Context, db *sqlx.DB) (inDB, inBinary int64, err error) {
|
||||
p, err := newProvider(db)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("read migrations dir: %w", err)
|
||||
return 0, 0, err
|
||||
}
|
||||
var top int64
|
||||
for _, e := range entries {
|
||||
name := e.Name()
|
||||
// Неразобранное имя — отказ, а не пропуск: страж «версия схемы не та»,
|
||||
// молча не заметивший миграцию, перестаёт страховать, не сказав об этом.
|
||||
idx := strings.IndexByte(name, '_')
|
||||
if idx <= 0 {
|
||||
return 0, fmt.Errorf("имя миграции %q не вида NNNNN_*.sql", name)
|
||||
}
|
||||
v, err := strconv.ParseInt(name[:idx], 10, 64)
|
||||
if err != nil {
|
||||
return 0, fmt.Errorf("имя миграции %q не вида NNNNN_*.sql", name)
|
||||
}
|
||||
if v > top {
|
||||
top = v
|
||||
}
|
||||
inDB, inBinary, err = p.GetVersions(ctx)
|
||||
if err != nil {
|
||||
return 0, 0, fmt.Errorf("read schema version: %w", err)
|
||||
}
|
||||
if top == 0 {
|
||||
return 0, errors.New("миграций не найдено")
|
||||
}
|
||||
return top, nil
|
||||
return inDB, inBinary, nil
|
||||
}
|
||||
|
||||
// Close закрывает соединение с БД.
|
||||
@@ -150,21 +195,38 @@ func readOnlyDSN(path string) string {
|
||||
// пересборка витрины откроет второе, и хранилище не должно зависеть от того,
|
||||
// что вызывающий этого не сделает.
|
||||
func migrate(db *sqlx.DB) error {
|
||||
sub, err := fs.Sub(migrationsFS, "migrations")
|
||||
ctx := context.Background()
|
||||
|
||||
inDB, inBinary, err := readSchemaVersion(ctx, db)
|
||||
if err != nil {
|
||||
return fmt.Errorf("goose migrations fs: %w", err)
|
||||
return err
|
||||
}
|
||||
if inDB > inBinary {
|
||||
return fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, inDB, inBinary)
|
||||
}
|
||||
|
||||
p, err := goose.NewProvider(goose.DialectSQLite3, db.DB, sub)
|
||||
p, err := newProvider(db)
|
||||
if err != nil {
|
||||
return fmt.Errorf("goose provider: %w", err)
|
||||
return err
|
||||
}
|
||||
if _, err := p.Up(context.Background()); err != nil {
|
||||
if _, err := p.Up(ctx); err != nil {
|
||||
return fmt.Errorf("goose up: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func newProvider(db *sqlx.DB) (*goose.Provider, error) {
|
||||
sub, err := fs.Sub(migrationsFS, "migrations")
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("goose migrations fs: %w", err)
|
||||
}
|
||||
p, err := goose.NewProvider(goose.DialectSQLite3, db.DB, sub)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("goose provider: %w", err)
|
||||
}
|
||||
return p, nil
|
||||
}
|
||||
|
||||
// Now — единая точка генерации времени: UTC, секундная точность.
|
||||
// Секунды дают фиксированную ширину RFC 3339, а значит лексикографическая
|
||||
// сортировка TEXT совпадает с хронологией.
|
||||
|
||||
@@ -0,0 +1,77 @@
|
||||
package store
|
||||
|
||||
import "context"
|
||||
|
||||
// pickBest выбирает победителя среди кандидатов на одни координаты.
|
||||
//
|
||||
// **Победитель — функция МНОЖЕСТВА кандидатов, а не порядка их поступления.**
|
||||
// Попарная свёртка этого не даёт: полнота (или покрытие) — частичный порядок,
|
||||
// тай-брейк — тотальный, и вместе они образуют нетранзитивное отношение победы,
|
||||
// то есть цикл. При цикле повторная свёртка одной и той же доставки меняет
|
||||
// содержимое витрины, и она перестаёт быть свёрткой журнала. Проверено дважды:
|
||||
// сперва на точках, где нетранзитивность нашлась перебором троек, потом на
|
||||
// сущностях, где ту же ошибку повторили молча.
|
||||
//
|
||||
// Отсюда и общий помощник вместо второй рукописной копии: механизм один,
|
||||
// отношения разные. `architecture.md` уже обещает смену тай-брейка точек, когда
|
||||
// род метрики будет измерен, — то есть правка одного экземпляра при живом
|
||||
// втором запланирована заранее, и расхождение правил слияния ломает детерминизм
|
||||
// свёртки молча.
|
||||
//
|
||||
// dominates(a, b) обязан быть СТРОГИМ превосходством: a не хуже b и b не не
|
||||
// хуже a. Иначе взаимно покрывающие друг друга кандидаты выбьют друг друга, и
|
||||
// множество непревзойдённых окажется пустым.
|
||||
//
|
||||
// less обязан быть ТОТАЛЬНЫМ строгим порядком: при неединственном минимуме
|
||||
// победителем оказывается просто первый в срезе, то есть порядок элементов на
|
||||
// проводе, а он у HAE нестабилен.
|
||||
//
|
||||
// Отбор непревзойдённых КВАДРАТИЧЕН по числу кандидатов, и это названо вслух,
|
||||
// потому что число кандидатов приходит из чужого тела. Отсюда `ctx`: цикл, чья
|
||||
// стоимость определяется размером входа, обязан видеть отмену. Без него тело с
|
||||
// двадцатью тысячами версий одного ключа занимало бы единственного воркера
|
||||
// свёртки дольше, чем длится его же дедлайн, — то есть дедлайн, заведённый
|
||||
// ровно против такого случая, не значил бы ничего. Замерено: n=4000 — 3.9 с,
|
||||
// n=8000 — вчетверо больше.
|
||||
//
|
||||
// Вызывающий обязан сокращать множество до входа сюда: совпавших кандидатов
|
||||
// схлопывать, а число различных — ограничивать. Помощник этого не делает
|
||||
// сам — что считать «тем же» кандидатом, знает только он.
|
||||
//
|
||||
// Возвращает индекс победителя и индексы непревзойдённых — вторые нужны тем,
|
||||
// кто считает несравнимость среди них. Пустой срез кандидатов — нарушенный
|
||||
// инвариант вызывающего: оба сегодняшних вызова собирают множество из карты и
|
||||
// пустого дать не могут.
|
||||
func pickBest[T any](ctx context.Context, cands []T, dominates func(a, b T) bool, less func(a, b T) bool) (winner int, maximal []int, err error) {
|
||||
if len(cands) == 0 {
|
||||
return 0, nil, errNoCandidates
|
||||
}
|
||||
|
||||
maximal = make([]int, 0, len(cands))
|
||||
for i := range cands {
|
||||
if err := ctx.Err(); err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
beaten := false
|
||||
for j := range cands {
|
||||
if i == j {
|
||||
continue
|
||||
}
|
||||
if dominates(cands[j], cands[i]) {
|
||||
beaten = true
|
||||
break
|
||||
}
|
||||
}
|
||||
if !beaten {
|
||||
maximal = append(maximal, i)
|
||||
}
|
||||
}
|
||||
|
||||
winner = maximal[0]
|
||||
for _, i := range maximal[1:] {
|
||||
if less(cands[i], cands[winner]) {
|
||||
winner = i
|
||||
}
|
||||
}
|
||||
return winner, maximal, nil
|
||||
}
|
||||
Reference in New Issue
Block a user