Files
av 8db2ec7ff4 Цена читающего маршрута: чекпойнт WAL по таймеру и условный запрос
- рядом с воркером свёртки живёт горутина, раз в минуту разбирающая журнал
  пассивным чекпойнтом; «журнал не разбирается» видно строкой владельцу, а не
  только по `df`. Признак — пара чисел, а не флаг занятости: тот молчит под
  удерживаемым читателем (`busy=0` при 6256 страницах и пяти перенесённых), а
  при занятой блокировке отдаёт `-1` вместо ответа, и `-1 >= -1` читалось бы как
  «разобрано целиком»
- каталог отвечает `304` на `If-None-Match`, не открывая снимок витрины. Метка
  собрана из всего, от чего зависит ответ: версии витрины (`data_version` с
  закреплённого соединения плюс поколение — значение локально для соединения и
  не переживает переоткрытия), горизонта измерения и области действия ресурса.
  Версия снимается до и после сборки: снятая после пометила бы устаревший снимок
  свежим номером
- предел и дедлайн ответа отложены в задачу Read API точек вместе с измеренной
  ценой первого запроса; попутно починен флаки-тест чужой задачи, искавший
  значение точки в сыром буфере записи лога
2026-08-02 20:42:22 +03:00

273 lines
15 KiB
Go

