package store import ( "context" "database/sql" "errors" "fmt" "sort" "strings" "time" ) // Запросы ряда точек. Вынесены константами по той же причине, что и запросы // каталога: по ним проверяется план выполнения. const ( // Объекты выбранного слоя. Точный префикс первичного ключа // (metric, layer, hour_utc); `payload` здесь и нужен. seriesPointsQuery = ` SELECT hour_utc, units, payload FROM bucket WHERE metric = ? AND layer = ? AND hour_utc BETWEEN ? AND ? ORDER BY hour_utc` ) // layerSpansQuery строит выборку охватов слоёв метрики внутри периода. // // Границы берутся ТОЧНЫЕ (`first_ts`/`last_ts`), а не по `hour_utc`: объект // адресуется часом, а ряд отбирается точной меткой, и на периоде короче часа // эти два множества расходятся. Слой, выбранный по часам, отдал бы пустой ряд // при непустых данных соседнего слоя — час объекта попадает в период, а его // единственная точка в период не попадает. // // СЛОИ ПЕРЕЧИСЛЕНЫ ЯВНО, и это не украшение запроса, а его цена. Индекс // `bucket_catalog` идёт `(metric, layer, hour_utc, …)`; без предиката по слою // SQLite не может сузить поиск по `hour_utc` внутри индекса и просматривает // ВСЕ строки метрики за всю историю, применяя период построчным фильтром. План // при этом выглядит успешным (`SEARCH … USING COVERING INDEX`), а цена растёт // вместе с возрастом сервиса при любой ширине запроса: измерено эксплуатационным // проходом ревью на копии схемы — 2.06 мс при 52 560 строках метрики против // 13.9 мс при 350 400, и 0.026 мс с этим перечислением. Словарь слоёв задаёт // вызывающий: правило принадлежит домену, хранилище лишь выбирает по нему. func layerSpansQuery(layers int) string { return ` SELECT layer, min(first_ts), max(last_ts) FROM bucket WHERE metric = ? AND layer IN (` + strings.TrimSuffix(strings.Repeat("?,", layers), ",") + `) AND hour_utc BETWEEN ? AND ? GROUP BY layer ORDER BY layer` } // errEmptyWindow — период задан вывернутым. Нарушенный инвариант вызывающего: // транспорт обязан отвергнуть такой запрос раньше. var errEmptyWindow = errors.New("окно ряда: from не раньше to") // errNoLayers — словарь слоёв не задан. Без него выборка охватов выродилась бы // в скан всей истории метрики, а не отдала бы пустой результат: молчаливая // деградация цены хуже отказа. var errNoLayers = errors.New("окно ряда: словарь слоёв пуст") // LayerSpan — охват одного слоя метрики внутри запрошенного периода. // // First и Last — метки первой и последней ТОЧКИ объектов слоя, попавших в // границы часов запроса. Это границы данных, а не обещание покрытия: внутри // законно есть дыры. type LayerSpan struct { Layer string First time.Time Last time.Time } // SeriesPoint — точка ряда вместе с единицами объекта, из которого она // прочитана. // // Единицы едут с точкой, а не с рядом: они хранятся на часовом объекте, и // метрика, чьи объекты разошлись единицами, обязана показать это строкой, а не // выбрать одно из двух молча. type SeriesPoint struct { Point Units string } // SeriesWindow — что читать. Все правила задаёт вызывающий: хранилище выбирает // строки, а не решает, какие из них правильные. type SeriesWindow struct { Metric string // From включительно, To исключительно. From, To time.Time // Layer — явно запрошенный слой; пустая строка означает «выбери сам», // и тогда зовётся pick. Layer string // Layers — словарь слоёв, среди которых вообще имеет смысл искать. Задаёт // вызывающий: перечень слоёв — знание домена, а хранилищу он нужен, чтобы // выборка охватов не превращалась в скан всей истории метрики. Layers []string // Measure — окно измерения рода агрегации той же метрики. Measure CatalogWindow } // SeriesSnapshot — весь вход ответа точек, снятый ОДНОЙ транзакцией чтения. // // Единый снимок здесь не аккуратность, а условие непротиворечивости: приём идёт // непрерывно, и фоновая свёртка вправе закоммитить между выбором слоя и чтением // точек. Тогда слой выбран по одному состоянию витрины, ряд прочитан по // второму, а род измерен по третьему — ответ внутренне противоречив и от свежего // неотличим. Пара проб версии такой ответ ОБНАРУЖИТ (метки не будет), но не // предотвратит: тело всё равно уедет. Тот же довод записан у входа каталога. type SeriesSnapshot struct { // Охватов слоёв здесь НЕТ намеренно: они приходят в pick аргументом, и это // единственные ворота решения о слое. Поле наружу предлагало бы те же данные // любому будущему вызывающему и приглашало бы решить слой пост-фактум, мимо // правила, — двое ворот к одному решению. // // Layer — слой, из которого собран ряд. Пустая строка означает, что слоя // нет: выбирать было не из чего либо запрошенный слой пуст. Layer string // Points — ряд, отобранный до точных границ периода и упорядоченный. Points []SeriesPoint // Pairs — общие часы метрики, от самых свежих к старым. Pairs []HourPair } // ReadSeries снимает вход ответа точек одной транзакцией чтения. // // Правило выбора слоя остаётся В ДОМЕНЕ и приходит сюда функцией pick, // вызываемой ВНУТРИ транзакции. Форма не изобретена: VersionedRead уже // принимает работу колбэком, а ReadCatalog уже получает параметры правила // структурой. Перенести само правило сюда значило бы вернуть в хранилище // решение, которое из него специально убирали. // // pick зовётся только когда слой не задан явно и есть из чего выбирать; вернуть // он может пустую строку — это законный исход «ряда нет». func (s *Store) ReadSeries(ctx context.Context, w SeriesWindow, pick func([]LayerSpan) string) (SeriesSnapshot, error) { // Границы В ТЕКСТ ОШИБКИ НЕ ИДУТ. Текст доезжает до записи лога вызывающего, // а границы периода — часть запроса о здоровье; сегодня транспорт отвергает // такой запрос раньше, но второй вызывающий (адаптер MCP) откроет этот путь. if !w.From.Before(w.To) { return SeriesSnapshot{}, errEmptyWindow } if len(w.Layers) == 0 { return SeriesSnapshot{}, errNoLayers } // `ReadOnly` у `modernc.org/sqlite` выбирает `BEGIN` вместо // `BEGIN IMMEDIATE` и записи НЕ ЗАПРЕЩАЕТ (прочитан исходник драйвера // зафиксированной версии). То, что этот путь не пишет, держится ревью, а не // драйвером; полагаться на флаг как на защиту нельзя. tx, err := s.db.BeginTx(ctx, &sql.TxOptions{ReadOnly: true}) if err != nil { return SeriesSnapshot{}, fmt.Errorf("begin read tx: %w", err) } defer func() { _ = tx.Rollback() }() // Границы по часам ШИРЕ запроса: точка 10:59 живёт в объекте 10:00, и // огрубление до часа — единственный способ её не потерять. Точный отбор // идёт ниже, по меткам самих точек. fromHour := FormatTime(w.From.UTC().Truncate(time.Hour)) toHour := FormatTime(w.To.UTC().Truncate(time.Hour)) spans, err := readLayerSpans(ctx, tx, w, fromHour, toHour) if err != nil { return SeriesSnapshot{}, err } out := SeriesSnapshot{Layer: w.Layer} if out.Layer == "" && pick != nil { out.Layer = pick(spans) } if out.Layer != "" { // Слой ВЫБРАННЫЙ, а не запрошенный: они расходятся ровно тогда, когда // правило сработало, — то есть в самом частом случае. if out.Points, err = readSeriesPoints(ctx, tx, w, out.Layer, fromHour, toHour); err != nil { return SeriesSnapshot{}, err } } pairs, err := commonHours(ctx, tx, w.Metric, w.Measure) if err != nil { return SeriesSnapshot{}, err } if len(pairs) > 0 { if err := readHourPairs(ctx, tx, w.Metric, w.Measure, pairs); err != nil { return SeriesSnapshot{}, err } } out.Pairs = pairs return out, nil } func readLayerSpans(ctx context.Context, tx *sql.Tx, w SeriesWindow, fromHour, toHour string) ([]LayerSpan, error) { args := make([]any, 0, len(w.Layers)+3) args = append(args, w.Metric) for _, l := range w.Layers { args = append(args, l) } args = append(args, fromHour, toHour) rows, err := tx.QueryContext(ctx, layerSpansQuery(len(w.Layers)), args...) if err != nil { return nil, fmt.Errorf("select layer spans: %w", err) } defer func() { _ = rows.Close() }() out := make([]LayerSpan, 0, 4) for rows.Next() { var sp LayerSpan var first, last string if err := rows.Scan(&sp.Layer, &first, &last); err != nil { return nil, fmt.Errorf("scan layer span: %w", err) } if sp.First, err = ParseTime(first); err != nil { return nil, err } if sp.Last, err = ParseTime(last); err != nil { return nil, err } out = append(out, sp) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("select layer spans: %w", err) } return out, nil } // readSeriesPoints разжимает объекты выбранного слоя и отбирает точки до точных // границ периода. // // Принадлежность точки периоду определяется её НАЧАЛОМ — тем же правилом, каким // час объекта берётся по началу точки. Цена названа в спеке: интервал, // начавшийся раньше from, в ответ не входит. func readSeriesPoints(ctx context.Context, tx *sql.Tx, w SeriesWindow, layer, fromHour, toHour string) ([]SeriesPoint, error) { rows, err := tx.QueryContext(ctx, seriesPointsQuery, w.Metric, layer, fromHour, toHour) if err != nil { return nil, fmt.Errorf("select series points: %w", err) } defer func() { _ = rows.Close() }() // Непустой срез, а не nil: пустой ряд обязан уехать клиенту как `[]`. out := make([]SeriesPoint, 0, 64) for rows.Next() { var hour, units string var payload []byte if err := rows.Scan(&hour, &units, &payload); err != nil { return nil, fmt.Errorf("scan series bucket: %w", err) } points, err := decodePayload(payload) if err != nil { return nil, err } for _, p := range points { if p.Start.Before(w.From) || !p.Start.Before(w.To) { continue } out = append(out, SeriesPoint{Point: p, Units: units}) } } if err := rows.Err(); err != nil { return nil, fmt.Errorf("select series points: %w", err) } // Порядок утверждается здесь, а не наследуется от порядка объектов: на нём // стоят и байтовое утверждение формы ответа, и метка условного запроса. // Пара (начало, конец) внутри одного слоя одной метрики есть ключ // идентичности, поэтому порядок ею определён однозначно. sort.SliceStable(out, func(i, j int) bool { if !out[i].Start.Equal(out[j].Start) { return out[i].Start.Before(out[j].Start) } return out[i].End.Before(out[j].End) }) return out, nil }