diff --git a/appview/database/BACKUPS.md b/appview/database/BACKUPS.md new file mode 100644 index 0000000..29321ae --- /dev/null +++ b/appview/database/BACKUPS.md @@ -0,0 +1,163 @@ +# AppView Database — Backups and Restore + +Postgres is the single source of truth for every social record indexed from +the AT Protocol firehose. If the database is gone or corrupt, the AppView is +gone. Backups only count once they have been restored successfully, so this +document covers the full round trip. + +## What we rely on + +| Layer | Purpose | Retention | +|---|---|---| +| Railway managed snapshots | Point-in-time recovery for operator error. | Railway default; see dashboard. | +| Nightly logical dump (`pg_dump`) | Portable, auditable, restorable to any Postgres. | 30 days in object storage. | +| Pre-deploy manual snapshot | Roll-back target when a migration goes wrong. | Kept until the following clean deploy. | + +A database without both managed snapshots **and** an off-platform logical dump +is not considered backed up. Railway snapshots are fast to restore but tie you +to Railway; the logical dump is the exit hatch. + +## 1. Railway managed snapshots + +1. **Confirm the snapshot schedule** in the Railway dashboard under + `effem-appview → Database → Backups`. Take a screenshot and paste the cadence + and retention window into this document whenever the policy changes. +2. **Manual snapshot before any deploy that ships a migration file.** From the + Railway CLI: + + ```sh + railway run --service effem-db -- pg_dump --schema-only > /tmp/pre-deploy-schema.sql + # then click "Create snapshot" in the Railway dashboard for the full backup + ``` + + Record the snapshot name in the PR description that ships the migration. +3. **Retention**: trust Railway's policy; do not assume snapshots older than + the window are available. + +## 2. Nightly logical dump + +Run a daily `pg_dump --format=custom` and ship the artifact off-platform. + +### Schedule + +Use a Railway cron service (or equivalent) to run every day at 04:00 UTC: + +```sh +#!/bin/sh +set -euo pipefail + +TS=$(date -u +%Y%m%dT%H%M%SZ) +OUT="/tmp/effem-${TS}.dump" + +pg_dump --format=custom --no-owner --no-privileges \ + "$EFFEM_DATABASE_URL" > "$OUT" + +aws s3 cp "$OUT" "s3://effem-backups/postgres/${TS}.dump" \ + --storage-class STANDARD_IA + +rm -f "$OUT" +``` + +### Retention and lifecycle + +- **30 days** in `STANDARD_IA` on the hot bucket. +- **90 days** in `GLACIER_IR` for compliance lookback. +- Lifecycle managed via the bucket's S3/R2 lifecycle rules, not by hand. + +### Monitoring + +The cron must emit a success/failure event. Wire it to the same alert channel +as the app (Phase 1, Step 6). A silent failing backup is worse than no backup +because it creates false confidence. + +Alert triggers: + +- Last successful upload older than 48 hours. +- Backup file size drops by >50% vs. 7-day rolling average (likely a partial + dump or an empty database pointed at wrong URL). + +## 3. Restore drill (quarterly) + +A backup you have not restored is not a backup. Put this on the calendar for +the first Friday of each quarter. + +### Procedure + +1. Provision a scratch Postgres 16 instance (Docker locally is fine): + + ```sh + docker run --rm -d --name effem-restore-test \ + -e POSTGRES_PASSWORD=restore -e POSTGRES_DB=effem \ + -p 54329:5432 postgres:16-alpine + ``` + +2. Pull the latest dump: + + ```sh + aws s3 cp s3://effem-backups/postgres/$(aws s3 ls s3://effem-backups/postgres/ \ + | sort | tail -1 | awk '{print $4}') /tmp/latest.dump + ``` + +3. Restore: + + ```sh + pg_restore --no-owner --no-privileges --clean --if-exists \ + --dbname="postgres://postgres:restore@localhost:54329/effem" \ + /tmp/latest.dump + ``` + +4. Boot the AppView against the restored database: + + ```sh + EFFEM_DATABASE_URL="postgres://postgres:restore@localhost:54329/effem" \ + go run ./cmd/effem-appview --auth-required=false + ``` + +5. **Confirm `RunMigrations` is a no-op.** The startup log should show no + "applying migration" lines. If it tries to apply `0001_initial.sql`, the + dump was taken from a database that did not yet have `schema_migrations` + seeded — fix the seed row in production (see + `migrations/README.md`) and re-drill. + +6. Run a smoke query against the `/xrpc/xyz.effem.feed.getSubscriptions` + endpoint with a DID you know exists in the backup. Confirm rows come back. + +7. Tear down the scratch container. + +### Drill success criteria + +- Restore completed in under 30 minutes for a dump up to 10 GB. +- AppView booted without errors. +- At least one read endpoint returned the expected rows. + +Record the drill date, duration, and any issues in a new entry at the bottom +of this file. If the drill fails, restoring becomes the highest-priority +follow-up until it succeeds. + +## 4. Pre-deploy snapshot policy + +Every deploy that includes a new `migrations/*.sql` file must: + +1. Manual Railway snapshot, name it `pre-` (e.g. + `pre-0002`). +2. Link the snapshot name in the PR description. +3. Keep the snapshot until the deploy has been stable for at least 24 hours. + +For rollback: + +1. Identify the last healthy snapshot. +2. In Railway, restore the database from that snapshot to a **new** database + instance. +3. Swap `EFFEM_DATABASE_URL` on the service to point at the restored instance. +4. Redeploy the previous AppView image (the one without the failed migration). +5. Once stable, delete the broken instance. + +Do not restore in-place over the live database unless you have also stopped +every writer (AppView + firehose consumer) — a partial restore under load +leaves the cluster in an unrecoverable state. + +## 5. Drill log + +| Date | Dump size | Restore duration | Notes | +|---|---|---|---| +| _(first drill)_ | | | | diff --git a/appview/database/migrations.go b/appview/database/migrations.go index b2ec5a8..b647a9d 100644 --- a/appview/database/migrations.go +++ b/appview/database/migrations.go @@ -1,97 +1,236 @@ package database import ( + "context" + "crypto/sha256" + "database/sql" + "embed" + "encoding/hex" "fmt" - "log/slog" + "io/fs" + "path" + "sort" + "strings" + "time" "gorm.io/gorm" ) -// dropStaleColumns removes leftover columns created by GORM's broken snake_case -// conversion of acronym field names (DID → d_id, ATURI → a_t_u_r_i, etc.). -// The correct columns (did, at_uri, subject_did) already exist; these are duplicates. -func dropStaleColumns(db *gorm.DB) { - drops := []struct { - table, column string - }{ - {"subscriptions", "d_id"}, - {"comments", "d_id"}, - {"comments", "a_t_u_r_i"}, - {"recommendations", "d_id"}, - {"podcast_lists", "d_id"}, - {"bookmarks", "d_id"}, - {"profiles", "d_id"}, - {"episode_states", "d_id"}, - {"blocks", "d_id"}, - {"blocks", "a_t_u_r_i"}, - {"blocks", "subject_d_i_d"}, - {"reports", "d_id"}, - } - - for _, d := range drops { - // Use raw SQL to check — GORM's HasColumn applies its own naming convention - // which defeats the purpose of checking for the misnamed column. - sql := fmt.Sprintf( - `ALTER TABLE %q DROP COLUMN IF EXISTS %q`, - d.table, d.column, - ) - if err := db.Exec(sql).Error; err != nil { - slog.Warn("drop column failed", "table", d.table, "column", d.column, "err", err) - } - } +//go:embed migrations/*.sql +var embeddedMigrations embed.FS + +// advisoryLockID is a fixed 64-bit integer used with pg_advisory_lock so +// concurrent AppView instances cannot apply migrations at the same time. +const advisoryLockID int64 = 838001447 + +// noTransactionDirective lets a migration opt out of the wrapping transaction. +// Needed for statements that Postgres does not allow inside a tx, such as +// CREATE INDEX CONCURRENTLY. Place this marker on the first non-blank line +// (after any leading -- comments): +// +// -- migrate:no-transaction +const noTransactionDirective = "migrate:no-transaction" + +type migration struct { + Version string + Name string + SQL string + Checksum string + NoTransaction bool } -// widenIntColumns changes integer columns to bigint where PodcastIndex IDs -// can exceed int4 max (~2.1 billion). Safe to run repeatedly. -func widenIntColumns(db *gorm.DB) { - alters := []struct { - table, column string - }{ - {"subscriptions", "feed_id"}, - {"comments", "feed_id"}, - {"comments", "episode_id"}, - {"recommendations", "feed_id"}, - {"recommendations", "episode_id"}, - {"bookmarks", "feed_id"}, - {"bookmarks", "episode_id"}, - {"podcast_stats", "feed_id"}, - {"episode_stats", "episode_id"}, - {"episode_stats", "feed_id"}, - {"episode_states", "feed_id"}, - {"episode_states", "episode_id"}, - } - - for _, a := range alters { - if !db.Migrator().HasTable(a.table) { +// RunMigrations applies every pending migration embedded in migrations/*.sql. +// It is safe to call on every startup; already-applied migrations are skipped. +// The advisory lock ensures that if two instances boot simultaneously only one +// runs migrations while the others wait. +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(), 10*time.Minute) + defer cancel() + + // Pin every statement to a single connection so the advisory lock and + // its matching unlock run in the same session. A pool-level Exec could + // otherwise release the lock on a different connection, leaking it + // until that connection is dropped. + conn, err := sqlDB.Conn(ctx) + if err != nil { + return fmt.Errorf("acquire migration connection: %w", err) + } + defer conn.Close() + + if _, err := conn.ExecContext(ctx, ` + 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() + )`); err != nil { + return fmt.Errorf("create schema_migrations table: %w", err) + } + + if _, err := conn.ExecContext(ctx, "SELECT pg_advisory_lock($1)", advisoryLockID); err != nil { + return fmt.Errorf("acquire migration lock: %w", err) + } + defer func() { + unlockCtx, unlockCancel := context.WithTimeout(context.Background(), 10*time.Second) + defer unlockCancel() + _, _ = conn.ExecContext(unlockCtx, "SELECT pg_advisory_unlock($1)", advisoryLockID) + }() + + migrations, err := loadMigrations(embeddedMigrations) + if err != nil { + return err + } + + for _, m := range migrations { + var existingChecksum string + row := conn.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: recorded %s, file %s (never edit an applied migration)", + m.Name, existingChecksum, m.Checksum, + ) + } continue + case sql.ErrNoRows: + // fall through to apply + default: + return fmt.Errorf("query migration %s: %w", m.Version, err) } - sql := fmt.Sprintf(`ALTER TABLE %q ALTER COLUMN %q TYPE bigint`, a.table, a.column) - if err := db.Exec(sql).Error; err != nil { - slog.Warn("widen column failed (may already be bigint)", "table", a.table, "column", a.column, "err", err) + + if err := applyMigration(ctx, conn, m); err != nil { + return err } } + return nil } -func AutoMigrateAll(db *gorm.DB) error { - dropStaleColumns(db) - widenIntColumns(db) - - if err := db.AutoMigrate( - &FirehoseCursor{}, - &Subscription{}, - &Comment{}, - &Recommendation{}, - &PodcastList{}, - &Bookmark{}, - &Profile{}, - &PICache{}, - &PodcastStats{}, - &EpisodeStats{}, - &EpisodeState{}, - &Block{}, - &Report{}, +// applyMigration runs a single migration. Most migrations run inside a +// transaction with short statement/lock timeouts so a stuck DDL cannot +// block deploys indefinitely. Migrations marked with the no-transaction +// directive run outside a tx with a more generous session-level timeout, +// for statements like CREATE INDEX CONCURRENTLY that cannot be in a tx. +func applyMigration(ctx context.Context, conn *sql.Conn, m migration) error { + if m.NoTransaction { + if _, err := conn.ExecContext(ctx, "SET statement_timeout = '30min'"); err != nil { + return fmt.Errorf("set statement_timeout for %s: %w", m.Name, err) + } + defer func() { + resetCtx, resetCancel := context.WithTimeout(context.Background(), 5*time.Second) + defer resetCancel() + _, _ = conn.ExecContext(resetCtx, "SET statement_timeout = DEFAULT") + }() + if _, err := conn.ExecContext(ctx, m.SQL); err != nil { + return fmt.Errorf("execute migration %s: %w", m.Name, err) + } + if _, err := conn.ExecContext(ctx, + "INSERT INTO schema_migrations(version, name, checksum) VALUES ($1, $2, $3)", + m.Version, m.Name, m.Checksum, + ); err != nil { + return fmt.Errorf("record migration %s: %w", m.Name, err) + } + return nil + } + + tx, err := conn.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin tx for %s: %w", m.Name, err) + } + if _, err := tx.ExecContext(ctx, "SET LOCAL statement_timeout = '5min'"); err != nil { + _ = tx.Rollback() + return fmt.Errorf("set statement_timeout for %s: %w", m.Name, err) + } + if _, err := tx.ExecContext(ctx, "SET LOCAL lock_timeout = '30s'"); err != nil { + _ = tx.Rollback() + return fmt.Errorf("set lock_timeout for %s: %w", m.Name, err) + } + if _, err := tx.ExecContext(ctx, m.SQL); err != nil { + _ = tx.Rollback() + return fmt.Errorf("execute migration %s: %w", m.Name, err) + } + if _, err := tx.ExecContext(ctx, + "INSERT INTO schema_migrations(version, name, checksum) VALUES ($1, $2, $3)", + m.Version, m.Name, m.Checksum, ); err != nil { - return fmt.Errorf("auto migrate: %w", err) + _ = 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 } + +// loadMigrations reads every *.sql file from the migrations/ directory of the +// embedded filesystem and returns them sorted by version prefix. +func loadMigrations(fsys fs.FS) ([]migration, error) { + entries, err := fs.ReadDir(fsys, "migrations") + if err != nil { + return nil, fmt.Errorf("read migrations dir: %w", err) + } + + var list []migration + for _, entry := range entries { + // fs.FS paths use forward slashes — path.Join, not filepath.Join. + if entry.IsDir() || path.Ext(entry.Name()) != ".sql" { + continue + } + name := entry.Name() + version := strings.SplitN(strings.TrimSuffix(name, ".sql"), "_", 2)[0] + if version == "" { + return nil, fmt.Errorf("migration %s has no version prefix", name) + } + + raw, err := fs.ReadFile(fsys, path.Join("migrations", name)) + if err != nil { + return nil, fmt.Errorf("read migration %s: %w", name, err) + } + sum := sha256.Sum256(raw) + list = append(list, migration{ + Version: version, + Name: name, + SQL: strings.TrimSpace(string(raw)), + Checksum: hex.EncodeToString(sum[:]), + NoTransaction: hasNoTransactionDirective(raw), + }) + } + + sort.Slice(list, func(i, j int) bool { return list[i].Version < list[j].Version }) + + // Defensive: reject duplicate version prefixes to avoid silently skipping a file. + seen := make(map[string]string, len(list)) + for _, m := range list { + if prior, ok := seen[m.Version]; ok { + return nil, fmt.Errorf("duplicate migration version %s (%s and %s)", m.Version, prior, m.Name) + } + seen[m.Version] = m.Name + } + + return list, nil +} + +// hasNoTransactionDirective returns true if one of the leading comment lines +// contains the directive marker. The directive must appear before any real +// SQL; everything after the first non-comment, non-blank line is ignored. +func hasNoTransactionDirective(raw []byte) bool { + for _, line := range strings.Split(string(raw), "\n") { + line = strings.TrimSpace(line) + if line == "" { + continue + } + if strings.HasPrefix(line, "--") { + if strings.Contains(line, noTransactionDirective) { + return true + } + continue + } + return false + } + return false +} diff --git a/appview/database/migrations/0001_initial.sql b/appview/database/migrations/0001_initial.sql new file mode 100644 index 0000000..b62ece8 --- /dev/null +++ b/appview/database/migrations/0001_initial.sql @@ -0,0 +1,185 @@ +-- 0001_initial.sql +-- Baseline schema for the Effem AppView. +-- Matches the post-cleanup state produced by the former AutoMigrateAll +-- (widened integer columns, no stale misnamed columns). + +CREATE TABLE "firehose_cursor" ( + "id" BIGSERIAL PRIMARY KEY, + "seq" BIGINT NOT NULL, + "updated_at" TIMESTAMPTZ +); + +CREATE TABLE "subscriptions" ( + "id" BIGSERIAL PRIMARY KEY, + "did" VARCHAR(255) NOT NULL, + "rkey" VARCHAR(512) NOT NULL, + "feed_id" BIGINT NOT NULL, + "feed_url" VARCHAR(2048), + "podcast_guid" VARCHAR(512), + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_subscriptions_did_rkey" ON "subscriptions" ("did", "rkey"); +CREATE INDEX "idx_subscriptions_did" ON "subscriptions" ("did"); +CREATE INDEX "idx_subscriptions_feed_id" ON "subscriptions" ("feed_id"); + +CREATE TABLE "comments" ( + "id" BIGSERIAL PRIMARY KEY, + "did" VARCHAR(255) NOT NULL, + "rkey" VARCHAR(512) NOT NULL, + "cid" VARCHAR(256), + "at_uri" VARCHAR(1024) NOT NULL, + "feed_id" BIGINT NOT NULL, + "episode_id" BIGINT NOT NULL, + "episode_guid" VARCHAR(512), + "podcast_guid" VARCHAR(512), + "text" TEXT NOT NULL, + "timestamp_s" INTEGER, + "reply_root" VARCHAR(1024), + "reply_root_cid" VARCHAR(256), + "reply_parent" VARCHAR(1024), + "reply_parent_cid" VARCHAR(256), + "facets" JSONB, + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_comments_did_rkey" ON "comments" ("did", "rkey"); +CREATE INDEX "idx_comments_did" ON "comments" ("did"); +CREATE INDEX "idx_comments_at_uri" ON "comments" ("at_uri"); +CREATE INDEX "idx_comments_episode" ON "comments" ("feed_id", "episode_id"); +CREATE INDEX "idx_comments_timestamp_s" ON "comments" ("timestamp_s"); +CREATE INDEX "idx_comments_reply_root" ON "comments" ("reply_root"); + +CREATE TABLE "recommendations" ( + "id" BIGSERIAL PRIMARY KEY, + "did" VARCHAR(255) NOT NULL, + "rkey" VARCHAR(512) NOT NULL, + "feed_id" BIGINT NOT NULL, + "episode_id" BIGINT NOT NULL, + "episode_guid" VARCHAR(512), + "podcast_guid" VARCHAR(512), + "text" TEXT, + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_recommendations_did_rkey" ON "recommendations" ("did", "rkey"); +CREATE INDEX "idx_recommendations_did" ON "recommendations" ("did"); +CREATE INDEX "idx_recommendations_episode" ON "recommendations" ("feed_id", "episode_id"); + +CREATE TABLE "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, + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_lists_did_rkey" ON "podcast_lists" ("did", "rkey"); +CREATE INDEX "idx_lists_did" ON "podcast_lists" ("did"); + +CREATE TABLE "bookmarks" ( + "id" BIGSERIAL PRIMARY KEY, + "did" VARCHAR(255) NOT NULL, + "rkey" VARCHAR(512) NOT NULL, + "feed_id" BIGINT NOT NULL, + "episode_id" BIGINT NOT NULL, + "episode_guid" VARCHAR(512), + "podcast_guid" VARCHAR(512), + "timestamp_s" INTEGER, + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_bookmarks_did_rkey" ON "bookmarks" ("did", "rkey"); +CREATE INDEX "idx_bookmarks_did" ON "bookmarks" ("did"); +CREATE INDEX "idx_bookmarks_episode" ON "bookmarks" ("feed_id", "episode_id"); +CREATE INDEX "idx_bookmarks_timestamp_s" ON "bookmarks" ("timestamp_s"); + +CREATE TABLE "profiles" ( + "did" VARCHAR(255) PRIMARY KEY, + "display_name" VARCHAR(640), + "description" TEXT, + "favorite_genres" JSONB, + "indexed_at" TIMESTAMPTZ +); + +CREATE TABLE "pi_cache" ( + "id" BIGSERIAL PRIMARY KEY, + "cache_key" VARCHAR(512), + "response" JSONB NOT NULL, + "expires_at" TIMESTAMPTZ NOT NULL, + "created_at" TIMESTAMPTZ, + "updated_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_pi_cache_cache_key" ON "pi_cache" ("cache_key"); +CREATE INDEX "idx_pi_cache_expires" ON "pi_cache" ("expires_at"); + +CREATE TABLE "podcast_stats" ( + "feed_id" BIGINT PRIMARY KEY, + "subscriber_count" INTEGER DEFAULT 0, + "comment_count" INTEGER DEFAULT 0, + "recommendation_count" INTEGER DEFAULT 0, + "last_updated" TIMESTAMPTZ +); + +CREATE TABLE "episode_stats" ( + "episode_id" BIGINT PRIMARY KEY, + "feed_id" BIGINT NOT NULL, + "comment_count" INTEGER DEFAULT 0, + "recommendation_count" INTEGER DEFAULT 0, + "bookmark_count" INTEGER DEFAULT 0, + "last_updated" TIMESTAMPTZ +); +CREATE INDEX "idx_episode_stats_feed" ON "episode_stats" ("feed_id"); + +CREATE TABLE "episode_states" ( + "id" BIGSERIAL PRIMARY KEY, + "did" VARCHAR(255) NOT NULL, + "rkey" VARCHAR(512) NOT NULL, + "feed_id" BIGINT NOT NULL, + "episode_id" BIGINT NOT NULL, + "episode_guid" VARCHAR(512), + "podcast_guid" VARCHAR(512), + "position_s" INTEGER, + "duration_s" INTEGER, + "played" BOOLEAN DEFAULT false, + "saved" BOOLEAN DEFAULT false, + "hidden" BOOLEAN DEFAULT false, + "updated_at" VARCHAR(64), + "created_at" VARCHAR(64) NOT NULL, + "indexed_at" TIMESTAMPTZ +); +CREATE UNIQUE INDEX "idx_episode_states_did_rkey" ON "episode_states" ("did", "rkey"); +CREATE INDEX "idx_episode_states_did" ON "episode_states" ("did"); +CREATE UNIQUE INDEX "idx_episode_states_episode" ON "episode_states" ("feed_id", "episode_id"); + +CREATE TABLE "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 +); +CREATE UNIQUE INDEX "idx_blocks_did_rkey" ON "blocks" ("did", "rkey"); +CREATE INDEX "idx_blocks_did" ON "blocks" ("did"); +CREATE UNIQUE INDEX "idx_blocks_did_subject" ON "blocks" ("did", "subject_did"); +CREATE INDEX "idx_blocks_at_uri" ON "blocks" ("at_uri"); +CREATE INDEX "idx_blocks_subject" ON "blocks" ("subject_did"); + +CREATE TABLE "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 +); +CREATE UNIQUE INDEX "idx_reports_did_rkey" ON "reports" ("did", "rkey"); +CREATE INDEX "idx_reports_did" ON "reports" ("did"); +CREATE INDEX "idx_reports_status" ON "reports" ("status"); diff --git a/appview/database/migrations/README.md b/appview/database/migrations/README.md new file mode 100644 index 0000000..ca5ba0d --- /dev/null +++ b/appview/database/migrations/README.md @@ -0,0 +1,100 @@ +# AppView Database Migrations + +This directory holds the forward-only SQL migrations applied on AppView startup +by `database.RunMigrations` in `../migrations.go`. + +## How the runner works + +On every boot, `RunMigrations`: + +1. Pins a single `*sql.Conn` from the pool so the advisory lock travels with the same session. +2. Ensures `schema_migrations` exists. +3. Acquires `pg_advisory_lock(838001447)` so only one instance applies migrations at a time. +4. Loads every `*.sql` file from this directory (ordered by the numeric prefix before the first `_`). +5. For each file, checks `schema_migrations.version`: + - Applied with matching SHA-256 checksum → skip. + - Applied with a **different** checksum → refuses to start (somebody edited a committed migration). + - Not applied → runs the SQL and records it. + +Each file runs inside a transaction with `statement_timeout = 5min` and `lock_timeout = 30s`. A bad migration fails fast instead of blocking deploys. + +## Adding a new migration + +1. Create a new file with the next sequential 4-digit prefix, e.g. `0002_add_report_resolution.sql`. The prefix is the version; everything after the first `_` is a human-readable name. +2. Write idempotent, forward-only SQL. Prefer additive changes (`ADD COLUMN`, new tables) over in-place alters on hot tables. +3. Commit the file. The checksum is computed from the file's exact bytes — edits to an already-applied migration will be rejected on the next boot. +4. Never delete a migration file. Undo a change with a new forward migration. + +### Large tables + +A plain `CREATE INDEX` takes an `ACCESS EXCLUSIVE` lock and blocks writes for the duration. For indexes on `comments`, `recommendations`, or anything similarly hot, use `CREATE INDEX CONCURRENTLY`. That cannot run inside a transaction, so mark the file with the opt-out directive as the first non-blank line: + +```sql +-- migrate:no-transaction + +CREATE INDEX CONCURRENTLY idx_comments_reply_root_rkey + ON comments (reply_root, rkey); +``` + +Opted-out migrations run with `statement_timeout = 30min` at the session level. If the statement fails partway through, drop the invalid index and re-run the migration — the `schema_migrations` row is only written after success. + +### Zero-downtime patterns + +- Add a nullable column → backfill in batches in application code or a later migration → mark `NOT NULL` in a follow-up migration once the backfill is complete. +- Renames: add the new column, dual-write from the app, backfill, cut reads over, drop the old column in a final migration. +- `DROP COLUMN` on a big table is cheap logically (metadata-only) but may trigger autovacuum churn — schedule it during a low-traffic window. + +## Bootstrapping an existing database + +`0001_initial.sql` captures the post-cleanup schema that the retired +`AutoMigrateAll` produced. A database that was last migrated by `AutoMigrateAll` +is structurally equivalent, so the simplest bootstrap is to mark `0001` as +already applied and let the runner continue from `0002` onward. + +On the running database, one time only: + +```sql +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() +); + +INSERT INTO schema_migrations (version, name, checksum) +VALUES ('0001', '0001_initial.sql', '') +ON CONFLICT (version) DO NOTHING; +``` + +Compute the checksum on the checked-in file: + +```sh +shasum -a 256 appview/database/migrations/0001_initial.sql | awk '{print $1}' +``` + +After that, deploy the AppView. `RunMigrations` will see `0001` as applied, verify the checksum matches the file you committed, and apply anything numbered `0002` or higher. + +### Fresh databases + +A brand-new empty Postgres runs `0001_initial.sql` automatically — no manual seeding needed. Every later migration follows. + +### If your pre-existing schema disagrees with 0001 + +Run a structural diff before shipping the first follow-up migration: + +```sh +pg_dump --schema-only "$EFFEM_DATABASE_URL" > /tmp/live.sql +diff -u appview/database/migrations/0001_initial.sql /tmp/live.sql +``` + +If the diff is material (missing index, different column type, etc.), either: + +- Replace the committed `0001_initial.sql` with the live `pg_dump` output, recompute the checksum, update the seed row. Do this **before** anyone else pulls — the file is part of the forever contract. +- Or leave `0001` alone and issue a `0002_reconcile.sql` that brings fresh databases in line with the live schema. + +## Policy + +- **Forward-only.** No `down` migrations. +- **Never edit an applied migration.** Checksums enforce this. Add a new file instead. +- **One logical change per file.** Schema and backfill split into separate files lets you monitor and roll each independently. +- **Test migrations on a throwaway copy of production before deploying.** `pg_dump` → restore → run → verify. diff --git a/appview/database/migrations_test.go b/appview/database/migrations_test.go new file mode 100644 index 0000000..dc1648b --- /dev/null +++ b/appview/database/migrations_test.go @@ -0,0 +1,149 @@ +package database + +import ( + "testing" + "testing/fstest" +) + +func TestHasNoTransactionDirective(t *testing.T) { + t.Parallel() + cases := []struct { + name string + sql string + want bool + }{ + {"empty file", "", false}, + {"no comments", "CREATE INDEX foo ON bar(baz);", false}, + { + "directive on first line", + "-- migrate:no-transaction\nCREATE INDEX CONCURRENTLY foo ON bar(baz);", + true, + }, + { + "directive after a blank line", + "\n-- migrate:no-transaction\nCREATE INDEX CONCURRENTLY foo ON bar(baz);", + true, + }, + { + "directive after other comments", + "-- file header\n-- notes: foo\n-- migrate:no-transaction\nCREATE INDEX CONCURRENTLY foo ON bar(baz);", + true, + }, + { + "directive after SQL is ignored", + "CREATE TABLE foo (id int);\n-- migrate:no-transaction\n", + false, + }, + { + "directive inside a string literal is not a comment line", + "CREATE TABLE foo (note text default 'migrate:no-transaction');", + false, + }, + } + + for _, tc := range cases { + tc := tc + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + if got := hasNoTransactionDirective([]byte(tc.sql)); got != tc.want { + t.Fatalf("want %v, got %v", tc.want, got) + } + }) + } +} + +func TestLoadMigrationsSorting(t *testing.T) { + t.Parallel() + fsys := fstest.MapFS{ + "migrations/0002_second.sql": {Data: []byte("CREATE TABLE b ();")}, + "migrations/0001_first.sql": {Data: []byte("CREATE TABLE a ();")}, + "migrations/0010_tenth.sql": {Data: []byte("CREATE TABLE j ();")}, + } + ms, err := loadMigrations(fsys) + if err != nil { + t.Fatalf("loadMigrations: %v", err) + } + want := []string{"0001", "0002", "0010"} + if len(ms) != len(want) { + t.Fatalf("want %d migrations, got %d", len(want), len(ms)) + } + for i, v := range want { + if ms[i].Version != v { + t.Errorf("position %d: want version %s, got %s", i, v, ms[i].Version) + } + } +} + +func TestLoadMigrationsRejectsDuplicateVersion(t *testing.T) { + t.Parallel() + fsys := fstest.MapFS{ + "migrations/0001_first.sql": {Data: []byte("CREATE TABLE a ();")}, + "migrations/0001_second.sql": {Data: []byte("CREATE TABLE b ();")}, + } + if _, err := loadMigrations(fsys); err == nil { + t.Fatal("expected duplicate-version error, got nil") + } +} + +func TestLoadMigrationsSkipsNonSQL(t *testing.T) { + t.Parallel() + fsys := fstest.MapFS{ + "migrations/0001_first.sql": {Data: []byte("CREATE TABLE a ();")}, + "migrations/README.md": {Data: []byte("# ignore me")}, + } + ms, err := loadMigrations(fsys) + if err != nil { + t.Fatalf("loadMigrations: %v", err) + } + if len(ms) != 1 { + t.Fatalf("want 1 migration, got %d", len(ms)) + } +} + +func TestLoadMigrationsChecksumDetectsChange(t *testing.T) { + t.Parallel() + first := fstest.MapFS{ + "migrations/0001_first.sql": {Data: []byte("CREATE TABLE a ();")}, + } + second := fstest.MapFS{ + "migrations/0001_first.sql": {Data: []byte("CREATE TABLE a (id int);")}, + } + a, err := loadMigrations(first) + if err != nil { + t.Fatalf("first load: %v", err) + } + b, err := loadMigrations(second) + if err != nil { + t.Fatalf("second load: %v", err) + } + if a[0].Checksum == b[0].Checksum { + t.Fatalf("checksums should differ after content change") + } +} + +func TestLoadMigrationsRejectsMissingVersion(t *testing.T) { + t.Parallel() + fsys := fstest.MapFS{ + "migrations/.sql": {Data: []byte("CREATE TABLE a ();")}, + } + if _, err := loadMigrations(fsys); err == nil { + t.Fatal("expected missing-version error, got nil") + } +} + +func TestEmbeddedBaselineParses(t *testing.T) { + t.Parallel() + ms, err := loadMigrations(embeddedMigrations) + if err != nil { + t.Fatalf("embedded migrations did not load: %v", err) + } + if len(ms) == 0 { + t.Fatal("expected at least the 0001_initial.sql baseline") + } + if ms[0].Version != "0001" { + t.Fatalf("first migration should be 0001, got %s", ms[0].Version) + } + if ms[0].Checksum == "" { + t.Fatal("baseline checksum is empty") + } +} diff --git a/appview/server.go b/appview/server.go index fd1f226..c0e8866 100644 --- a/appview/server.go +++ b/appview/server.go @@ -57,8 +57,8 @@ func NewServer(cfg Config) (*Server, error) { return nil, fmt.Errorf("connecting to database: %w", err) } - if err := database.AutoMigrateAll(db); err != nil { - return nil, fmt.Errorf("auto migrating database: %w", err) + if err := database.RunMigrations(db); err != nil { + return nil, fmt.Errorf("running database migrations: %w", err) } piClient := podcastindex.NewClient(cfg.PIKey, cfg.PISecret) @@ -76,6 +76,17 @@ func NewServer(cfg Config) (*Server, error) { e := echo.New() e.HideBanner = true + + // Railway's load balancer terminates TLS and forwards via X-Forwarded-For + // from its private network. TrustPrivateNet unwraps that header; the other + // options are defaults but made explicit. Without this, c.RealIP() returns + // the Railway proxy IP and every rate-limit bucket collapses into one. + e.IPExtractor = echo.ExtractIPFromXFFHeader( + echo.TrustLoopback(true), + echo.TrustLinkLocal(true), + echo.TrustPrivateNet(true), + ) + e.Use(middleware.CORSWithConfig(middleware.CORSConfig{ AllowOrigins: cfg.CORSOrigins, AllowMethods: []string{http.MethodGet, http.MethodOptions}, @@ -83,11 +94,17 @@ func NewServer(cfg Config) (*Server, error) { })) e.Use(middleware.Recover()) e.Use(middleware.RequestLoggerWithConfig(middleware.RequestLoggerConfig{ - LogStatus: true, - LogURI: true, - LogMethod: true, + LogStatus: true, + LogURI: true, + LogMethod: true, + LogRemoteIP: true, LogValuesFunc: func(c echo.Context, v middleware.RequestLoggerValues) error { - logger.Info("request", "method", v.Method, "uri", v.URI, "status", v.Status) + logger.Info("request", + "method", v.Method, + "uri", v.URI, + "status", v.Status, + "ip", v.RemoteIP, + ) return nil }, }))