diff --git a/appview/database/migrations.go b/appview/database/migrations.go index a4d2a76..9c8599a 100644 --- a/appview/database/migrations.go +++ b/appview/database/migrations.go @@ -1,186 +1,28 @@ package database import ( - "context" - "crypto/sha256" - "database/sql" - "embed" - "encoding/hex" "fmt" - "io/fs" - "path/filepath" - "sort" - "strings" - "time" "gorm.io/gorm" ) -//go:embed migrations/*.sql -var embeddedMigrations embed.FS - -const advisoryLockID int64 = 838001447 - -type migration struct { - Version string - Name string - SQL string - Checksum string -} - -func RunMigrations(db *gorm.DB) error { - sqlDB, err := db.DB() - if err != nil { - return fmt.Errorf("get sql db: %w", err) - } - - ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second) - defer cancel() - - if err := ensureMigrationsTable(ctx, sqlDB); err != nil { - return err - } - - if _, err := sqlDB.ExecContext(ctx, "SELECT pg_advisory_lock($1)", advisoryLockID); err != nil { - return fmt.Errorf("acquire migration advisory lock: %w", err) - } - defer func() { - _, _ = sqlDB.ExecContext(context.Background(), "SELECT pg_advisory_unlock($1)", advisoryLockID) - }() - - migrations, err := loadMigrationsFromFS(embeddedMigrations) - if err != nil { - return err - } - - for _, m := range migrations { - if err := applyMigration(ctx, sqlDB, m); err != nil { - return err - } - } - return nil -} - -func ensureMigrationsTable(ctx context.Context, db *sql.DB) error { - const stmt = ` -CREATE TABLE IF NOT EXISTS schema_migrations ( - version TEXT PRIMARY KEY, - name TEXT NOT NULL, - checksum TEXT NOT NULL, - applied_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -)` - if _, err := db.ExecContext(ctx, stmt); err != nil { - return fmt.Errorf("create schema_migrations table: %w", err) - } - return nil -} - -func loadMigrationsFromFS(fsys fs.FS) ([]migration, error) { - entries, err := fs.ReadDir(fsys, "migrations") - if err != nil { - return nil, fmt.Errorf("read migrations dir: %w", err) - } - - list := make([]migration, 0, len(entries)) - seenVersions := map[string]string{} - for _, entry := range entries { - if entry.IsDir() { - continue - } - name := entry.Name() - if filepath.Ext(name) != ".sql" { - continue - } - - version := parseMigrationVersion(name) - if version == "" { - return nil, fmt.Errorf("invalid migration filename: %s", name) - } - if prev, exists := seenVersions[version]; exists { - return nil, fmt.Errorf("duplicate migration version %s in %s and %s", version, prev, name) - } - seenVersions[version] = name - - raw, err := fs.ReadFile(fsys, filepath.Join("migrations", name)) - if err != nil { - return nil, fmt.Errorf("read migration %s: %w", name, err) - } - sqlText := strings.TrimSpace(string(raw)) - if sqlText == "" { - return nil, fmt.Errorf("migration %s is empty", name) - } - - sum := sha256.Sum256(raw) - list = append(list, migration{ - Version: version, - Name: name, - SQL: sqlText, - Checksum: hex.EncodeToString(sum[:]), - }) - } - - sort.Slice(list, func(i, j int) bool { - return list[i].Version < list[j].Version - }) - return list, nil -} - -func parseMigrationVersion(name string) string { - base := strings.TrimSuffix(name, filepath.Ext(name)) - if base == "" { - return "" - } - parts := strings.SplitN(base, "_", 2) - version := parts[0] - if version == "" { - return "" - } - for _, ch := range version { - if ch < '0' || ch > '9' { - return "" - } - } - return version -} - -func applyMigration(ctx context.Context, db *sql.DB, m migration) error { - var existingChecksum string - row := db.QueryRowContext(ctx, "SELECT checksum FROM schema_migrations WHERE version = $1", m.Version) - switch err := row.Scan(&existingChecksum); err { - case nil: - if existingChecksum != m.Checksum { - return fmt.Errorf("checksum mismatch for migration %s (%s): expected %s got %s", m.Version, m.Name, existingChecksum, m.Checksum) - } - return nil - case sql.ErrNoRows: - // proceed and apply - default: - return fmt.Errorf("query migration %s: %w", m.Version, err) - } - - tx, err := db.BeginTx(ctx, nil) - if err != nil { - return fmt.Errorf("begin migration tx for %s: %w", m.Name, err) - } - defer func() { - _ = tx.Rollback() - }() - - if _, err := tx.ExecContext(ctx, m.SQL); err != nil { - return fmt.Errorf("execute migration %s: %w", m.Name, err) - } - if _, err := tx.ExecContext( - ctx, - "INSERT INTO schema_migrations(version, name, checksum, applied_at) VALUES ($1, $2, $3, NOW())", - m.Version, - m.Name, - m.Checksum, +func AutoMigrateAll(db *gorm.DB) error { + if err := db.AutoMigrate( + &FirehoseCursor{}, + &Subscription{}, + &Comment{}, + &Recommendation{}, + &PodcastList{}, + &Bookmark{}, + &Profile{}, + &PICache{}, + &PodcastStats{}, + &EpisodeStats{}, + &EpisodeState{}, + &Block{}, + &Report{}, ); err != nil { - 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 fmt.Errorf("auto migrate: %w", err) } return nil } diff --git a/appview/database/migrations/0001_initial.sql b/appview/database/migrations/0001_initial.sql deleted file mode 100644 index 7db81ce..0000000 --- a/appview/database/migrations/0001_initial.sql +++ /dev/null @@ -1,126 +0,0 @@ -CREATE TABLE IF NOT EXISTS firehose_cursor ( - id BIGSERIAL PRIMARY KEY, - seq BIGINT NOT NULL, - updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); - -CREATE TABLE IF NOT EXISTS subscriptions ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - feed_url VARCHAR(2048), - podcast_guid VARCHAR(512), - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_subscriptions_did_rkey ON subscriptions (did, rkey); -CREATE INDEX IF NOT EXISTS idx_subscriptions_did ON subscriptions (did); -CREATE INDEX IF NOT EXISTS idx_subscriptions_feed_id ON subscriptions (feed_id); - -CREATE TABLE IF NOT EXISTS comments ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - at_uri VARCHAR(1024) NOT NULL, - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - episode_id INTEGER NOT NULL CHECK (episode_id > 0), - episode_guid VARCHAR(512), - podcast_guid VARCHAR(512), - text TEXT NOT NULL, - timestamp_s INTEGER, - reply_root VARCHAR(1024), - reply_parent VARCHAR(1024), - facets JSONB, - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_comments_did_rkey ON comments (did, rkey); -CREATE INDEX IF NOT EXISTS idx_comments_did ON comments (did); -CREATE INDEX IF NOT EXISTS idx_comments_episode ON comments (feed_id, episode_id); -CREATE INDEX IF NOT EXISTS idx_comments_reply_root ON comments (reply_root); -CREATE INDEX IF NOT EXISTS idx_comments_at_uri ON comments (at_uri); -CREATE INDEX IF NOT EXISTS idx_comments_timestamp_s ON comments (timestamp_s); - -CREATE TABLE IF NOT EXISTS recommendations ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - episode_id INTEGER NOT NULL CHECK (episode_id > 0), - episode_guid VARCHAR(512), - podcast_guid VARCHAR(512), - text TEXT, - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_recommendations_did_rkey ON recommendations (did, rkey); -CREATE INDEX IF NOT EXISTS idx_recommendations_did ON recommendations (did); -CREATE INDEX IF NOT EXISTS idx_recommendations_episode ON recommendations (feed_id, episode_id); - -CREATE TABLE IF NOT EXISTS podcast_lists ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - name VARCHAR(500) NOT NULL, - description TEXT, - podcasts JSONB NOT NULL DEFAULT '[]'::jsonb, - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_lists_did_rkey ON podcast_lists (did, rkey); -CREATE INDEX IF NOT EXISTS idx_lists_did ON podcast_lists (did); - -CREATE TABLE IF NOT EXISTS bookmarks ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - episode_id INTEGER NOT NULL CHECK (episode_id > 0), - episode_guid VARCHAR(512), - podcast_guid VARCHAR(512), - timestamp_s INTEGER, - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_bookmarks_did_rkey ON bookmarks (did, rkey); -CREATE INDEX IF NOT EXISTS idx_bookmarks_did ON bookmarks (did); -CREATE INDEX IF NOT EXISTS idx_bookmarks_episode ON bookmarks (feed_id, episode_id); -CREATE INDEX IF NOT EXISTS idx_bookmarks_timestamp_s ON bookmarks (timestamp_s); - -CREATE TABLE IF NOT EXISTS profiles ( - did VARCHAR(255) PRIMARY KEY, - display_name VARCHAR(640), - description TEXT, - favorite_genres JSONB, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); - -CREATE TABLE IF NOT EXISTS pi_cache ( - id BIGSERIAL PRIMARY KEY, - cache_key VARCHAR(512) NOT NULL, - response JSONB NOT NULL, - expires_at TIMESTAMPTZ NOT NULL, - created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(), - updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_pi_cache_cache_key ON pi_cache (cache_key); -CREATE INDEX IF NOT EXISTS idx_pi_cache_expires ON pi_cache (expires_at); - -CREATE TABLE IF NOT EXISTS podcast_stats ( - feed_id INTEGER PRIMARY KEY CHECK (feed_id > 0), - subscriber_count INTEGER NOT NULL DEFAULT 0, - comment_count INTEGER NOT NULL DEFAULT 0, - recommendation_count INTEGER NOT NULL DEFAULT 0, - last_updated TIMESTAMPTZ NOT NULL DEFAULT NOW() -); - -CREATE TABLE IF NOT EXISTS episode_stats ( - episode_id INTEGER PRIMARY KEY CHECK (episode_id > 0), - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - comment_count INTEGER NOT NULL DEFAULT 0, - recommendation_count INTEGER NOT NULL DEFAULT 0, - bookmark_count INTEGER NOT NULL DEFAULT 0, - last_updated TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE INDEX IF NOT EXISTS idx_episode_stats_feed ON episode_stats (feed_id); diff --git a/appview/database/migrations/0002_episode_states.sql b/appview/database/migrations/0002_episode_states.sql deleted file mode 100644 index c3400ca..0000000 --- a/appview/database/migrations/0002_episode_states.sql +++ /dev/null @@ -1,21 +0,0 @@ -CREATE TABLE IF NOT EXISTS episode_states ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - feed_id INTEGER NOT NULL CHECK (feed_id > 0), - episode_id INTEGER NOT NULL CHECK (episode_id > 0), - episode_guid VARCHAR(512), - podcast_guid VARCHAR(512), - position_s INTEGER, - duration_s INTEGER, - played BOOLEAN NOT NULL DEFAULT FALSE, - saved BOOLEAN NOT NULL DEFAULT FALSE, - hidden BOOLEAN NOT NULL DEFAULT FALSE, - updated_at VARCHAR(64), - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_episode_states_did_rkey ON episode_states (did, rkey); -CREATE INDEX IF NOT EXISTS idx_episode_states_did ON episode_states (did); -CREATE UNIQUE INDEX IF NOT EXISTS idx_episode_states_episode ON episode_states (did, feed_id, episode_id); -CREATE INDEX IF NOT EXISTS idx_episode_states_saved ON episode_states (did, saved); diff --git a/appview/database/migrations/0003_blocks_reports.sql b/appview/database/migrations/0003_blocks_reports.sql deleted file mode 100644 index 4d91901..0000000 --- a/appview/database/migrations/0003_blocks_reports.sql +++ /dev/null @@ -1,29 +0,0 @@ -CREATE TABLE IF NOT EXISTS blocks ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - at_uri VARCHAR(1024) NOT NULL, - subject_did VARCHAR(255) NOT NULL, - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_blocks_did_rkey ON blocks (did, rkey); -CREATE INDEX IF NOT EXISTS idx_blocks_did ON blocks (did); -CREATE INDEX IF NOT EXISTS idx_blocks_subject ON blocks (subject_did); -CREATE UNIQUE INDEX IF NOT EXISTS idx_blocks_did_subject ON blocks (did, subject_did); -CREATE INDEX IF NOT EXISTS idx_blocks_at_uri ON blocks (at_uri); - -CREATE TABLE IF NOT EXISTS reports ( - id BIGSERIAL PRIMARY KEY, - did VARCHAR(255) NOT NULL, - rkey VARCHAR(512) NOT NULL, - subject_uri VARCHAR(1024) NOT NULL, - reason_type VARCHAR(64) NOT NULL, - reason TEXT, - status VARCHAR(64) NOT NULL DEFAULT 'pending', - created_at VARCHAR(64) NOT NULL, - indexed_at TIMESTAMPTZ NOT NULL DEFAULT NOW() -); -CREATE UNIQUE INDEX IF NOT EXISTS idx_reports_did_rkey ON reports (did, rkey); -CREATE INDEX IF NOT EXISTS idx_reports_did ON reports (did); -CREATE INDEX IF NOT EXISTS idx_reports_status ON reports (status); diff --git a/appview/database/migrations_test.go b/appview/database/migrations_test.go deleted file mode 100644 index 0f2c87c..0000000 --- a/appview/database/migrations_test.go +++ /dev/null @@ -1,49 +0,0 @@ -package database - -import ( - "testing" - "testing/fstest" -) - -func TestLoadMigrationsFromFSSortsByVersion(t *testing.T) { - fsys := fstest.MapFS{ - "migrations/0002_second.sql": {Data: []byte("SELECT 2;")}, - "migrations/0001_first.sql": {Data: []byte("SELECT 1;")}, - } - - migrations, err := loadMigrationsFromFS(fsys) - if err != nil { - t.Fatalf("expected no error, got %v", err) - } - if len(migrations) != 2 { - t.Fatalf("expected 2 migrations, got %d", len(migrations)) - } - if migrations[0].Version != "0001" || migrations[1].Version != "0002" { - t.Fatalf("unexpected order: %#v", migrations) - } - if migrations[0].Checksum == "" || migrations[1].Checksum == "" { - t.Fatal("expected non-empty checksums") - } -} - -func TestLoadMigrationsFromFSRejectsBadFilename(t *testing.T) { - fsys := fstest.MapFS{ - "migrations/invalid_name.sql": {Data: []byte("SELECT 1;")}, - } - - if _, err := loadMigrationsFromFS(fsys); err == nil { - t.Fatal("expected filename validation error") - } -} - -func TestLoadMigrationsFromFSRejectsDuplicateVersions(t *testing.T) { - fsys := fstest.MapFS{ - "migrations/0001_first.sql": {Data: []byte("SELECT 1;")}, - "migrations/0001_again.sql": {Data: []byte("SELECT 1;")}, - "migrations/0002_second.sql": {Data: []byte("SELECT 2;")}, - } - - if _, err := loadMigrationsFromFS(fsys); err == nil { - t.Fatal("expected duplicate version validation error") - } -} diff --git a/appview/server.go b/appview/server.go index 5a89138..e4ef070 100644 --- a/appview/server.go +++ b/appview/server.go @@ -41,8 +41,8 @@ func NewServer(cfg Config) (*Server, error) { return nil, fmt.Errorf("connecting to database: %w", err) } - if err := database.RunMigrations(db); err != nil { - return nil, fmt.Errorf("running database migrations: %w", err) + if err := database.AutoMigrateAll(db); err != nil { + return nil, fmt.Errorf("auto migrating database: %w", err) } piClient := podcastindex.NewClient(cfg.PIKey, cfg.PISecret)