// Package store — SQLite-витрина healthlog: доставки и (позже) разобранные
// точки. Витрина производна от сырого архива и пересобирается из него.
package store
import (
"context"
"embed"
"errors"
"fmt"
"io/fs"
"net/url"
"time"
"github.com/jmoiron/sqlx"
"github.com/pressly/goose/v3"
_ "modernc.org/sqlite" // чистый Go-драйвер SQLite, без cgo
)
//go:embed migrations/*.sql
var migrationsFS embed.FS
// Store — доступ к витрине.
type Store struct {
db *sqlx.DB
// probe — закреплённое соединение версии витрины, заводится по первому
// запросу версии. См. version.go: значение `data_version` локально для
// соединения, и брать его из пула нельзя.
probe versionProbe
}
// Open открывает БД по пути, сверяет версию схемы и накатывает миграции.
//
// Версия базы ВЫШЕ версии бинаря — отказ, а не повод мигрировать. Иначе откат
// бинаря проходит молча: старый бинарь поверх новой схемы стартует успешно,
// незнакомые секции игнорирует и доставки за окно отката помечает
// разобранными — ничто не намекает, что для этого окна нужна пересборка. Класс
// «молчание», и цена его растёт вместе с ретеншеном: после удаления тел окно
// становится невосстановимым.
//
// Цена самого отказа названа вслух, потому что она реальна: сервис не
// поднимется, а телефон шлёт непрерывно и молча — доставка, не попавшая в
// архив, в журнал не попадает вовсе. Выбор сделан так потому, что откат бинаря
// это действие оператора, который в этот момент рядом и видит отказ сразу
// (контейнер уходит в цикл перезапуска), а дыры плотных метрик за время простоя
// закроют широкий и глубокий проходы синхронизации. Не закроют `stateOfMind` —
// у него доставки HAE единственный источник; это и есть цена. Она меньше цены
// молчания, которое портит витрину за всё окно отката незаметно.
//
// Версия базы НИЖЕ версии бинаря отказом не является: ради этого случая
// миграции и существуют. Асимметрия только здесь — OpenForRead остаётся
// строгим.
func Open(dbPath string) (*Store, error) {
db, err := sqlx.Connect("sqlite", dsn(dbPath))
if err != nil {
return nil, fmt.Errorf("open sqlite %q: %w", dbPath, err)
}
if err := migrate(db); err != nil {
_ = db.Close()
return nil, err
}
return &Store{db: db}, nil
}
// ErrSchemaMismatch — версия схемы базы не та, которую знает бинарь.
var ErrSchemaMismatch = errors.New("версия схемы базы не совпадает с версией бинаря")
// OpenForRead открывает базу только для чтения и **без наката миграций**.
//
// Обычный Open мигрирует безусловно, а миграции здесь меняют не только схему, но
// и данные: та, что ввела частичный разбор, переписала parse_status у всех строк.
// Значит утилита, которой достаточно прочитать учёт, обычным открытием нарушала
// бы обещание «рабочую базу не трогаем», — и хуже: свежий бинарь мигрировал бы
// схему под работающим старым сервисом, который держит запросы к прежней.
//
// Расхождение версий — отказ с указанием обеих, а не повод мигрировать.
func OpenForRead(dbPath string) (*Store, error) {
db, err := sqlx.Connect("sqlite", readOnlyDSN(dbPath))
if err != nil {
return nil, fmt.Errorf("open sqlite %q read-only: %w", dbPath, err)
}
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
}
if !ok {
_ = db.Close()
return nil, fmt.Errorf("%w: журнала миграций в базе нет", ErrSchemaMismatch)
}
inDB, inBinary, err := readSchemaVersion(ctx, db)
if err != nil {
_ = db.Close()
return nil, err
}
// Строгое равенство, в отличие от Open: у чтения нет способа догнать схему,
// а база старее бинаря отдала бы колонки, которых в ней ещё нет. Так уже
// нормировано пересборкой, и настоящее правило её не ослабляет.
if inDB != inBinary {
_ = db.Close()
return nil, fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, inDB, inBinary)
}
return &Store{db: db}, nil
}
// 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, 0, err
}
inDB, inBinary, err = p.GetVersions(ctx)
if err != nil {
return 0, 0, fmt.Errorf("read schema version: %w", err)
}
return inDB, inBinary, nil
}
// Close закрывает соединение с БД.
func (s *Store) Close() error {
// Щуп версии закрывается сам и раньше пула: закреплённое соединение
// переживает `db.Close()` и продолжает отвечать на запросы — проверено.
s.closeProbe()
if err := s.db.Close(); err != nil {
return fmt.Errorf("close sqlite: %w", err)
}
return nil
}
// `journal_size_limit` держит верхнюю границу файла журнала. Чекпойнт
// возвращает страницы в базу, но файл оставляет на пике: измерено — 51 МБ до и
// после успешного переноса 12502 страниц. С лимитом первая же следующая запись
// усекает файл до предела.
//
// Величина взята из чужой практики (гайды по SQLite в проде ставят 26–64 МБ) и
// собственного измерения: суточный поток даёт около 23 МБ архива, то есть
// журнал такого размера означает не всплеск, а удерживаемый снимок. С пределом
// тела приёма она НЕ связана, хотя и совпадает по порядку: журнал растёт от
// чтения, а не от размера доставки, и менять её вслед за `ingest.max_body_mb`
// незачем.
const journalSizeLimit = 64 << 20
// JournalSizeLimitMB — тот же предел для строки старта: владелец обязан видеть
// в логе, с какими параметрами обслуживается журнал, а число живёт здесь.
const JournalSizeLimitMB = journalSizeLimit >> 20
// dsn собирает строку подключения: WAL для параллельного чтения во время
// записи, busy_timeout — чтобы конкурентная запись ждала, а не падала.
//
// `_txlock=immediate` обязателен, а не «на всякий случай»: запись объекта —
// это read-modify-write, и при отложенной блокировке повышение с чтения на
// запись после уже прочитанного снимка даёт SQLITE_BUSY_SNAPSHOT, которого
// busy_timeout не покрывает. Измерено на восьми писателях: 242 успешных
// слияния из 800 против 800 из 800.
//
// `journal_size_limit` — см. константу выше.
func dsn(path string) string {
q := url.Values{}
q.Add("_pragma", "journal_mode(WAL)")
q.Add("_pragma", "busy_timeout(5000)")
q.Add("_pragma", "foreign_keys(on)")
q.Add("_pragma", fmt.Sprintf("journal_size_limit(%d)", journalSizeLimit))
q.Add("_txlock", "immediate")
return "file:" + path + "?" + q.Encode()
}
// readOnlyDSN — подключение только для чтения.
//
// `journal_mode` здесь не задаётся: сменить его на read-only соединении нельзя,
// а читать базу в режиме WAL это не мешает. `_txlock=immediate` тоже не нужен —
// он лечит повышение блокировки с чтения на запись, которого здесь не бывает.
func readOnlyDSN(path string) string {
q := url.Values{}
q.Add("mode", "ro")
q.Add("_pragma", "busy_timeout(5000)")
return "file:" + path + "?" + q.Encode()
}
// migrate накатывает миграции.
//
// Через Provider, а не через пакетные функции: goose.SetBaseFS и
// goose.SetDialect пишут глобальное состояние пакета, и два одновременных
// Open дают гонку — её ловит детектор. Сегодня открытие одно, на старте, но
// пересборка витрины откроет второе, и хранилище не должно зависеть от того,
// что вызывающий этого не сделает.
func migrate(db *sqlx.DB) error {
ctx := context.Background()
inDB, inBinary, err := readSchemaVersion(ctx, db)
if err != nil {
return err
}
if inDB > inBinary {
return fmt.Errorf("%w: база %d, бинарь %d", ErrSchemaMismatch, inDB, inBinary)
}
p, err := newProvider(db)
if err != nil {
return err
}
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 совпадает с хронологией.
func Now() time.Time { return time.Now().UTC().Truncate(time.Second) }
// FormatTime приводит время к каноническому виду хранения: RFC 3339, UTC.
func FormatTime(t time.Time) string { return t.UTC().Format(time.RFC3339) }
// ParseTime разбирает метку времени из БД.
func ParseTime(s string) (time.Time, error) {
t, err := time.Parse(time.RFC3339, s)
if err != nil {
return time.Time{}, fmt.Errorf("parse time %q: %w", s, err)
}
return t.UTC(), nil
}