// Package sqlite — хранилище сервиса: база на своей схеме и файлы записей своим // каталогом. // // Пакет назван по драйверу, а не по роли: соседи в `internal/adapter` названы // тем же способом — `converter`, `metaviewer`, `recognizer`, — и «repo/sqlite» // читается как «репозитории поверх SQLite» без знания кода. package sqlite import ( "context" "database/sql" "errors" "fmt" "net/url" "os" "path/filepath" "strconv" // Драйвер регистрируется загрузкой пакета. CGO ему не нужен — этим он и // выбран: сборка бинарника остаётся без компилятора C. _ "modernc.org/sqlite" ) // driverName — имя, под которым драйвер регистрируется в `database/sql`. const driverName = "sqlite" // DatabaseFile — имя файла базы в каталоге данных. Рядом с ним драйвер кладёт // журнал упреждающей записи и его указатель, поэтому каталог данных занят базой // целиком, а не одним файлом. const DatabaseFile = "transcriber.db" // Settings — числа, которыми настраивается база. Оба приходят настройкой, а не // константой кода: крутят их при одном и том же отказе — «база занята» под // несколькими воркерами, — и подбор ответа на такой отказ не должен требовать // пересборки образа. type Settings struct { // BusyTimeoutMs — сколько ждать занятую базу, миллисекунды. BusyTimeoutMs int // ReadConnections — сколько соединений держит читающий пул. ReadConnections int } // Validate проверяет числа базы. Ноль и отрицательное — опечатка, а не режим: // нулевое ожидание отдаёт «база занята» первому же воркеру, а нулевой пул // чтения означает пул без предела, то есть настройку, которой не управляют. func (s Settings) Validate() error { if s.BusyTimeoutMs <= 0 { return errors.New("storage: ожидание занятой базы задаётся положительным числом миллисекунд") } if s.ReadConnections <= 0 { return errors.New("storage: число соединений читающего пула задаётся положительным числом") } return nil } // Обращения к базе идут с **собственным** контекстом, а не с контекстом // запроса, и это решение, а не недосмотр. Репозитории отменять нечего: операции // местные и короткие, а единственное ожидание — занятая база — задано числом. За // отмену при этом платили бы дважды: шаг, прерванный остановкой сервиса, // перестал бы освобождать захват и писать причину остановки — то есть отмена // ломала бы ровно ту уборку, ради которой она и делается. // // Отмена, которой сервис распоряжается по-настоящему, доходит туда, где она // стоит денег и времени: до `ffmpeg` и до платного распознавания. // DB — база сервиса двумя пулами. // // Пишущий пул держит **одно** соединение: драйвер пишет единственным // соединением, и несколько воркеров, пришедших писать разом мимо этого правила, // получают отказ по занятости — на записи результата шага, то есть после // оплаченной работы. Пул с одним соединением обращает их в очередь. // // Читающий пул отдельный: в журнале упреждающей записи читатели не мешают // писателю, и список записей не ждёт, пока конвейер сохранит свой шаг. type DB struct { // writer — единственное пишущее соединение. Через него идёт всякая // операция, которая читает состояние и следом его пишет: транзакцию, // начатую на читающем соединении, SQLite до пишущей не повышает и отвечает // отказом по занятости немедленно — заданное числом ожидание такой отказ не // лечит, ждать там нечего. writer *sql.DB // reader — пул чтения. reader *sql.DB } // Writer отдаёт пишущее соединение. func (db *DB) Writer() *sql.DB { return db.writer } // Reader отдаёт читающий пул. func (db *DB) Reader() *sql.DB { return db.reader } // Open открывает базу в каталоге данных, заводя каталог, если его ещё нет. // // Настройки соединения задаются **строкой подключения обоих пулов**, а не // запросом после открытия. Соблюдение внешних ключей в SQLite — настройка // соединения, а не базы, и по умолчанию она выключена; пул раздаёт соединения и // заводит новые по мере надобности, поэтому запрос, выполненный один раз, // настроил бы одно соединение из многих, а остальные остались бы с умолчанием — // молча. func Open(dataDir string, settings Settings) (*DB, error) { if err := settings.Validate(); err != nil { return nil, err } if err := os.MkdirAll(dataDir, 0o750); err != nil { return nil, fmt.Errorf("failed to create data directory: %w", err) } path := filepath.Join(dataDir, DatabaseFile) // Пишущее соединение начинает транзакцию сразу пишущей (`immediate`): // операция, которая читает и следом пишет, иначе взяла бы читающую // транзакцию и упёрлась бы в отказ при первой же записи. writer, err := open(path, settings, "immediate") if err != nil { return nil, err } writer.SetMaxOpenConns(1) writer.SetMaxIdleConns(1) reader, err := open(path, settings, "deferred") if err != nil { return nil, errors.Join(fmt.Errorf("failed to open read pool: %w", err), writer.Close()) } reader.SetMaxOpenConns(settings.ReadConnections) reader.SetMaxIdleConns(settings.ReadConnections) db := &DB{writer: writer, reader: reader} // Пробное обращение делается сразу: `sql.Open` соединения не открывает, и // негодная строка подключения вылезла бы не на старте, а на первом запросе — // то есть отказом каждого запроса вместо одной строки о причине. if err := writer.PingContext(context.Background()); err != nil { return nil, errors.Join(fmt.Errorf("failed to open database: %w", err), db.Close()) } return db, nil } // open заводит один пул с общими настройками соединения. func open(path string, settings Settings, txlock string) (*sql.DB, error) { query := url.Values{} query.Add("_pragma", "busy_timeout("+strconv.Itoa(settings.BusyTimeoutMs)+")") query.Add("_pragma", "journal_mode(WAL)") query.Add("_pragma", "foreign_keys(1)") query.Set("_txlock", txlock) db, err := sql.Open(driverName, "file:"+path+"?"+query.Encode()) if err != nil { return nil, fmt.Errorf("failed to open database: %w", err) } return db, nil } // Close закрывает оба пула. Повторный вызов паники не даёт: закрытие уже // закрытого пула отказом не считается. func (db *DB) Close() error { var errs []error if db.reader != nil { if err := db.reader.Close(); err != nil { errs = append(errs, fmt.Errorf("failed to close read pool: %w", err)) } db.reader = nil } if db.writer != nil { if err := db.writer.Close(); err != nil { errs = append(errs, fmt.Errorf("failed to close write pool: %w", err)) } db.writer = nil } return errors.Join(errs...) }