package store import ( "context" "database/sql" "errors" "fmt" "strings" "time" ) // UncoveredSection — строка перечня непокрытых секций: имя, которое поток // приносил, и границы его встреч по журналу. // // Число доставок и границы считаются по ВСЕМ статусам разбора: список непокрытых // секций переживает отказ свёртки, и молчать о таком имени значило бы терять как // раз подозрительное. type UncoveredSection struct { // Name — верхнеуровневый ключ `data`, дословно как прислал HAE. Приходит из // чужого тела: длина ограничена разбором, содержимое — ничем. Name string // Deliveries — сколько доставок принесло это имя. Deliveries int64 // FirstSeen и FirstDeliveryID — первая встреча ПО ЖУРНАЛУ; LastSeen и // LastDeliveryID — последняя. Идентификатор нужен, чтобы достать тело из // архива и посмотреть форму секции глазами. FirstSeen time.Time FirstDeliveryID string LastSeen time.Time LastDeliveryID string } // usableList — условие «списку непокрытых секций есть что сказать». // // Отдельной константой, потому что стоит во всех трёх запросах и означает одно: // по умолчанию колонка держит `[]`, а строки с ним до json_each доходить не // должны. // // `json_valid` здесь не перестраховка. `json_each` над неразбираемым значением // отвечает ошибкой и валит ВЕСЬ запрос, а не пропускает одну строку (проверено // на закреплённом драйвере: `SQL logic error: malformed JSON`). Одна такая // строка — из ручной правки, из будущей миграции данных мимо `json.Marshal`, // из восстановления базы чужим инструментом — обесценила бы и сверку (вечное // «сверка не состоялась» на каждой доставке), и перечень (полный отказ команды // без указания виновника). Штатный путь записи такого значения не производит, и // именно поэтому отказ был бы необъясним. const usableList = `uncovered_sections <> '[]' AND uncovered_sections <> ''` // listOf — содержимое колонки, приведённое к разбираемому виду. // // Проверка стоит ВНУТРИ аргумента `json_each`, а не условием в `WHERE`: // табличная функция получает значение строки раньше, чем применится фильтр, и // порядок этот SQLite не обещает. Условие в `WHERE` работало бы, пока // планировщик проталкивает его вниз, и молча перестало бы — с ошибкой, роняющей // весь запрос, а не строку. const listOf = `json_each(CASE WHEN json_valid(d.uncovered_sections) THEN d.uncovered_sections ELSE '[]' END)` // SectionsSeenBefore отвечает, какие из имён уже встречались в доставках, // стоящих в журнале СТРОГО РАНЬШЕ указанной. // // Строгое сравнение пары `(received_at, id)` решает три вещи разом: собственная // строка доставки в счёт не идёт при любом порядке записи исхода; повторная // свёртка той же доставки даёт тот же ответ, поэтому пересборка воспроизводит те // же события; и судит имя порядок ЖУРНАЛА, а не порядок прогона — строка учёта // становится видимой воркеру только после записи тела, так что при конкурентном // приёме доставка может стать видимой после более новой. // // Индекса по `uncovered_sections` в схеме нет, поэтому проход по журналу // полный, и форма запроса выбрана ЗАМЕРОМ (числа и метод — в `design.md` // изменения `aktivnaya-proverka-novyh-sekcij`, решение 3; повторяется прогоном // `tmp/seenmeasure`). Здесь только правило, которое из замера следует: // // - ранний выход `LIMIT 1` срабатывает, лишь когда первая встреча имени лежит // БЛИЗКО К НАЧАЛУ журнала, — строки просматриваются от старых к новым. // Секция, приезжающая давно, стоит десятки микросекунд; секция, появившаяся // только что, — почти полный проход на каждой доставке, пока её не покроет // отдельная задача; // - имён больше одного спрашиваются ОДНИМ запросом: тридцать два запроса // подряд стоили секунду с лишним на доставку, одна выборка — столько же, // сколько один полный проход, независимо от числа имён. // // Вызывается ВНЕ транзакции записи: проход по растущему журналу внутри неё // удерживал бы блокировку, а конкурирующий приём отвечает `500` по доставке, чьё // тело уже на диске. func (s *Store) SectionsSeenBefore(ctx context.Context, names []string, before DeliveryRef) (map[string]struct{}, error) { unique := make([]string, 0, len(names)) dedup := make(map[string]struct{}, len(names)) for _, name := range names { if _, ok := dedup[name]; ok { continue } dedup[name] = struct{}{} unique = append(unique, name) } if len(unique) == 0 { return nil, nil } at := FormatTime(before.ReceivedAt) if len(unique) == 1 { return s.sectionSeenOnce(ctx, unique[0], at, before.ID) } return s.sectionsSeenBatch(ctx, unique, at, before.ID) } // sectionSeenOnce спрашивает про одно имя с ранним выходом. func (s *Store) sectionSeenOnce(ctx context.Context, name, at, id string) (map[string]struct{}, error) { const q = ` SELECT 1 FROM delivery d, ` + listOf + ` j WHERE ` + usableList + ` AND (d.received_at, d.id) < (?, ?) AND j.value = ? LIMIT 1` var one int err := s.db.QueryRowxContext(ctx, q, at, id, name).Scan(&one) switch { case err == nil: return map[string]struct{}{name: {}}, nil case errors.Is(err, sql.ErrNoRows): return nil, nil default: // Имя в текст ошибки не попадает целиком: оно приходит из чужого тела, а // ошибка уходит в лог уровня выше DEBUG. Та же граница, что у координат // наблюдения категориального значения. return nil, fmt.Errorf("section seen before (%s): %w", clipCoord(name), err) } } // sectionsSeenBatch спрашивает про все имена одним проходом по журналу. func (s *Store) sectionsSeenBatch(ctx context.Context, names []string, at, id string) (map[string]struct{}, error) { q := ` SELECT DISTINCT j.value FROM delivery d, ` + listOf + ` j WHERE ` + usableList + ` AND (d.received_at, d.id) < (?, ?) AND j.value IN (?` + strings.Repeat(", ?", len(names)-1) + `)` args := make([]any, 0, len(names)+2) args = append(args, at, id) for _, name := range names { args = append(args, name) } rows, err := s.db.QueryContext(ctx, q, args...) if err != nil { return nil, fmt.Errorf("sections seen before: %w", err) } defer func() { _ = rows.Close() }() seen := make(map[string]struct{}, len(names)) for rows.Next() { var name string if err := rows.Scan(&name); err != nil { return nil, fmt.Errorf("scan section seen before: %w", err) } seen[name] = struct{}{} } if err := rows.Err(); err != nil { return nil, fmt.Errorf("sections seen before: %w", err) } return seen, nil } // UncoveredSections возвращает перечень непокрытых секций — не больше limit // строк — и общее число различных имён в журнале. // // Предел объявлен, а не подразумевается: граница разбора в 32 имени действует на // ОДНУ доставку, а различных имён журнал накопит сколько угодно — достаточно // версии HAE, кладущей в ключ переменную часть. Второе возвращаемое значение // нужно, чтобы вывод мог назвать остаток числом, а не молча оборваться. // // Границы встреч берутся минимумом и максимумом склейки `received_at` и `id`: // метка хранится в RFC3339 фиксированной ширины, поэтому склейка сравнивается // лексикографически ровно как пара, а одна выборка с двумя агрегатами не // опирается на то, из какой строки SQLite возьмёт голые колонки. func (s *Store) UncoveredSections(ctx context.Context, limit int) ([]UncoveredSection, int64, error) { if limit <= 0 { return nil, 0, fmt.Errorf("предел строк перечня должен быть положительным, дано %d", limit) } const countQ = ` SELECT count(DISTINCT j.value) FROM delivery d, ` + listOf + ` j WHERE ` + usableList var total int64 if err := s.db.GetContext(ctx, &total, countQ); err != nil { return nil, 0, fmt.Errorf("count uncovered sections: %w", err) } const q = ` SELECT j.value, count(DISTINCT d.id), min(d.received_at || ' ' || d.id), max(d.received_at || ' ' || d.id) FROM delivery d, ` + listOf + ` j WHERE ` + usableList + ` GROUP BY j.value ORDER BY j.value LIMIT ?` rows, err := s.db.QueryContext(ctx, q, limit) if err != nil { return nil, 0, fmt.Errorf("select uncovered sections: %w", err) } defer func() { _ = rows.Close() }() var out []UncoveredSection for rows.Next() { var u UncoveredSection var first, last string if err := rows.Scan(&u.Name, &u.Deliveries, &first, &last); err != nil { return nil, 0, fmt.Errorf("scan uncovered section: %w", err) } if u.FirstSeen, u.FirstDeliveryID, err = splitJournalKey(first); err != nil { return nil, 0, err } if u.LastSeen, u.LastDeliveryID, err = splitJournalKey(last); err != nil { return nil, 0, err } out = append(out, u) } if err := rows.Err(); err != nil { return nil, 0, fmt.Errorf("select uncovered sections: %w", err) } return out, total, nil } // splitJournalKey разбирает склейку метки журнала и идентификатора доставки. func splitJournalKey(key string) (time.Time, string, error) { at, id, ok := strings.Cut(key, " ") if !ok { return time.Time{}, "", fmt.Errorf("ключ журнала без разделителя: %q", key) } t, err := ParseTime(at) if err != nil { return time.Time{}, "", err } return t, id, nil }