/
codespawn
/
tgstream
Обзор
Документация
Войти
/
codespawn
/
tgstream
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
16
CI/CD
Аналитика
Безопасность
master
internal/adapter/sqlite/sqlite.go
377 строк
13 KB
codespawn
fix(sqlite): serialize writers to prevent SQLITE_BUSY under concurrency
29 июл 2026, 20:00
29 июл 2026, 20:00
d3172ee
Код
Авторство
О чём код?
// Package sqlite реализует порт port.Store как L2-хранилище: персистентный // источник истины для истории постов (AGENTS.md §6). Реализован на чистом // Go-драйвере modernc.org/sqlite (без CGO), что упрощает сборку и `-race`. // // Модель данных (AGENTS.md §6): // - таблицы channels и posts; upsert идемпотентен по составному ключу // (channel, id) — повторный скрап добавляет новые посты и обновляет // существующие (Views/EditedAt/текст/медиа); // - Channel.FetchedAt хранится как RFC3339 — основа TTL для FeedService; // - Media[] сериализуется в JSON-колонку (драйверу/адаптеру разрешён // encoding/json; домен/port чисты по depguard). package sqlite import ( "context" "database/sql" "encoding/json" "errors" "fmt" "strconv" "sync" "time" "gitverse.ru/codespawn/tgstream/internal/domain" "gitverse.ru/codespawn/tgstream/internal/port" _ "modernc.org/sqlite" // регистрация драйвера "sqlite" ) // Store — реализация port.Store поверх SQLite. type Store struct { db *sql.DB // writeMu сериализует писателей. SQLite под WAL допускает ровно одного // писателя; при нескольких соединениях в пуле (SetMaxOpenConns не // ограничен) и busy_timeout=5000ms конкурентные записи в проде падали с // SQLITE_BUSY — тяжёлая транзакция (UpsertPosts с батчем постов) // удерживала write-lock дольше таймаута ожидания других писателей. // Мьютекс делает писателя единственным на уровне приложения; busy_timeout // остаётся как defence-in-depth. Чтения под WAL конкурентны и мьютекс // не трогают. writeMu sync.Mutex } // Config настраивает хранилище. type Config struct { // Path — путь к файлу БД; "" или ":memory:" → изолированная in-memory БД. // Для :memory: пул сводится к одному соединению (SetMaxOpenConns(1)), иначе // разные соединения получили бы РАЗНЫЕ пустые in-memory БД (классический // gotcha database/sql + :memory:). Path string // BusyTimeout — сколько (мс) SQLite-писатель ждёт release write-lock'а // прежде чем вернуть SQLITE_BUSY (передаётся как busy_timeout в DSN). // По умолчанию (<=0) — 5000ms. Это per-connection настройка, поэтому // живёт в DSN, а не в одном Exec (см. комментарий в New). Defence-in-depth // поверх writeMu: при корректной блокировке писателей таймаут не нужен, // но страхует от будущих write-путей, забывших про мьютекс. BusyTimeout time.Duration } // New открывает БД, применяет прагмы (вежливая конкуренция писателей) и схему. func New(ctx context.Context, cfg Config) (*Store, error) { path := cfg.Path if path == "" { path = ":memory:" } isMemory := path == ":memory:" busyMs := int(cfg.BusyTimeout / time.Millisecond) if busyMs <= 0 { busyMs = 5000 } var dsn string if isMemory { // Единое соединение = единая приватная in-memory БД на этот Store // (иначе разные соединения получили бы РАЗНЫЕ пустые БД — классический // gotcha database/sql + :memory:). dsn = ":memory:" } else { // busy_timeout — НАСТРОЙКА СОЕДИНЕНИЯ (а не БД): Exec-прагма применится // лишь к одному соединению пула, остальные получат 0 и тут же упадут с // SQLITE_BUSY. Поэтому кладётся в DSN (_pragma) — так он действует на // КАЖДОЕ соединение. journal_mode, напротив, хранится в файле БД — // его задаём одним Exec ниже. dsn = "file:" + path + "?_pragma=busy_timeout(" + strconv.Itoa(busyMs) + ")" } db, err := sql.Open("sqlite", dsn) if err != nil { return nil, fmt.Errorf("sqlite: open %q: %w", path, err) } if isMemory { db.SetMaxOpenConns(1) } else if _, err := db.ExecContext(ctx, `PRAGMA journal_mode=WAL;`); err != nil { // WAL: конкурентные чтения + один писатель. Свойство файла БД, одним Exec достаточно. _ = db.Close() return nil, fmt.Errorf("sqlite: pragma journal_mode: %w", err) } s := &Store{db: db} if err := s.migrate(ctx); err != nil { _ = db.Close() return nil, err } return s, nil } // Close освобождает ресурсы БД. func (s *Store) Close() error { if s.db == nil { return nil } return s.db.Close() } const schema = ` CREATE TABLE IF NOT EXISTS channels ( name TEXT PRIMARY KEY, title TEXT NOT NULL DEFAULT '', description TEXT NOT NULL DEFAULT '', subscribers INTEGER NOT NULL DEFAULT 0, fetched_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS posts ( channel TEXT NOT NULL, id TEXT NOT NULL, title TEXT NOT NULL DEFAULT '', text TEXT NOT NULL DEFAULT '', html TEXT NOT NULL DEFAULT '', published_at TEXT NOT NULL, edited_at TEXT, views INTEGER NOT NULL DEFAULT 0, url TEXT NOT NULL DEFAULT '', media TEXT NOT NULL DEFAULT '[]', PRIMARY KEY (channel, id) ) WITHOUT ROWID; CREATE INDEX IF NOT EXISTS idx_posts_channel_time ON posts(channel, published_at DESC); ` func (s *Store) migrate(ctx context.Context) error { if _, err := s.db.ExecContext(ctx, schema); err != nil { return fmt.Errorf("sqlite: migrate: %w", err) } return nil } // GetChannel возвращает метаинформацию канала. Неизвестный → ErrChannelNotFound. func (s *Store) GetChannel(ctx context.Context, name string) (domain.Channel, error) { var ( ch domain.Channel title string description string subscribers int fetchedAtStr string ) err := s.db.QueryRowContext(ctx, ` SELECT name, title, description, subscribers, fetched_at FROM channels WHERE name = ?`, name, ).Scan(&ch.Name, &title, &description, &subscribers, &fetchedAtStr) if errors.Is(err, sql.ErrNoRows) { return domain.Channel{}, domain.ErrChannelNotFound } if err != nil { return domain.Channel{}, fmt.Errorf("sqlite: get channel %q: %w", name, err) } fetchedAt, err := time.Parse(time.RFC3339Nano, fetchedAtStr) if err != nil { return domain.Channel{}, fmt.Errorf("sqlite: parse fetched_at %q: %w", fetchedAtStr, err) } ch.Title = title ch.Description = description ch.Subscribers = subscribers ch.FetchedAt = fetchedAt return ch, nil } // GetPosts возвращает посты канала: от новых к старым, с фильтром Since (строго // новее) и ограничением Limit (Limit<=0 — без ограничения). func (s *Store) GetPosts(ctx context.Context, channel string, q port.PostQuery) ([]domain.Post, error) { query := `SELECT id, channel, title, text, html, published_at, edited_at, views, url, media FROM posts WHERE channel = ?` args := []any{channel} if !q.Since.IsZero() { query += ` AND published_at > ?` args = append(args, q.Since.UTC().Format(time.RFC3339Nano)) } query += ` ORDER BY published_at DESC` if q.Limit > 0 { query += ` LIMIT ?` args = append(args, q.Limit) } rows, err := s.db.QueryContext(ctx, query, args...) if err != nil { return nil, fmt.Errorf("sqlite: get posts %q: %w", channel, err) } defer func() { _ = rows.Close() }() var out []domain.Post for rows.Next() { var ( p domain.Post title string text string html string publishedStr string edited sql.NullString views int url string mediaJSON string ) if err := rows.Scan( &p.ID, &p.Channel, &title, &text, &html, &publishedStr, &edited, &views, &url, &mediaJSON, ); err != nil { return nil, fmt.Errorf("sqlite: scan post: %w", err) } publishedAt, err := time.Parse(time.RFC3339Nano, publishedStr) if err != nil { return nil, fmt.Errorf("sqlite: parse published_at %q: %w", publishedStr, err) } p.Title = title p.Text = text p.HTML = html p.PublishedAt = publishedAt p.Views = views p.URL = url if edited.Valid { t, err := time.Parse(time.RFC3339Nano, edited.String) if err != nil { return nil, fmt.Errorf("sqlite: parse edited_at %q: %w", edited.String, err) } p.EditedAt = &t } media, err := unmarshalMedia(mediaJSON) if err != nil { return nil, fmt.Errorf("sqlite: unmarshal media: %w", err) } p.Media = media out = append(out, p) } if err := rows.Err(); err != nil { return nil, fmt.Errorf("sqlite: rows: %w", err) } return out, nil } // UpsertChannel перезаписывает метаинформацию канала (апсерт по name). func (s *Store) UpsertChannel(ctx context.Context, ch domain.Channel) error { s.writeMu.Lock() defer s.writeMu.Unlock() _, err := s.db.ExecContext(ctx, ` INSERT INTO channels (name, title, description, subscribers, fetched_at) VALUES (?, ?, ?, ?, ?) ON CONFLICT(name) DO UPDATE SET title = excluded.title, description = excluded.description, subscribers = excluded.subscribers, fetched_at = excluded.fetched_at`, ch.Name, ch.Title, ch.Description, ch.Subscribers, ch.FetchedAt.UTC().Format(time.RFC3339Nano), ) if err != nil { return fmt.Errorf("sqlite: upsert channel %q: %w", ch.Name, err) } return nil } const upsertPostSQL = ` INSERT INTO posts (channel, id, title, text, html, published_at, edited_at, views, url, media) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(channel, id) DO UPDATE SET title = excluded.title, text = excluded.text, html = excluded.html, published_at = excluded.published_at, edited_at = excluded.edited_at, views = excluded.views, url = excluded.url, media = excluded.media` // UpsertPosts добавляет/обновляет посты канала одной транзакцией (идемпотентно // по (channel, id)). Пустой слайс — no-op. func (s *Store) UpsertPosts(ctx context.Context, channel string, posts []domain.Post) error { if len(posts) == 0 { return nil } s.writeMu.Lock() defer s.writeMu.Unlock() tx, err := s.db.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("sqlite: begin: %w", err) } defer func() { _ = tx.Rollback() }() // no-op после Commit stmt, err := tx.PrepareContext(ctx, upsertPostSQL) if err != nil { return fmt.Errorf("sqlite: prepare upsert post: %w", err) } defer func() { _ = stmt.Close() }() for _, p := range posts { mediaJSON, err := marshalMedia(p.Media) if err != nil { return fmt.Errorf("sqlite: marshal media for %s/%s: %w", channel, p.ID, err) } var edited any if p.EditedAt != nil { edited = p.EditedAt.UTC().Format(time.RFC3339Nano) } if _, err := stmt.ExecContext(ctx, channel, p.ID, p.Title, p.Text, p.HTML, p.PublishedAt.UTC().Format(time.RFC3339Nano), edited, p.Views, p.URL, mediaJSON, ); err != nil { return fmt.Errorf("sqlite: upsert post %s/%s: %w", channel, p.ID, err) } } if err := tx.Commit(); err != nil { return fmt.Errorf("sqlite: commit: %w", err) } return nil } // Trim оставляет не более max свежих постов канала, отсортированных по // published_at DESC (новее→старее); остальные удаляются. max<=0 — no-op. // Используется composition root для соблюдения STORE_MAX_PER_CHANNEL // (AGENTS.md §10): история канала не растёт безусловно при повторных // скрапах. Вызывается адаптером-декоратором в app-слое после UpsertPosts. func (s *Store) Trim(ctx context.Context, channel string, max int) error { if max <= 0 { return nil } s.writeMu.Lock() defer s.writeMu.Unlock() _, err := s.db.ExecContext(ctx, ` DELETE FROM posts WHERE channel = ? AND id NOT IN ( SELECT id FROM posts WHERE channel = ? ORDER BY published_at DESC LIMIT ? )`, channel, channel, max) if err != nil { return fmt.Errorf("sqlite: trim %q to %d: %w", channel, max, err) } return nil } // --- сериализация Media --- func marshalMedia(m []domain.Media) (string, error) { if len(m) == 0 { return "[]", nil } b, err := json.Marshal(m) if err != nil { return "", err } return string(b), nil } func unmarshalMedia(s string) ([]domain.Media, error) { if s == "" || s == "[]" { return nil, nil } var m []domain.Media if err := json.Unmarshal([]byte(s), &m); err != nil { return nil, err } return m, nil } // Гарантия реализации порта во время компиляции. var _ port.Store = (*Store)(nil)