первая встреча непокрытой секции стала наблюдаемым событием
- свёртка спрашивает журнал, встречалось ли имя строго раньше по паре (received_at, id), и пишет WARN с атрибутом uncovered_new; повторные молчат. Признак выводится, а не хранится — реестр был бы второй копией факта - добавлена подкоманда `healthlog uncovered`: перечень накопленного, чтение только на чтение, экранированные имена и названные границы носителя - синк документации: ADR о выводе новизны из журнала, две записи в журнал дефектов, два правила промоутом в конвенции, терминал оператора назван адресатом недоверенного входа
This commit is contained in:
@@ -0,0 +1,243 @@
|
||||
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
|
||||
}
|
||||
Reference in New Issue
Block a user