- рядом с воркером свёртки живёт горутина, раз в минуту разбирающая журнал пассивным чекпойнтом; «журнал не разбирается» видно строкой владельцу, а не только по `df`. Признак — пара чисел, а не флаг занятости: тот молчит под удерживаемым читателем (`busy=0` при 6256 страницах и пяти перенесённых), а при занятой блокировке отдаёт `-1` вместо ответа, и `-1 >= -1` читалось бы как «разобрано целиком» - каталог отвечает `304` на `If-None-Match`, не открывая снимок витрины. Метка собрана из всего, от чего зависит ответ: версии витрины (`data_version` с закреплённого соединения плюс поколение — значение локально для соединения и не переживает переоткрытия), горизонта измерения и области действия ресурса. Версия снимается до и после сборки: снятая после пометила бы устаревший снимок свежим номером - предел и дедлайн ответа отложены в задачу Read API точек вместе с измеренной ценой первого запроса; попутно починен флаки-тест чужой задачи, искавший значение точки в сыром буфере записи лога
273 lines
15 KiB
Go
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
|
|
}
|