package store import ( "context" "database/sql" "errors" "fmt" "strconv" "sync" "git.vakhrushev.me/av/healthlog/internal/ident" ) // dataVersionQuery — счётчик коммитов, видимых соединению. // // SQLite меняет его, когда изменения в базу закоммитило ДРУГОЕ соединение, и // оставляет неизменным для коммитов самого соединения. Отсюда требование к // щупу: он ничего не пишет. const dataVersionQuery = `PRAGMA data_version` // versionProbe — закреплённое соединение, с которого читается версия витрины. // // Соединение своё, а не из пула, по измеренной причине: значение `data_version` // локально для соединения. На одном и том же состоянии базы два соединения // одного пула отвечают разными числами, а любое СВЕЖЕЕ соединение отвечает // одним и тем же значением независимо от содержимого базы. Версия из пула // поэтому давала бы не только ложную инвалидацию (разные метки на // неизменившейся витрине — не страшно), но и одинаковые метки на разных // состояниях — то есть `304` на изменившиеся данные. // // Поколение выдаётся соединению и меняется вместе с ним. Без него метка не // переживала бы рестарт: счётчик после переоткрытия начинается заново, и одно и // то же значение до и после означало бы разные состояния витрины. Побочная // выгода: поколение меняется и при выкатке нового бинаря, так что смена ФОРМЫ // ответа при неизменившихся данных тоже обнуляет метки клиентов. // // Мьютекс здесь не только ради поля: `sql.Conn` не предназначен для // одновременного использования из нескольких горутин, а запросов чтения бывает // сколько угодно. Запрос при этом мгновенный — очереди на нём не образуется. // На щупе выполняется РОВНО ОДИН вид запроса и только через // `QueryRowContext(...).Scan(...)`. `QueryContext` и `BeginTx` на нём не // зовутся никогда: незакрытые `Rows` или открытая транзакция удержали бы // читающий снимок до конца жизни процесса — а щуп единственное долгоживущее // соединение процесса. Тогда пассивный чекпойнт перестал бы продвигаться // вовсе, и вторая половина задачи убила бы первую при полностью исправном // обслуживании. Заодно повисли бы все читающие запросы: щуп у них общий. type versionProbe struct { mu sync.Mutex conn *sql.Conn generation string // closed — хранилище закрыто, щупа больше не будет. Без флага гонка // «запрос версии против Close» воскресила бы соединение уже после закрытия // пула, и последнее соединение к базе осталось бы открытым: SQLite не // сделал бы финальный чекпойнт, а рядом с базой остался бы `-wal`. closed bool } // StateVersion отдаёт версию витрины: метку, которая меняется при любом // коммите в базу и не меняется, пока коммитов не было. // // Имя не называет PRAGMA намеренно: метка это пара «поколение + счётчик», а не // голое значение `data_version`, и область её сравнимости задаёт хранилище, а // не SQLite. Со «версией схемы» (ErrSchemaMismatch) она не пересекается ничем. // // Равные метки означают, что между их снятием в базу никто ничего не записал, — // на этом и держится условный запрос читающих маршрутов. Обратное неверно: // метка меняется от любой записи, включая учёт доставки, витрину не менявшей. // Это ложная инвалидация, то есть безопасная сторона. func (s *Store) StateVersion(ctx context.Context) (string, error) { s.probe.mu.Lock() defer s.probe.mu.Unlock() // Две попытки, а не цикл: единственная восстановимая беда — умершее // соединение, и лечится она ровно одним пересозданием. Повторять дальше // значило бы ходить по кругу за отказом, который не в соединении. var lastErr error for range 2 { conn, generation, err := s.probe.acquire(ctx, s.db.DB) if err != nil { return "", err } var counter int64 err = conn.QueryRowContext(ctx, dataVersionQuery).Scan(&counter) if err == nil { return generation + "-" + strconv.FormatInt(counter, 10), nil } lastErr = err // Непригодность соединения — это ТОЛЬКО `ErrConnDone`, и список узок // намеренно. Измерено на этом драйвере: отмена контекста запроса щуп не // убивает — следующий запрос на нём проходит; непригодным соединение // становится после явного закрытия. Считать смертью щупа любую ошибку // нельзя: занятость базы и обрыв запроса клиентом (обычные события, для // которых в проекте заведён `Transient`) меняли бы поколение, и все // потребители получали бы полный ответ вместо `304` — то есть механизм // схлопывался бы ровно под нагрузкой, ради которой заведён. if !errors.Is(err, sql.ErrConnDone) { break } // Вместе с соединением выбрасывается и поколение: переиспользовать его // нельзя — счётчик у нового соединения начнётся заново, и старая метка // совпала бы с новой на другом состоянии. s.probe.release() } return "", fmt.Errorf("read data version: %w", lastErr) } // VersionedRead выполняет чтение и отдаёт версию витрины, которой это чтение // подписано. Пустая версия означает «подписать нечем» — не отказ. // // Правило живёт здесь, в одном экземпляре, потому что нарушить его можно ровно // одним способом и этот способ опасен: версия, снятая ПОСЛЕ чтения, пометила бы // устаревший снимок свежей меткой и заперла бы клиента на нём навсегда. Версия, // снятая только ДО, допускает два разных ответа под одной меткой. Поэтому проба // делается дважды, а метка выдаётся, только если между пробами в базу никто не // коммитил. // // Read API точек и MCP заявлены потребителями той же машинерии: вторая её // реализация «по образцу» отличалась бы от первой ровно на этот порядок, и ни // один тест каталога этого не увидел бы. // // Отказ пробы версией не является и запрос не роняет: читающий маршрут // деградирует до полного ответа, а не до отказа. Настоящий отказ базы всплывёт // самим чтением, которое идёт следом, и будет назван один раз им. func (s *Store) VersionedRead(ctx context.Context, read func(context.Context) error) (string, error) { before, probeErr := s.StateVersion(ctx) if err := read(ctx); err != nil { return "", err } if probeErr != nil { return "", nil } after, err := s.StateVersion(ctx) if err != nil || after != before { return "", nil } return before, nil } // acquire отдаёт закреплённое соединение, заводя его при первом обращении. // Вызывается под мьютексом. // // Лениво, а не при открытии базы: щуп нужен читающим маршрутам, а `reindex` и // утилиты учёта открывают ту же базу и версию не спрашивают ни разу. func (p *versionProbe) acquire(ctx context.Context, db *sql.DB) (*sql.Conn, string, error) { if p.conn != nil { return p.conn, p.generation, nil } if p.closed { return nil, "", ErrClosed } conn, err := db.Conn(ctx) if err != nil { return nil, "", fmt.Errorf("pin version probe connection: %w", err) } p.conn = conn p.generation = ident.NewID() return p.conn, p.generation, nil } // release закрывает закреплённое соединение и забывает поколение. // Вызывается под мьютексом. func (p *versionProbe) release() { if p.conn == nil { return } // Ошибка закрытия непригодного соединения ничего не меняет: следующий // заход возьмёт новое. _ = p.conn.Close() p.conn = nil p.generation = "" } // closeProbe закрывает щуп. Отдельно от пула и ДО него: закреплённое // соединение переживает `db.Close()` и продолжает отвечать на запросы // (проверено), то есть само по себе не закрывается ничем. func (s *Store) closeProbe() { s.probe.mu.Lock() defer s.probe.mu.Unlock() s.probe.closed = true s.probe.release() }