Files
NaviWatcher/internal/database/database.go
Vladimir Zagainov b0f69d3a4f feat: implement rate limiting and caching layer for MusicBrainz provider
- Add golang.org/x/time/rate dependency for token-bucket rate limiting
- Replace custom channel-based rate limiter with rate.NewLimiter(1, 1)
- Add context.Context support to doGet for cancellation
- Add cached_at column to external_releases via migration 005
- Implement cache hit/miss queries with TTL-based filtering
- Add CacheStats type for tracking cached RGIDs
- Update ExternalRelease struct with CachedAt field
- Add rate limiting tests (1 req/sec enforcement, burst behavior)
- Add cache tests (hit, miss, expired, mixed, empty artist)
- Update migration count test for new migration
2026-05-26 12:15:28 +03:00

198 lines
5.2 KiB
Go

package database
import (
"database/sql"
"fmt"
"time"
_ "github.com/mattn/go-sqlite3"
)
// DB wraps sql.DB with migration support.
type DB struct {
conn *sql.DB
}
// New opens a SQLite database at dbPath and runs schema migrations.
func New(dbPath string) (*DB, error) {
conn, err := sql.Open("sqlite3", dbPath)
if err != nil {
return nil, fmt.Errorf("open database: %w", err)
}
// Enable WAL mode for better concurrent read performance.
if _, err := conn.Exec("PRAGMA journal_mode=WAL"); err != nil {
conn.Close()
return nil, fmt.Errorf("set WAL mode: %w", err)
}
// Enable foreign key enforcement.
if _, err := conn.Exec("PRAGMA foreign_keys=ON"); err != nil {
conn.Close()
return nil, fmt.Errorf("enable foreign keys: %w", err)
}
// Set busy timeout to handle concurrent write contention.
if _, err := conn.Exec("PRAGMA busy_timeout=5000"); err != nil {
conn.Close()
return nil, fmt.Errorf("set busy timeout: %w", err)
}
db := &DB{conn: conn}
if err := db.migrate(); err != nil {
conn.Close()
return nil, fmt.Errorf("migrate: %w", err)
}
return db, nil
}
// Close closes the database connection.
func (db *DB) Close() error {
return db.conn.Close()
}
// Conn returns the underlying sql.DB for use by other packages.
func (db *DB) Conn() *sql.DB {
return db.conn
}
// Begin starts a new database transaction.
func (db *DB) Begin() (*sql.Tx, error) {
return db.conn.Begin()
}
// migrate runs all pending schema migrations in order.
func (db *DB) migrate() error {
// Create the migrations tracking table first, unconditionally.
if _, err := db.conn.Exec(`CREATE TABLE IF NOT EXISTS _migrations (
version INTEGER PRIMARY KEY,
name TEXT NOT NULL,
applied_at DATETIME DEFAULT CURRENT_TIMESTAMP
);`); err != nil {
return fmt.Errorf("create migrations table: %w", err)
}
migrations := []struct {
name string
sql string
}{
{
name: "001_create_artist_settings",
sql: `CREATE TABLE IF NOT EXISTS artist_settings (
id TEXT PRIMARY KEY,
name TEXT NOT NULL,
ignore_singles BOOLEAN DEFAULT 0,
ignore_compilations BOOLEAN DEFAULT 0,
monitored BOOLEAN DEFAULT 1
);`,
},
{
name: "002_create_external_releases",
sql: `CREATE TABLE IF NOT EXISTS external_releases (
rgid TEXT PRIMARY KEY,
artist_id TEXT NOT NULL REFERENCES artist_settings(id),
title TEXT NOT NULL,
type TEXT,
release_date TEXT,
is_ignored BOOLEAN DEFAULT 0
);`,
},
{
name: "003_create_local_albums",
sql: `CREATE TABLE IF NOT EXISTS local_albums (
id TEXT PRIMARY KEY,
artist_id TEXT NOT NULL REFERENCES artist_settings(id),
title TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_local_albums_artist_id ON local_albums(artist_id);`,
},
{
name: "004_create_notifications_sent",
sql: `CREATE TABLE IF NOT EXISTS notifications_sent (
rgid TEXT NOT NULL REFERENCES external_releases(rgid),
sent_at DATETIME DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (rgid, sent_at)
);`,
},
{
name: "005_add_cached_at_to_external_releases",
sql: `ALTER TABLE external_releases ADD COLUMN cached_at DATETIME;`,
},
}
for _, m := range migrations {
applied, err := db.isMigrationApplied(m.name)
if err != nil {
return fmt.Errorf("check migration %s: %w", m.name, err)
}
if applied {
continue
}
tx, err := db.conn.Begin()
if err != nil {
return fmt.Errorf("begin transaction for migration %s: %w", m.name, err)
}
if _, err := tx.Exec(m.sql); err != nil {
tx.Rollback()
return fmt.Errorf("apply migration %s: %w", m.name, err)
}
if _, err := tx.Exec("INSERT INTO _migrations (name) VALUES (?)", m.name); err != nil {
tx.Rollback()
return fmt.Errorf("record migration %s: %w", m.name, err)
}
if err := tx.Commit(); err != nil {
return fmt.Errorf("commit migration %s: %w", m.name, err)
}
}
return nil
}
// isMigrationApplied checks whether a migration with the given name has already been applied.
func (db *DB) isMigrationApplied(name string) (bool, error) {
var count int
err := db.conn.QueryRow("SELECT COUNT(*) FROM _migrations WHERE name = ?", name).Scan(&count)
if err != nil {
return false, err
}
return count > 0, nil
}
// ArtistSettings represents a row in the artist_settings table.
type ArtistSettings struct {
ID string `json:"id"`
Name string `json:"name"`
IgnoreSingles bool `json:"ignore_singles"`
IgnoreCompilations bool `json:"ignore_compilations"`
Monitored bool `json:"monitored"`
}
// LocalAlbum represents a row in the local_albums table.
type LocalAlbum struct {
ID string `json:"id"`
ArtistID string `json:"artist_id"`
Title string `json:"title"`
}
// ExternalRelease represents a row in the external_releases table.
type ExternalRelease struct {
RGID string `json:"rgid"`
ArtistID string `json:"artist_id"`
Title string `json:"title"`
Type string `json:"type"`
ReleaseDate string `json:"release_date"`
IsIgnored bool `json:"is_ignored"`
CachedAt time.Time `json:"cached_at"`
}
// NotificationSent represents a row in the notifications_sent table.
type NotificationSent struct {
RGID string `json:"rgid"`
SentAt time.Time `json:"sent_at"`
}