diff --git a/CHANGELOG.md b/CHANGELOG.md index 5fb19c0..c59aeae 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +### Changed + +- ORM: **each migration now runs in a real transaction, and the runner is safe to race.** `Migrator::run`/`rollback` used to fake atomicity with `execute("BEGIN") … execute("COMMIT")` — but both stores are connection *pools*, so under any concurrent traffic those statements land on different connections: measured on the old runner, **7 of 60 failing migrations left their half-applied schema behind** (and on Postgres the parked `BEGIN` poisoned a pooled connection). Each migration's body and its history row now commit in one single-connection transaction via the `Transactional` seam, so `run`, `rollback`, `sync`, and `Model::migrate` now require `Backend + Transactional` — which also makes it a compile error to run migrations from inside an open transaction handle. Concurrent runners serialize on the backend's named advisory lock (`sutegi:migrations`: a dedicated crash-released session on Postgres, the process registry on SQLite), waited for by **polling** rather than a server-side `pg_advisory_lock()` wait — a parked waiter holds a snapshot, which deadlocks against a `CREATE INDEX CONCURRENTLY` holder — and each migration re-checks the history table *inside its own write transaction* (`BEGIN IMMEDIATE` on SQLite), so even two OS processes on one SQLite file — where no shared lock exists — skip instead of double-applying. Verified by a reliability suite that races 8 runners on a pooled file DB, 6 fresh handles on Postgres, and replays the exact old failure under load. + +### Added + +- ORM: **migration preflight guards — nothing runs until the whole plan is validated.** Duplicate versions (two files, or a file shadowing a coded migration), empty or non-portable version/name strings (also closes `write_migration_file` writing outside its directory via a `../`-shaped version), an applied migration that was *renamed* in code (the checksum guard already caught edits; `repair` now re-stamps names too), and an **out-of-order pending migration** — one sorting before a version already applied, the merged-stale-branch hazard whose DDL would run against a schema later migrations already reshaped. Out-of-order is a hard error naming the versions; teams that genuinely interleave opt in with `Migrator::allow_out_of_order()`. The guard anchors only on versions the migrator itself defines, so two apps sharing one database (and one `_sutegi_migrations` table) don't read each other's history as staleness. +- ORM: **`Migrator::rollback` preflights the whole batch before undoing anything.** It used to discover a forward-only or code-deleted migration *mid-batch*, erroring with the newer half already rolled back — the one state with no clean way forward or back. Now that discovery happens up front and the database is untouched. +- ORM: **`Migrator::plan_run(&db)` — the dry run.** The pending migrations in apply order, each with the exact SQL a declarative migration would execute (rendered for the backend's dialect against the live schema, table-rebuild expansions included); closure bodies report `None` rather than pretending. Read-only, so it's safe to wire into a deploy pipeline's review step. +- ORM: **`Migration::no_transaction()`** for DDL that refuses to run inside a transaction — Postgres `CREATE INDEX CONCURRENTLY` being the canonical case. The trade is explicit and documented: a crash between the body and its history row re-runs the body next time, so such migrations must be idempotent. `Migrator::lock_timeout(...)` tunes how long a runner waits for a busy migration lock (default 300 s) instead of hanging a deploy forever, and `MigrationOps` gained `dialect()` so a closure migration can write dialect-specific SQL without guessing. + +### Fixed + +- ORM: **dev-mode `sync` (and `Model::migrate`) is now atomic.** A SQLite column widening is a four-statement table rebuild (create-new / copy / drop / rename); a failure or crash mid-rebuild could previously strand the copy — or worse, sit between the drop and the rename. The whole sync now runs inside one transaction. + ### Added - Web: **`App::listener` — non-HTTP socket loops as first-class app citizens.** A UDP ingest port, a raw TCP protocol, a discovery beacon: `std::net` could always run them on a hand-spawned thread, but that thread was invisible to the app — it outlived graceful drains, saw none of the shared state, and appeared nowhere an agent could discover it. `app.listener(name, doc, run)` registers a closure that runs on its own named thread for the life of the server and receives a `ListenerCtx`: `should_stop()` (the same flag `run_graceful` flips on SIGTERM) plus the typed `state::()` / `db::()` access handlers and tools already have. Shutdown is now whole-app: `run`/`run_until` stop accepting, drain in-flight HTTP requests, then **join listener threads before returning**, so a rolling deploy waits for your loop's last iteration. The contract is cooperative — block with a socket read timeout and poll `should_stop()`, because the join waits for the loop to notice. A panicking listener is caught and reported on stderr instead of dying silently; `/__introspect` gains a `listeners` block (name + doc), keeping the non-HTTP surface agent-discoverable; `App::service()` never spawns listeners, so in-process tests and benches stay socket-free. See `docs/LISTENERS.md`. diff --git a/crates/sutegi-orm/src/backend.rs b/crates/sutegi-orm/src/backend.rs index 4049fff..eca7910 100644 --- a/crates/sutegi-orm/src/backend.rs +++ b/crates/sutegi-orm/src/backend.rs @@ -591,9 +591,10 @@ pub trait Model { /// **add any columns/indexes/foreign keys the model gained** — the fix for /// the old create-if-missing behaviour that silently ignored new fields. /// Additive and non-destructive; it errors (pointing at `migrate gen`) on a - /// change that needs a real migration. Use a [`Migrator`](crate::migrate) - /// for production. - fn migrate(conn: &B) -> Result<(), String> { + /// change that needs a real migration, and runs inside one transaction so + /// a failure never leaves a table half-rebuilt. Use a + /// [`Migrator`](crate::migrate) for production. + fn migrate(conn: &B) -> Result<(), String> { crate::migrate::sync_table(conn, &Self::schema()) } diff --git a/crates/sutegi-orm/src/migrate.rs b/crates/sutegi-orm/src/migrate.rs index 6277a75..68c6fb0 100644 --- a/crates/sutegi-orm/src/migrate.rs +++ b/crates/sutegi-orm/src/migrate.rs @@ -9,8 +9,9 @@ //! A [`Migration`] is a `version` (a sortable id like `20260701_120000`), a //! human `name`, an `up` closure, and an optional `down`. The closures receive //! a [`MigrationOps`] handle — the object-safe subset of [`Backend`] (raw -//! `execute`/`query` plus schema `migrate`) — so a migration can create tables -//! from a [`TableSchema`] *or* run arbitrary DDL/DML. +//! `execute`/`query` plus schema `migrate` and the SQL [`Dialect`]) — so a +//! migration can create tables from a [`TableSchema`] *or* run arbitrary +//! DDL/DML. //! //! ```ignore //! use sutegi::orm::migrate::{Migration, Migrator}; @@ -25,11 +26,37 @@ //! } //! ``` //! -//! Each migration runs inside its own `BEGIN`/`COMMIT` (rolled back on error), -//! so a failing migration leaves neither a half-applied schema nor a history -//! row behind. +//! ## The reliability contract +//! +//! Running migrations is the one moment an app rewrites its own foundations, +//! so [`Migrator::run`] and [`Migrator::rollback`] are built so that **no +//! outcome leaves the database in a state the migrator cannot account for**: +//! +//! - **Each migration is atomic.** Its body and its history row commit in one +//! *real* transaction pinned to a single connection +//! ([`Transactional::run_in_tx`]) — never `BEGIN`/`COMMIT` strings sprayed +//! across a connection pool, where each statement can land on a different +//! connection and "rollback" rolls back nothing. A failing migration leaves +//! neither a half-applied schema nor a history row. +//! - **Concurrent runners serialize.** The run holds the backend's named +//! advisory lock (`sutegi:migrations`) — cluster-wide on Postgres via a +//! dedicated session that auto-releases on crash, process-wide on SQLite. +//! Where the lock cannot reach (two OS processes on one SQLite file), each +//! migration *re-checks the history table inside its own write transaction* +//! before running, so a lost race means a skip, never a double-apply. +//! - **Nothing runs before the plan is validated.** Duplicate or malformed +//! versions, an out-of-order pending migration (older than something already +//! applied — the merged-stale-branch hazard), an applied migration whose +//! file was edited (checksum) or renamed, and a rollback batch containing a +//! forward-only or code-deleted migration are all rejected **up front**, +//! before the database is touched at all. +//! +//! A migration that *must not* run in a transaction (Postgres +//! `CREATE INDEX CONCURRENTLY`) can opt out with +//! [`Migration::no_transaction`]; it trades atomicity for that capability and +//! must be written idempotently. -use crate::backend::Backend; +use crate::backend::{Backend, CapScope, Isolation, LockGuard, Transactional}; use crate::schema_diff::{apply, diff, render, Plan, SchemaOp}; use crate::value::{Dialect, TableSchema, Value}; use sutegi_json::Json; @@ -38,6 +65,14 @@ use sutegi_json::Json; /// are spelled the same on SQLite and Postgres. const HISTORY_TABLE: &str = "_sutegi_migrations"; +/// The advisory-lock name serializing concurrent migration runners. +const MIGRATION_LOCK: &str = "sutegi:migrations"; + +/// How long [`Migrator::run`]/[`rollback`](Migrator::rollback) wait for the +/// migration lock before giving up (override with +/// [`Migrator::lock_timeout`]). +const DEFAULT_LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(300); + /// The object-safe slice of [`Backend`] a migration body is handed. /// /// [`Backend`] itself is not object-safe (it has generic `fetch`/`paginate` @@ -52,6 +87,9 @@ pub trait MigrationOps { fn query(&self, sql: &str, params: &[Value]) -> Result, String>; /// Create a table from a schema if it does not already exist. fn migrate_schema(&self, schema: &TableSchema) -> Result<(), String>; + /// The SQL dialect on the other side — for a closure migration that must + /// write dialect-specific DDL by hand. + fn dialect(&self) -> Dialect; } impl MigrationOps for B { @@ -64,6 +102,29 @@ impl MigrationOps for B { fn migrate_schema(&self, schema: &TableSchema) -> Result<(), String> { Backend::migrate(self, schema) } + fn dialect(&self) -> Dialect { + Backend::dialect(self) + } +} + +/// [`MigrationOps`] over a `&dyn Backend` — the sized adapter that lets a +/// migration body run against the transaction handle `run_in_tx` provides +/// (an unsized `dyn Backend` can't coerce to `dyn MigrationOps` directly). +struct DynOps<'a>(&'a dyn Backend); + +impl MigrationOps for DynOps<'_> { + fn execute(&self, sql: &str, params: &[Value]) -> Result { + self.0.execute(sql, params) + } + fn query(&self, sql: &str, params: &[Value]) -> Result, String> { + self.0.query(sql, params) + } + fn migrate_schema(&self, schema: &TableSchema) -> Result<(), String> { + self.0.migrate(schema) + } + fn dialect(&self) -> Dialect { + self.0.dialect() + } } /// The signature of a migration's `up`/`down` step: it is handed the @@ -87,6 +148,7 @@ pub struct Migration { version: String, name: String, body: Body, + transactional: bool, } impl Migration { @@ -97,6 +159,7 @@ impl Migration { version: version.into(), name: name.into(), body: Body::Closure { up, down: None }, + transactional: true, } } @@ -114,6 +177,7 @@ impl Migration { up, down: Some(down), }, + transactional: true, } } @@ -129,15 +193,32 @@ impl Migration { version: version.into(), name: name.into(), body: Body::Ops(ops), + transactional: true, } } + /// Opt this migration out of the per-migration transaction — for DDL that + /// refuses to run inside one, like Postgres `CREATE INDEX CONCURRENTLY`. + /// + /// The trade is explicit: if the process dies between the body finishing + /// and the history row committing, the next run executes the body + /// **again** — so a non-transactional migration must be idempotent + /// (`IF NOT EXISTS` its DDL, key its backfills). + pub fn no_transaction(mut self) -> Migration { + self.transactional = false; + self + } + pub fn version(&self) -> &str { &self.version } pub fn name(&self) -> &str { &self.name } + /// True unless [`no_transaction`](Migration::no_transaction) opted out. + pub fn runs_in_transaction(&self) -> bool { + self.transactional + } /// True if this migration can be rolled back (declarative migrations always /// can; closure migrations only if they were given a `down`). pub fn reversible_migration(&self) -> bool { @@ -205,17 +286,17 @@ impl Migration { } /// Run the forward step against `conn`. - fn run_up(&self, conn: &B) -> Result<(), String> { + fn run_up(&self, conn: &dyn Backend) -> Result<(), String> { match &self.body { - Body::Closure { up, .. } => up(conn), + Body::Closure { up, .. } => up(&DynOps(conn)), Body::Ops(ops) => exec_ops(conn, ops), } } /// Run the reverse step against `conn` (errors for a forward-only closure). - fn run_down(&self, conn: &B) -> Result<(), String> { + fn run_down(&self, conn: &dyn Backend) -> Result<(), String> { match &self.body { - Body::Closure { down: Some(d), .. } => d(conn), + Body::Closure { down: Some(d), .. } => d(&DynOps(conn)), Body::Closure { down: None, .. } => Err(format!( "cannot roll back {} ({}): migration is forward-only", self.version, self.name @@ -228,11 +309,21 @@ impl Migration { } } +/// True for the version/name strings the migrator accepts: non-empty ASCII +/// letters, digits, `_`, `-`, `.`. One path component by construction, so a +/// version can never traverse out of the migrations directory, and `<` on the +/// string is a sane apply order. +fn valid_ident(s: &str) -> bool { + !s.is_empty() + && s.chars() + .all(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '.')) +} + /// Execute a list of schema ops against `conn`: render each to the backend's /// dialect and run it, threading a shadow schema forward so SQLite's table /// rebuilds see the correct pre-op state. The shadow starts from the live /// database so ops that touch pre-existing tables render correctly. -fn exec_ops(conn: &B, ops: &[SchemaOp]) -> Result<(), String> { +fn exec_ops(conn: &dyn Backend, ops: &[SchemaOp]) -> Result<(), String> { let dialect = conn.dialect(); let mut shadow = conn.introspect()?; for op in ops { @@ -257,16 +348,39 @@ pub struct MigrationStatus { pub orphan: bool, } +/// What [`Migrator::plan_run`] would do for one pending migration, without +/// doing it. +#[derive(Clone, Debug, PartialEq)] +pub struct PlannedMigration { + pub version: String, + pub name: String, + /// The SQL a declarative migration would execute, rendered for the + /// backend's dialect against the current schema. `None` for a closure + /// migration — its statements only exist at run time. (Rendering after a + /// closure is best-effort: the closure's schema effects are unknowable + /// without running it.) + pub statements: Option>, +} + /// An ordered set of migrations plus the run/rollback/status machinery. -#[derive(Default)] pub struct Migrator { migrations: Vec, + allow_out_of_order: bool, + lock_timeout: std::time::Duration, +} + +impl Default for Migrator { + fn default() -> Migrator { + Migrator::new() + } } impl Migrator { pub fn new() -> Migrator { Migrator { migrations: Vec::new(), + allow_out_of_order: false, + lock_timeout: DEFAULT_LOCK_TIMEOUT, } } @@ -278,6 +392,25 @@ impl Migrator { self } + /// Accept a pending migration whose version sorts *before* one already + /// applied. Off by default: an out-of-order pending migration usually + /// means a stale branch was merged, and applying it changes a schema that + /// later migrations already built on — [`run`](Migrator::run) errors and + /// names the versions instead. Opt in when parallel teams genuinely ship + /// interleaved versions. + pub fn allow_out_of_order(mut self) -> Migrator { + self.allow_out_of_order = true; + self + } + + /// How long [`run`](Migrator::run)/[`rollback`](Migrator::rollback) wait + /// for the `sutegi:migrations` advisory lock before erroring (default + /// 300 s). Raise it when a fleet's slowest migration outlives it. + pub fn lock_timeout(mut self, timeout: std::time::Duration) -> Migrator { + self.lock_timeout = timeout; + self + } + /// Load every `*.json` migration file in `dir` and register it. Files are /// parsed via [`Migration::from_json`]; the version/name come from the file /// contents, not the filename. Missing directory is not an error (no files @@ -309,6 +442,72 @@ impl Migrator { v } + /// Reject a malformed migration set before anything touches the database: + /// empty or non-portable version/name strings, and duplicate versions + /// (two files, a file shadowing a coded migration, a copy-paste slip) — + /// running under a duplicate would apply one body but record a version + /// that ambiguously names two. + fn validate(&self) -> Result<(), String> { + let mut seen = std::collections::BTreeMap::new(); + for m in &self.migrations { + if !valid_ident(&m.version) { + return Err(format!( + "invalid migration version {:?} ({}): use ASCII letters, digits, `_`, `-`, `.`", + m.version, m.name + )); + } + if !valid_ident(&m.name) { + return Err(format!( + "invalid migration name {:?} ({}): use ASCII letters, digits, `_`, `-`, `.`", + m.name, m.version + )); + } + if let Some(other) = seen.insert(m.version.as_str(), m.name.as_str()) { + return Err(format!( + "duplicate migration version {}: registered as both {:?} and {:?}", + m.version, other, m.name + )); + } + } + Ok(()) + } + + /// Hold the backend's migration lock for the duration of the returned + /// guard. Backends without advisory locks (`CapScope::None`) proceed + /// without one — the in-transaction history re-check still prevents + /// double-apply there. + /// + /// Waits by **polling `try_lock`**, never by a server-side blocking wait: + /// on Postgres a session parked inside `pg_advisory_lock()` holds a + /// snapshot for the whole wait, and the holder running a + /// [`no_transaction`](Migration::no_transaction) `CREATE INDEX + /// CONCURRENTLY` must wait for every such snapshot — a deadlock between + /// the waiters and the very migration they're waiting on. A polling + /// waiter is snapshot-free between attempts, so the holder always + /// finishes. + fn acquire_lock(&self, conn: &B) -> Result, String> { + if conn.capabilities().advisory_locks == CapScope::None { + return Ok(None); + } + let poll = std::time::Duration::from_millis(200); + let deadline = std::time::Instant::now() + self.lock_timeout; + loop { + if let Some(guard) = conn.try_lock(MIGRATION_LOCK)? { + return Ok(Some(guard)); + } + let now = std::time::Instant::now(); + if now >= deadline { + return Err(format!( + "could not acquire the migration lock within {:?} — another \ + migration runner appears to be active; retry once it finishes, \ + or raise Migrator::lock_timeout", + self.lock_timeout + )); + } + std::thread::sleep(poll.min(deadline - now)); + } + } + /// Create the history table if absent. Idempotent; tolerant of a /// concurrent pod winning the `IF NOT EXISTS` race. fn ensure_history(&self, conn: &B) -> Result<(), String> { @@ -356,23 +555,12 @@ impl Migrator { Ok(out) } - /// Apply every pending migration in version order, each in its own - /// transaction. Returns the versions applied (empty if already up to date). - /// - /// On Postgres a session-level advisory lock serializes concurrent runners - /// (many pods booting at once) so migrations can't race. An already-applied - /// migration whose checksum no longer matches its stored one is a hard error - /// — a file was edited after being applied. - pub fn run(&self, conn: &B) -> Result, String> { - let _guard = AdvisoryLock::acquire(conn)?; - self.ensure_history(conn)?; - let applied = self.applied(conn)?; - let done: std::collections::BTreeSet<&str> = - applied.iter().map(|r| r.version.as_str()).collect(); - let next_batch = applied.iter().map(|r| r.batch).max().unwrap_or(0) + 1; - let now = sutegi_crypto::now_secs(); - - // Tamper check: a migration already in history must still hash the same. + /// The tamper checks run before anything is applied: an already-applied + /// migration must still hash to its stored checksum (a file edited after + /// apply) and must still carry its recorded name (a rename after apply). + /// Both are fixed by restoring the file or, if the change was deliberate, + /// by [`repair`](Migrator::repair). + fn check_integrity(&self, applied: &[AppliedRow]) -> Result<(), String> { for m in self.sorted() { if let Some(row) = applied.iter().find(|r| r.version == m.version) { let current = m.checksum(); @@ -383,51 +571,220 @@ impl Migrator { m.version, m.name )); } + if !row.name.is_empty() && row.name != m.name { + return Err(format!( + "migration {} was renamed after being applied \ + ({:?} in the history, {:?} in code) — restore the name \ + or run `migrate repair`", + m.version, row.name, m.name + )); + } + } + } + Ok(()) + } + + /// The out-of-order guard: a pending migration sorting before the newest + /// applied version usually means a stale branch merged late, and its DDL + /// would run against a schema that later migrations already reshaped. + /// Rejected unless [`allow_out_of_order`](Migrator::allow_out_of_order). + /// + /// Only versions **this migrator defines** anchor the comparison — history + /// rows from another app sharing the database (or from migrations since + /// deleted from code) don't make every new migration read as stale. + fn check_order( + &self, + applied: &[AppliedRow], + done: &std::collections::BTreeSet<&str>, + ) -> Result<(), String> { + if self.allow_out_of_order { + return Ok(()); + } + let defined: std::collections::BTreeSet<&str> = + self.migrations.iter().map(|m| m.version.as_str()).collect(); + let newest_applied = match applied + .iter() + .map(|r| r.version.as_str()) + .filter(|v| defined.contains(v)) + .max() + { + Some(v) => v, + None => return Ok(()), + }; + let stale: Vec<&str> = self + .sorted() + .iter() + .filter(|m| !done.contains(m.version.as_str()) && m.version.as_str() < newest_applied) + .map(|m| m.version.as_str()) + .collect(); + if stale.is_empty() { + Ok(()) + } else { + Err(format!( + "out-of-order migration(s) [{}] sort before the newest applied \ + version ({newest_applied}) — a stale branch was probably merged; \ + renumber them, or opt in with Migrator::allow_out_of_order()", + stale.join(", ") + )) + } + } + + /// True if `version` has a history row — the in-transaction re-check that + /// makes a lost race a skip instead of a double-apply. + fn is_applied(&self, conn: &dyn Backend, version: &str) -> Result { + Ok(!conn + .query( + &format!("SELECT version FROM {HISTORY_TABLE} WHERE version = ?"), + &[Value::Text(version.to_string())], + )? + .is_empty()) + } + + /// Apply one pending migration atomically. Returns `false` if a concurrent + /// runner got there first (seen by the re-check inside the transaction). + fn apply_one( + &self, + conn: &B, + m: &Migration, + batch: i64, + now: i64, + ) -> Result { + let record = |ops: &dyn Backend| { + ops.execute( + &format!( + "INSERT INTO {HISTORY_TABLE} (version, name, batch, checksum, applied_at) \ + VALUES (?, ?, ?, ?, ?)" + ), + &[ + Value::Text(m.version.clone()), + Value::Text(m.name.clone()), + Value::Int(batch), + Value::Text(m.checksum()), + Value::Int(now), + ], + ) + .map(|_| ()) + }; + + if !m.transactional { + if self.is_applied(conn, &m.version)? { + return Ok(false); } + m.run_up(conn)?; + record(conn)?; + return Ok(true); } + let mut applied_now = false; + run_write_tx(conn, &mut |tx| { + if self.is_applied(tx, &m.version)? { + return Ok(()); + } + m.run_up(tx)?; + record(tx)?; + applied_now = true; + Ok(()) + })?; + Ok(applied_now) + } + + /// Apply every pending migration in version order, each atomically (its + /// body and history row in one single-connection transaction). Returns the + /// versions applied (empty if already up to date). + /// + /// Before anything runs, the whole plan is validated — duplicate/malformed + /// versions, checksum and rename tampering, out-of-order pending + /// migrations — and the backend's `sutegi:migrations` advisory lock is + /// held for the duration so concurrent runners (many pods booting at once) + /// serialize. See the module docs for the full reliability contract. + pub fn run(&self, conn: &B) -> Result, String> { + self.validate()?; + let _guard = self.acquire_lock(conn)?; + self.ensure_history(conn)?; + let applied = self.applied(conn)?; + self.check_integrity(&applied)?; + let done: std::collections::BTreeSet<&str> = + applied.iter().map(|r| r.version.as_str()).collect(); + self.check_order(&applied, &done)?; + let next_batch = applied.iter().map(|r| r.batch).max().unwrap_or(0) + 1; + let now = sutegi_crypto::now_secs(); + let mut ran = Vec::new(); for m in self.sorted() { if done.contains(m.version.as_str()) { continue; } - in_transaction(conn, || { - m.run_up(conn)?; - conn.execute( - &format!( - "INSERT INTO {HISTORY_TABLE} (version, name, batch, checksum, applied_at) \ - VALUES (?, ?, ?, ?, ?)" - ), - &[ - Value::Text(m.version.clone()), - Value::Text(m.name.clone()), - Value::Int(next_batch), - Value::Text(m.checksum()), - Value::Int(now), - ], - )?; - Ok(()) - }) - .map_err(|e| format!("migration {} ({}) failed: {e}", m.version, m.name))?; - ran.push(m.version.clone()); + let applied_now = self + .apply_one(conn, m, next_batch, now) + .map_err(|e| format!("migration {} ({}) failed: {e}", m.version, m.name))?; + if applied_now { + ran.push(m.version.clone()); + } } Ok(ran) } - /// Re-stamp the stored checksums to match the current migration files — - /// the escape hatch after a deliberate edit to an applied migration. Only - /// touches rows that are both applied and still defined in code. + /// What [`run`](Migrator::run) would do, without doing it: the pending + /// migrations in apply order, each with the SQL it would execute (rendered + /// against the live schema for declarative migrations; `None` for closure + /// bodies, which only exist at run time). Read-only apart from creating + /// the (empty) history table on a fresh database. + pub fn plan_run(&self, conn: &B) -> Result, String> { + self.validate()?; + self.ensure_history(conn)?; + let applied = self.applied(conn)?; + let done: std::collections::BTreeSet<&str> = + applied.iter().map(|r| r.version.as_str()).collect(); + + let dialect = conn.dialect(); + let mut shadow = conn.introspect()?; + let mut out = Vec::new(); + for m in self.sorted() { + if done.contains(m.version.as_str()) { + continue; + } + let statements = match m.ops_list() { + Some(ops) => { + let mut stmts = Vec::new(); + for op in ops { + stmts.extend(render(op, dialect, &shadow)?); + apply(&mut shadow, op)?; + } + Some(stmts) + } + None => None, + }; + out.push(PlannedMigration { + version: m.version.clone(), + name: m.name.clone(), + statements, + }); + } + Ok(out) + } + + /// Re-stamp the stored checksums and names to match the current migration + /// files — the escape hatch after a deliberate edit or rename of an + /// applied migration. Only touches rows that are both applied and still + /// defined in code. pub fn repair(&self, conn: &B) -> Result, String> { + self.validate()?; self.ensure_history(conn)?; let applied = self.applied(conn)?; let mut fixed = Vec::new(); for m in self.sorted() { if let Some(row) = applied.iter().find(|r| r.version == m.version) { let current = m.checksum(); - if row.checksum != current { + if row.checksum != current || row.name != m.name { conn.execute( - &format!("UPDATE {HISTORY_TABLE} SET checksum = ? WHERE version = ?"), - &[Value::Text(current), Value::Text(m.version.clone())], + &format!( + "UPDATE {HISTORY_TABLE} SET checksum = ?, name = ? WHERE version = ?" + ), + &[ + Value::Text(current), + Value::Text(m.name.clone()), + Value::Text(m.version.clone()), + ], )?; fixed.push(m.version.clone()); } @@ -436,11 +793,20 @@ impl Migrator { Ok(fixed) } - /// Roll back the most recent `batches` batch(es), newest first. Each - /// migration's `down` runs in its own transaction; a forward-only migration - /// aborts the rollback. Returns the versions rolled back. - pub fn rollback(&self, conn: &B, batches: usize) -> Result, String> { - let _guard = AdvisoryLock::acquire(conn)?; + /// Roll back the most recent `batches` batch(es), newest first, each + /// migration atomically (its `down` and its history delete in one + /// transaction). Returns the versions rolled back. + /// + /// The whole batch is **preflighted before anything is undone**: if any + /// victim is forward-only or no longer defined in code, the rollback + /// errors with the database untouched — never a half-rolled-back batch. + pub fn rollback( + &self, + conn: &B, + batches: usize, + ) -> Result, String> { + self.validate()?; + let _guard = self.acquire_lock(conn)?; self.ensure_history(conn)?; let applied = self.applied(conn)?; if applied.is_empty() || batches == 0 { @@ -461,30 +827,69 @@ impl Migrator { .collect(); victims.sort_by(|a, b| b.version.cmp(&a.version)); - let by_version = |v: &str| self.migrations.iter().find(|m| m.version == v); + // Preflight the whole batch before touching anything. + let mut plan: Vec<(&AppliedRow, &Migration)> = Vec::with_capacity(victims.len()); + for row in victims { + let migration = self + .migrations + .iter() + .find(|m| m.version == row.version) + .ok_or_else(|| { + format!( + "cannot roll back {}: no such migration in code — nothing was rolled back", + row.version + ) + })?; + if !migration.reversible_migration() { + return Err(format!( + "cannot roll back {} ({}): migration is forward-only — nothing was rolled back", + row.version, row.name + )); + } + plan.push((row, migration)); + } let mut rolled = Vec::new(); - for row in victims { - let migration = by_version(&row.version).ok_or_else(|| { - format!( - "cannot roll back {}: no such migration in code", - row.version - ) - })?; - in_transaction(conn, || { - migration.run_down(conn)?; - conn.execute( - &format!("DELETE FROM {HISTORY_TABLE} WHERE version = ?"), - &[Value::Text(row.version.clone())], - )?; - Ok(()) - }) - .map_err(|e| format!("rollback of {} ({}) failed: {e}", row.version, row.name))?; + for (row, migration) in plan { + self.rollback_one(conn, migration) + .map_err(|e| format!("rollback of {} ({}) failed: {e}", row.version, row.name))?; rolled.push(row.version.clone()); } Ok(rolled) } + /// Undo one migration atomically, skipping if a concurrent runner already + /// removed its history row. + fn rollback_one( + &self, + conn: &B, + m: &Migration, + ) -> Result<(), String> { + let erase = |ops: &dyn Backend| { + ops.execute( + &format!("DELETE FROM {HISTORY_TABLE} WHERE version = ?"), + &[Value::Text(m.version.clone())], + ) + .map(|_| ()) + }; + + if !m.transactional { + if !self.is_applied(conn, &m.version)? { + return Ok(()); + } + m.run_down(conn)?; + return erase(conn); + } + + run_write_tx(conn, &mut |tx| { + if !self.is_applied(tx, &m.version)? { + return Ok(()); + } + m.run_down(tx)?; + erase(tx) + }) + } + /// The status of every migration — code-defined and orphaned — sorted by /// version. pub fn status(&self, conn: &B) -> Result, String> { @@ -533,6 +938,7 @@ impl Migrator { ("name", Json::str(m.name.clone())), ("reversible", Json::Bool(m.reversible_migration())), ("declarative", Json::Bool(m.ops_list().is_some())), + ("transactional", Json::Bool(m.transactional)), ]) }) .collect(), @@ -566,12 +972,31 @@ impl Migrator { /// Build the shadow schema by **replaying** all migrations against a fresh /// scratch backend (e.g. an in-memory SQLite) and introspecting the result. /// Handles closure migrations, which [`shadow`](Migrator::shadow) can't fold. - pub fn shadow_via(&self, scratch: &B) -> Result, String> { + pub fn shadow_via( + &self, + scratch: &B, + ) -> Result, String> { self.run(scratch)?; Ok(normalize_all(scratch.introspect()?)) } } +/// Run `f` in a real single-connection transaction, taking the write lock up +/// front where the backend can express it (`BEGIN IMMEDIATE` on SQLite via +/// `Isolation::RepeatableRead`) so a cross-process racer serializes at BEGIN +/// instead of deadlocking on a mid-transaction lock upgrade. Backends without +/// isolation levels get a plain transaction. +fn run_write_tx( + conn: &B, + f: &mut dyn FnMut(&dyn Backend) -> Result<(), String>, +) -> Result<(), String> { + if conn.capabilities().isolation_levels { + conn.run_in_tx_with(Isolation::RepeatableRead, f) + } else { + conn.run_in_tx(f) + } +} + /// Fold ops into a schema set (thin wrapper over the diff engine's `apply`). fn apply_all_ops(schemas: &mut Vec, ops: &[SchemaOp]) -> Result<(), String> { for op in ops { @@ -604,7 +1029,7 @@ pub fn generate( /// Like [`generate`], but replays the migration history (including closures) /// against `scratch` to obtain the shadow — use when closure migrations exist. -pub fn generate_via( +pub fn generate_via( migrator: &Migrator, scratch: &B, desired: &[TableSchema], @@ -634,13 +1059,24 @@ fn is_drop(op: &SchemaOp) -> bool { /// columns/indexes/foreign keys, apply safe column widenings. Returns the /// summaries of what it applied. /// +/// The whole sync — introspection, planning, and every op — runs inside one +/// single-connection transaction, so a failure (or a crash mid-way through a +/// SQLite table rebuild) leaves the database exactly as it was. +/// /// It never drops anything (extra columns and tables are left untouched), and it /// refuses — with an error pointing at `migrate gen` — any change that could /// lose data or fail on existing rows (a `NOT NULL` column with no default, a /// lossy type change). This is the honest replacement for the old /// create-if-missing [`Model::migrate`](crate::Model::migrate): a convenience /// for local iteration, not a substitute for reviewed migrations in production. -pub fn sync(conn: &B, desired: &[TableSchema]) -> Result, String> { +pub fn sync( + conn: &B, + desired: &[TableSchema], +) -> Result, String> { + conn.transact(|tx| sync_in_tx(tx, desired)) +} + +fn sync_in_tx(conn: &dyn Backend, desired: &[TableSchema]) -> Result, String> { let dialect = conn.dialect(); let live = conn.introspect()?; @@ -689,7 +1125,10 @@ pub fn sync(conn: &B, desired: &[TableSchema]) -> Result /// Single-table [`sync`] — the engine behind the reimplemented /// [`Model::migrate`](crate::Model::migrate). -pub fn sync_table(conn: &B, schema: &TableSchema) -> Result<(), String> { +pub fn sync_table( + conn: &B, + schema: &TableSchema, +) -> Result<(), String> { sync(conn, std::slice::from_ref(schema)).map(|_| ()) } @@ -760,8 +1199,18 @@ pub fn drift_with_shadow( } /// Write a declarative migration to `/_.json` (creating the -/// directory if needed) and return the path. Errors for a closure migration. +/// directory if needed) and return the path. Errors for a closure migration, +/// and for a version/name that isn't a plain identifier (which could otherwise +/// escape `dir`). pub fn write_migration_file(dir: &str, migration: &Migration) -> Result { + if !valid_ident(migration.version()) || !valid_ident(migration.name()) { + return Err(format!( + "cannot write migration {:?} ({:?}): version and name must be ASCII \ + letters, digits, `_`, `-`, `.`", + migration.version(), + migration.name() + )); + } let json = migration .to_json() .ok_or("cannot write a closure migration to a file")?; @@ -779,41 +1228,6 @@ struct AppliedRow { checksum: String, } -/// A Postgres session advisory lock held for the duration of a run/rollback so -/// concurrent runners serialize. A no-op on SQLite (single-node). Released on -/// drop. -struct AdvisoryLock<'a, B: Backend> { - conn: Option<&'a B>, -} - -/// A fixed key (any constant) identifying the sutegi-migrations lock. -const ADVISORY_LOCK_KEY: i64 = 0x5537_4547_4900; // "SUTEGI" ish - -impl<'a, B: Backend> AdvisoryLock<'a, B> { - fn acquire(conn: &'a B) -> Result, String> { - if conn.dialect() == Dialect::Postgres { - conn.query( - "SELECT pg_advisory_lock(?)", - &[Value::Int(ADVISORY_LOCK_KEY)], - )?; - Ok(AdvisoryLock { conn: Some(conn) }) - } else { - Ok(AdvisoryLock { conn: None }) - } - } -} - -impl Drop for AdvisoryLock<'_, B> { - fn drop(&mut self) { - if let Some(conn) = self.conn { - let _ = conn.query( - "SELECT pg_advisory_unlock(?)", - &[Value::Int(ADVISORY_LOCK_KEY)], - ); - } - } -} - /// Render a status list as JSON (`[{version,name,applied,batch,orphan}]`). pub fn status_json(statuses: &[MigrationStatus]) -> Json { Json::arr( @@ -832,26 +1246,6 @@ pub fn status_json(statuses: &[MigrationStatus]) -> Json { ) } -/// Run `body` between `BEGIN` and `COMMIT`, rolling back on error. Uses the -/// backend's own `execute`, so it works identically on SQLite and Postgres -/// (both give transactional DDL). -fn in_transaction( - conn: &B, - body: impl FnOnce() -> Result<(), String>, -) -> Result<(), String> { - conn.execute("BEGIN", &[])?; - match body() { - Ok(()) => { - conn.execute("COMMIT", &[])?; - Ok(()) - } - Err(e) => { - let _ = conn.execute("ROLLBACK", &[]); - Err(e) - } - } -} - #[cfg(all(test, feature = "sqlite"))] mod tests { use super::*; @@ -983,6 +1377,10 @@ mod tests { Some("0001_create_users") ); assert_eq!(arr[0].get("reversible").and_then(Json::as_bool), Some(true)); + assert_eq!( + arr[0].get("transactional").and_then(Json::as_bool), + Some(true) + ); } // ---- P5: declarative ops migrations + generation ---- @@ -1197,4 +1595,223 @@ mod tests { let report = drift(&db, &migrator, &[todos_v1()]).unwrap(); assert!(!report.db_vs_migrations.is_empty()); } + + // ---- reliability guard rails ---- + + #[test] + fn duplicate_version_is_rejected_before_running() { + let db = Db::memory().unwrap(); + let m = Migrator::new() + .add(Migration::new("0001_x", "first", |_| Ok(()))) + .add(Migration::new("0001_x", "second", |_| Ok(()))); + let err = m.run(&db).unwrap_err(); + assert!(err.contains("duplicate migration version"), "got: {err}"); + // Nothing ran, not even the history table row. + assert!(m.status(&db).is_ok()); + } + + #[test] + fn malformed_version_is_rejected() { + let db = Db::memory().unwrap(); + for bad in ["", "0001/evil", "0001 x", "0001;drop"] { + let m = Migrator::new().add(Migration::new(bad, "n", |_| Ok(()))); + let err = m.run(&db).unwrap_err(); + assert!(err.contains("invalid migration version"), "{bad}: {err}"); + } + } + + #[test] + fn write_migration_file_rejects_path_escapes() { + let plan = crate::schema_diff::diff(&[], &[todos_v1()], Dialect::Sqlite); + let m = Migration::ops("../../0001", "create_todos", plan.ops); + let err = write_migration_file("/tmp/nowhere", &m).unwrap_err(); + assert!(err.contains("version and name"), "got: {err}"); + } + + #[test] + fn out_of_order_pending_is_rejected_unless_opted_in() { + let db = Db::memory().unwrap(); + Migrator::new() + .add(Migration::new("0002_later", "later", |db| { + db.execute("CREATE TABLE later (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + })) + .run(&db) + .unwrap(); + + // A stale branch lands 0001 after 0002 is already applied. + let stale = Migration::new("0001_stale", "stale", |db| { + db.execute("CREATE TABLE stale (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + }); + let m = Migrator::new() + .add(Migration::new("0002_later", "later", |_| Ok(()))) + .add(stale); + let err = m.run(&db).unwrap_err(); + assert!(err.contains("out-of-order"), "got: {err}"); + assert!(err.contains("0001_stale"), "got: {err}"); + // The stale migration did not run. + assert!(db.select(&QueryBuilder::table("stale")).is_err()); + + // Explicit opt-in applies it. + let m = Migrator::new() + .add(Migration::new("0002_later", "later", |_| Ok(()))) + .add(Migration::new("0001_stale", "stale", |db| { + db.execute("CREATE TABLE stale (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + })) + .allow_out_of_order(); + assert_eq!(m.run(&db).unwrap(), vec!["0001_stale"]); + } + + #[test] + fn renamed_applied_migration_trips_the_guard_and_repair_fixes_it() { + let db = Db::memory().unwrap(); + Migrator::new() + .add(Migration::new("0001_x", "old_name", |_| Ok(()))) + .run(&db) + .unwrap(); + + let renamed = Migrator::new().add(Migration::new("0001_x", "new_name", |_| Ok(()))); + let err = renamed.run(&db).unwrap_err(); + assert!(err.contains("renamed after being applied"), "got: {err}"); + + assert_eq!(renamed.repair(&db).unwrap(), vec!["0001_x"]); + assert!(renamed.run(&db).unwrap().is_empty()); + let status = renamed.status(&db).unwrap(); + assert_eq!(status[0].name, "new_name"); + assert!(!status[0].orphan); + } + + #[test] + fn rollback_preflights_the_whole_batch_before_undoing_anything() { + let db = Db::memory().unwrap(); + // One batch: a forward-only migration below a reversible one. + let m = Migrator::new() + .add(Migration::new("0001_forward", "forward", |db| { + db.execute("CREATE TABLE fwd (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + })) + .add(Migration::reversible( + "0002_rev", + "rev", + |db| { + db.execute("CREATE TABLE rev (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + }, + |db| db.execute("DROP TABLE rev", &[]).map(|_| ()), + )); + m.run(&db).unwrap(); + + // 0002 would be undone first — but 0001 is forward-only, so the + // preflight must refuse with NOTHING rolled back (0002 still applied). + let err = m.rollback(&db, 1).unwrap_err(); + assert!(err.contains("forward-only"), "got: {err}"); + assert!(db.select(&QueryBuilder::table("rev")).is_ok()); + let status = m.status(&db).unwrap(); + assert!(status.iter().all(|s| s.applied), "{status:?}"); + } + + #[test] + fn rollback_refuses_a_code_deleted_migration_without_undoing_others() { + let db = Db::memory().unwrap(); + migrator().run(&db).unwrap(); + + // Code lost 0001_create_users; its history row is now an orphan. + let partial = Migrator::new().add(Migration::reversible( + "0002_add_posts", + "add_posts", + |_| Ok(()), + |db| db.execute("DROP TABLE posts", &[]).map(|_| ()), + )); + let err = partial.rollback(&db, 1).unwrap_err(); + assert!(err.contains("no such migration in code"), "got: {err}"); + // 0002 was NOT rolled back on the way to discovering the orphan. + assert!(db.select(&QueryBuilder::table("posts")).is_ok()); + } + + #[test] + fn plan_run_previews_sql_without_executing() { + let db = Db::memory().unwrap(); + let v1 = crate::schema_diff::diff(&[], &[todos_v1()], Dialect::Sqlite); + let m = Migrator::new() + .add(Migration::ops("0001_todos", "create_todos", v1.ops)) + .add(Migration::new("0002_backfill", "backfill", |_| Ok(()))); + + let plan = m.plan_run(&db).unwrap(); + assert_eq!(plan.len(), 2); + assert_eq!(plan[0].version, "0001_todos"); + let stmts = plan[0].statements.as_ref().unwrap(); + assert!( + stmts + .iter() + .any(|s| s.contains("CREATE TABLE") && s.contains("todos")), + "{stmts:?}" + ); + // Closure bodies have no renderable SQL. + assert!(plan[1].statements.is_none()); + + // Nothing executed: the table does not exist, nothing is applied. + assert!(db.select(&QueryBuilder::table("todos")).is_err()); + assert!(m.status(&db).unwrap().iter().all(|s| !s.applied)); + + // After running, the plan is empty. + m.run(&db).unwrap(); + assert!(m.plan_run(&db).unwrap().is_empty()); + } + + #[test] + fn no_transaction_migration_applies_and_reruns_after_partial_failure() { + let db = Db::memory().unwrap(); + let ok = Migrator::new().add( + Migration::new("0001_idem", "idem", |db| { + db.execute( + "CREATE TABLE IF NOT EXISTS idem (id INTEGER PRIMARY KEY)", + &[], + ) + .map(|_| ()) + }) + .no_transaction(), + ); + assert_eq!(ok.run(&db).unwrap(), vec!["0001_idem"]); + assert!(ok.run(&db).unwrap().is_empty()); + + // A failing non-transactional migration keeps its side effects (the + // documented trade) but records nothing, so the fixed body re-runs. + let boom = Migrator::new().add( + Migration::new("0002_boom", "boom", |db| { + db.execute( + "CREATE TABLE IF NOT EXISTS half (id INTEGER PRIMARY KEY)", + &[], + )?; + Err("deliberate".into()) + }) + .no_transaction(), + ); + assert!(boom.run(&db).is_err()); + assert!(db.select(&QueryBuilder::table("half")).is_ok()); + let fixed = Migrator::new().add( + Migration::new("0002_boom", "boom", |db| { + db.execute( + "CREATE TABLE IF NOT EXISTS half (id INTEGER PRIMARY KEY)", + &[], + ) + .map(|_| ()) + }) + .no_transaction(), + ); + assert_eq!(fixed.run(&db).unwrap(), vec!["0002_boom"]); + } + + #[test] + fn migration_ops_exposes_the_dialect() { + let db = Db::memory().unwrap(); + let m = Migrator::new().add(Migration::new("0001_d", "dialect_probe", |ops| { + if ops.dialect() != Dialect::Sqlite { + return Err("expected sqlite".into()); + } + Ok(()) + })); + m.run(&db).unwrap(); + } } diff --git a/crates/sutegi-orm/tests/migrate_reliability.rs b/crates/sutegi-orm/tests/migrate_reliability.rs new file mode 100644 index 0000000..21ae6f6 --- /dev/null +++ b/crates/sutegi-orm/tests/migrate_reliability.rs @@ -0,0 +1,397 @@ +//! Reliability tests for the migration runner on the conditions production +//! actually has: a **pooled, file-backed** database (where a statement-level +//! `BEGIN`/`COMMIT` would spray across connections) and **concurrent runners** +//! (many pods, or two processes on one SQLite file). Every test here asserts +//! the same contract: whatever fails or races, the database is either fully +//! before a migration or fully after it — never in between. + +#![cfg(feature = "sqlite")] + +use sutegi_orm::db::Db; +use sutegi_orm::migrate::{Migration, Migrator}; +use sutegi_orm::{Backend, QueryBuilder}; + +/// A unique throwaway database file per test (pooled connections, WAL) — +/// `Db::memory()` pins its pool to one connection, which would hide every +/// cross-connection bug this suite exists to catch. +struct TempDb { + path: String, +} + +impl TempDb { + fn new(tag: &str) -> TempDb { + let nanos = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_nanos(); + let path = std::env::temp_dir().join(format!( + "sutegi_mig_rel_{tag}_{}_{nanos}.db", + std::process::id() + )); + TempDb { + path: path.to_str().unwrap().to_string(), + } + } + + fn open(&self) -> Db { + Db::open_pool(&self.path, 4).unwrap() + } +} + +impl Drop for TempDb { + fn drop(&mut self) { + for suffix in ["", "-wal", "-shm"] { + let _ = std::fs::remove_file(format!("{}{suffix}", self.path)); + } + } +} + +fn table_exists(db: &Db, table: &str) -> bool { + db.select(&QueryBuilder::table(table)).is_ok() +} + +fn history_count(db: &Db) -> i64 { + db.count(&QueryBuilder::table("_sutegi_migrations")) + .unwrap_or(0) +} + +#[test] +fn failing_migration_is_atomic_on_a_pooled_file_db() { + let tmp = TempDb::new("atomic"); + let db = tmp.open(); + + let m = Migrator::new().add(Migration::new("0001_boom", "boom", |ops| { + ops.execute("CREATE TABLE half_done (id INTEGER PRIMARY KEY)", &[])?; + ops.execute("INSERT INTO half_done (id) VALUES (1)", &[])?; + Err("deliberate failure".into()) + })); + + let err = m.run(&db).unwrap_err(); + assert!(err.contains("deliberate failure"), "got: {err}"); + + // The whole body rolled back on the ONE connection that ran it — no + // table, no rows, no history entry. + assert!(!table_exists(&db, "half_done")); + assert_eq!(history_count(&db), 0); +} + +#[test] +fn failing_rollback_is_atomic_on_a_pooled_file_db() { + let tmp = TempDb::new("rb_atomic"); + let db = tmp.open(); + + let m = Migrator::new().add(Migration::reversible( + "0001_t", + "t", + |ops| { + ops.execute("CREATE TABLE t (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + }, + |ops| { + ops.execute("DROP TABLE t", &[])?; + Err("down blew up after the drop".into()) + }, + )); + m.run(&db).unwrap(); + + let err = m.rollback(&db, 1).unwrap_err(); + assert!(err.contains("down blew up"), "got: {err}"); + + // The failed `down` rolled back wholesale: the table is still there and + // the migration is still recorded as applied. + assert!(table_exists(&db, "t")); + assert_eq!(history_count(&db), 1); + assert!(m.status(&db).unwrap()[0].applied); +} + +#[test] +fn partial_batch_failure_keeps_the_applied_prefix() { + let tmp = TempDb::new("prefix"); + let db = tmp.open(); + + let broken = Migrator::new() + .add(Migration::new("0001_ok", "ok", |ops| { + ops.execute("CREATE TABLE ok_t (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + })) + .add(Migration::new("0002_bad", "bad", |_| Err("nope".into()))); + + assert!(broken.run(&db).unwrap_err().contains("nope")); + // 0001 committed (its own transaction); 0002 left no trace. + assert!(table_exists(&db, "ok_t")); + assert_eq!(history_count(&db), 1); + + // Ship the fix: only the failed migration runs, in a fresh batch. + let fixed = Migrator::new() + .add(Migration::new("0001_ok", "ok", |_| Ok(()))) + .add(Migration::new("0002_bad", "bad", |ops| { + ops.execute("CREATE TABLE bad_t (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + })); + assert_eq!(fixed.run(&db).unwrap(), vec!["0002_bad"]); + assert!(table_exists(&db, "bad_t")); + assert_eq!(history_count(&db), 2); +} + +/// The migrations a racing-runner test applies. Plain `CREATE TABLE` (no +/// `IF NOT EXISTS`) and an unguarded marker INSERT: a double-apply cannot +/// pass silently — it either errors on the DDL or doubles the marker count. +fn racing_migrator() -> Migrator { + Migrator::new() + .add(Migration::new("0000_markers", "markers", |ops| { + ops.execute("CREATE TABLE markers (version TEXT NOT NULL)", &[]) + .map(|_| ()) + })) + .add(Migration::new("0001_a", "a", |ops| { + ops.execute("CREATE TABLE race_a (id INTEGER PRIMARY KEY)", &[])?; + ops.execute("INSERT INTO markers (version) VALUES ('0001_a')", &[]) + .map(|_| ()) + })) + .add(Migration::new("0002_b", "b", |ops| { + ops.execute("CREATE TABLE race_b (id INTEGER PRIMARY KEY)", &[])?; + ops.execute("INSERT INTO markers (version) VALUES ('0002_b')", &[]) + .map(|_| ()) + })) + .add(Migration::new("0003_c", "c", |ops| { + ops.execute("CREATE TABLE race_c (id INTEGER PRIMARY KEY)", &[])?; + ops.execute("INSERT INTO markers (version) VALUES ('0003_c')", &[]) + .map(|_| ()) + })) +} + +fn assert_applied_exactly_once(db: &Db, applied_by_runners: Vec>) { + // Across all runners, every version was applied by exactly one of them. + let mut all: Vec = applied_by_runners.into_iter().flatten().collect(); + all.sort(); + assert_eq!(all, vec!["0000_markers", "0001_a", "0002_b", "0003_c"]); + + // And the database agrees: one marker row per migration, one history row + // per migration, all tables present. + let markers = db.count(&QueryBuilder::table("markers")).unwrap(); + assert_eq!(markers, 3, "a migration body ran more than once"); + assert_eq!(history_count(db), 4); + for t in ["race_a", "race_b", "race_c"] { + assert!(table_exists(db, t)); + } +} + +#[test] +fn racing_runners_on_one_handle_apply_each_migration_exactly_once() { + let tmp = TempDb::new("race_one_handle"); + let db = tmp.open(); + + let results: Vec> = std::thread::scope(|s| { + (0..8) + .map(|_| { + let db = db.clone(); + s.spawn(move || racing_migrator().run(&db).unwrap()) + }) + .collect::>() + .into_iter() + .map(|h| h.join().unwrap()) + .collect() + }); + + assert_applied_exactly_once(&db, results); +} + +#[test] +fn racing_runners_on_separate_handles_apply_each_migration_exactly_once() { + // Separate `Db` handles = separate pools, like two processes sharing the + // file. Each opens its pool at boot (before migrating, as a real pod + // does — opening mid-race can itself hit SQLITE_BUSY); the process-scope + // lock registry plus the in-transaction history re-check must still give + // exactly-once. + let tmp = TempDb::new("race_two_handles"); + let handles: Vec = (0..4) + .map(|_| Db::open_pool(&tmp.path, 2).unwrap()) + .collect(); + + let results: Vec> = std::thread::scope(|s| { + handles + .iter() + .map(|db| s.spawn(move || racing_migrator().run(db).unwrap())) + .collect::>() + .into_iter() + .map(|h| h.join().unwrap()) + .collect() + }); + + assert_applied_exactly_once(&tmp.open(), results); +} + +#[test] +fn racing_rollbacks_undo_each_migration_exactly_once() { + let tmp = TempDb::new("race_rollback"); + let db = tmp.open(); + + let up = Migrator::new() + .add(Migration::reversible( + "0001_x", + "x", + |ops| { + ops.execute("CREATE TABLE x (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + }, + // A second DROP of the same table errors, so a double-rollback + // cannot pass silently. + |ops| ops.execute("DROP TABLE x", &[]).map(|_| ()), + )) + .add(Migration::reversible( + "0002_y", + "y", + |ops| { + ops.execute("CREATE TABLE y (id INTEGER PRIMARY KEY)", &[]) + .map(|_| ()) + }, + |ops| ops.execute("DROP TABLE y", &[]).map(|_| ()), + )); + up.run(&db).unwrap(); + + let results: Vec> = std::thread::scope(|s| { + (0..6) + .map(|_| { + let db = db.clone(); + s.spawn(move || { + Migrator::new() + .add(Migration::reversible( + "0001_x", + "x", + |_| Ok(()), + |ops| ops.execute("DROP TABLE x", &[]).map(|_| ()), + )) + .add(Migration::reversible( + "0002_y", + "y", + |_| Ok(()), + |ops| ops.execute("DROP TABLE y", &[]).map(|_| ()), + )) + .rollback(&db, 1) + .unwrap() + }) + }) + .collect::>() + .into_iter() + .map(|h| h.join().unwrap()) + .collect() + }); + + let mut all: Vec = results.into_iter().flatten().collect(); + all.sort(); + assert_eq!(all, vec!["0001_x", "0002_y"]); + assert!(!table_exists(&db, "x")); + assert!(!table_exists(&db, "y")); + assert_eq!(history_count(&db), 0); +} + +#[test] +fn a_held_migration_lock_times_out_with_guidance() { + let tmp = TempDb::new("lock_timeout"); + let db = tmp.open(); + + // Simulate a stuck runner: hold the migration lock from "elsewhere". + let held = db + .try_lock("sutegi:migrations") + .unwrap() + .expect("free lock"); + + let m = Migrator::new() + .add(Migration::new("0001_x", "x", |_| Ok(()))) + .lock_timeout(std::time::Duration::from_millis(80)); + let err = m.run(&db).unwrap_err(); + assert!(err.contains("migration lock"), "got: {err}"); + assert!(err.contains("lock_timeout"), "got: {err}"); + assert_eq!(history_count(&db), 0); + + // Once the stuck runner releases, the same migrator goes through. + drop(held); + assert_eq!(m.run(&db).unwrap(), vec!["0001_x"]); +} + +#[test] +fn concurrent_dev_syncs_converge_on_a_pooled_file_db() { + use sutegi_orm::{ColType, Column, TableSchema}; + + let tmp = TempDb::new("sync_race"); + let db = tmp.open(); + let schema = || { + vec![TableSchema::new("sync_t") + .column(Column::new("id", ColType::Integer).primary()) + .column(Column::new("title", ColType::Text))] + }; + + std::thread::scope(|s| { + let handles: Vec<_> = (0..6) + .map(|_| { + let db = db.clone(); + s.spawn(move || { + sutegi_orm::migrate::sync( + &db, + &[TableSchema::new("sync_t") + .column(Column::new("id", ColType::Integer).primary()) + .column(Column::new("title", ColType::Text))], + ) + }) + }) + .collect(); + for h in handles { + // Two syncs can race the same CREATE TABLE: losers may error on + // the duplicate, but no outcome may corrupt the schema. + let _ = h.join().unwrap(); + } + }); + + let live = db.introspect().unwrap(); + assert_eq!(live.len(), 1); + assert_eq!(live[0], schema()[0].normalized()); + // And a follow-up sync agrees there is nothing left to do. + assert!(sutegi_orm::migrate::sync(&db, &schema()) + .unwrap() + .is_empty()); +} + +#[test] +fn failing_migration_under_pool_contention_leaks_nothing() { + use std::sync::atomic::{AtomicBool, Ordering}; + + // The scenario that broke the statement-level BEGIN/COMMIT of old: + // concurrent traffic shuffles the pool's connections, so a transaction + // faked with `execute("BEGIN")` spans connections and its ROLLBACK rolls + // back nothing (measured on the old runner: 7/60 failing migrations left + // their half-applied schema behind under this exact load). + let tmp = TempDb::new("contention"); + let db = tmp.open(); + Backend::execute(&db, "CREATE TABLE noise (id INTEGER PRIMARY KEY)", &[]).unwrap(); + + let stop = AtomicBool::new(false); + let mut leaks = 0; + std::thread::scope(|s| { + for _ in 0..3 { + let db = db.clone(); + let stop = &stop; + s.spawn(move || { + while !stop.load(Ordering::Relaxed) { + let _ = db.select(&QueryBuilder::table("noise")); + } + }); + } + for i in 0..30 { + let m = Migrator::new().add(Migration::new(format!("{i:04}_boom"), "boom", |ops| { + ops.execute("CREATE TABLE half_done (id INTEGER PRIMARY KEY)", &[])?; + ops.execute("INSERT INTO half_done (id) VALUES (1)", &[])?; + Err("deliberate".into()) + })); + let _ = m.run(&db); + if table_exists(&db, "half_done") { + leaks += 1; + let _ = Backend::execute(&db, "DROP TABLE half_done", &[]); + } + } + stop.store(true, Ordering::Relaxed); + }); + assert_eq!( + leaks, 0, + "{leaks}/30 failing migrations left a half-applied schema behind" + ); + assert_eq!(history_count(&db), 0); +} diff --git a/crates/sutegi-orm/tests/pg_migrate.rs b/crates/sutegi-orm/tests/pg_migrate.rs new file mode 100644 index 0000000..f28718e --- /dev/null +++ b/crates/sutegi-orm/tests/pg_migrate.rs @@ -0,0 +1,263 @@ +//! Live migration-reliability tests on Postgres — the backend where the old +//! statement-level `BEGIN`/`COMMIT` was most dangerous, because every +//! `execute` checks a connection out of a pool: a faked transaction spanned +//! connections and left one parked in the pool mid-`BEGIN`. Runs only when +//! `SUTEGI_PG_TEST_URL` is set (same server the other `pg_*` suites use). +//! +//! The history table is shared with the other suites, so every assertion here +//! is scoped to this file's `pgmig_`-prefixed versions and tables — no global +//! counts, no dropping `_sutegi_migrations`. + +#![cfg(feature = "postgres")] + +use sutegi_orm::migrate::{Migration, Migrator}; +use sutegi_orm::pg::Pg; +use sutegi_orm::{Backend, Value}; + +fn db() -> Option { + let url = std::env::var("SUTEGI_PG_TEST_URL").ok()?; + Some(Pg::connect(&url, 4).unwrap()) +} + +/// Remove this test's tables and history rows so reruns start clean. +fn scrub(pg: &Pg, prefix: &str, tables: &[&str]) { + for t in tables { + pg.pool() + .batch(&format!("DROP TABLE IF EXISTS {t}")) + .unwrap(); + } + let _ = pg.execute( + "DELETE FROM _sutegi_migrations WHERE version LIKE ?", + &[Value::Text(format!("{prefix}%"))], + ); +} + +fn applied_versions(pg: &Pg, prefix: &str) -> Vec { + let mut v: Vec = pg + .query( + "SELECT version FROM _sutegi_migrations WHERE version LIKE ?", + &[Value::Text(format!("{prefix}%"))], + ) + .unwrap() + .iter() + .filter_map(|r| r.get("version").and_then(sutegi_json::Json::as_str)) + .map(str::to_string) + .collect(); + v.sort(); + v +} + +#[test] +fn failing_migration_is_atomic_on_the_pg_pool_under_traffic() { + use std::sync::atomic::{AtomicBool, Ordering}; + + let Some(pg) = db() else { + eprintln!("skipping: SUTEGI_PG_TEST_URL not set"); + return; + }; + scrub( + &pg, + "pgmig_atomic", + &["pgmig_atomic_half", "pgmig_atomic_noise"], + ); + Backend::execute( + &pg, + "CREATE TABLE pgmig_atomic_noise (id BIGINT PRIMARY KEY)", + &[], + ) + .unwrap(); + + // Concurrent readers shuffle the pool while migrations run and fail — + // the load that made a statement-level BEGIN span connections. + let stop = AtomicBool::new(false); + let problems = std::thread::scope(|s| { + for _ in 0..3 { + let pg = pg.clone(); + let stop = &stop; + s.spawn(move || { + while !stop.load(Ordering::Relaxed) { + let _ = pg.query("SELECT id FROM pgmig_atomic_noise", &[]); + } + }); + } + + let mut problems: Vec = Vec::new(); + for i in 0..10 { + let m = Migrator::new().add(Migration::new( + format!("pgmig_atomic_{i:03}"), + "boom", + |ops| { + ops.execute( + "CREATE TABLE pgmig_atomic_half (id BIGINT PRIMARY KEY)", + &[], + )?; + ops.execute("INSERT INTO pgmig_atomic_half (id) VALUES (1)", &[])?; + Err("deliberate".into()) + }, + )); + match m.run(&pg) { + Err(e) if e.contains("deliberate") => {} + other => problems.push(format!("iteration {i}: unexpected result {other:?}")), + } + if pg.query("SELECT 1 FROM pgmig_atomic_half", &[]).is_ok() { + problems.push(format!( + "iteration {i}: failing migration leaked a half-applied table" + )); + let _ = pg.pool().batch("DROP TABLE pgmig_atomic_half"); + } + } + stop.store(true, Ordering::Relaxed); + problems + }); + + assert!(problems.is_empty(), "{problems:?}"); + // And the pool came out healthy: no connection stuck inside a BEGIN. + assert!(applied_versions(&pg, "pgmig_atomic").is_empty()); + Backend::execute(&pg, "INSERT INTO pgmig_atomic_noise (id) VALUES (1)", &[]).unwrap(); + scrub( + &pg, + "pgmig_atomic", + &["pgmig_atomic_half", "pgmig_atomic_noise"], + ); +} + +fn racing_migrator() -> Migrator { + Migrator::new() + .add(Migration::new("pgmig_race_000", "markers", |ops| { + ops.execute( + "CREATE TABLE pgmig_race_markers (version TEXT NOT NULL)", + &[], + ) + .map(|_| ()) + })) + .add(Migration::new("pgmig_race_001", "a", |ops| { + ops.execute("CREATE TABLE pgmig_race_a (id BIGINT PRIMARY KEY)", &[])?; + ops.execute( + "INSERT INTO pgmig_race_markers (version) VALUES ('pgmig_race_001')", + &[], + ) + .map(|_| ()) + })) + .add(Migration::new("pgmig_race_002", "b", |ops| { + ops.execute("CREATE TABLE pgmig_race_b (id BIGINT PRIMARY KEY)", &[])?; + ops.execute( + "INSERT INTO pgmig_race_markers (version) VALUES ('pgmig_race_002')", + &[], + ) + .map(|_| ()) + })) +} + +#[test] +fn racing_pg_runners_apply_each_migration_exactly_once() { + let Some(pg) = db() else { + eprintln!("skipping: SUTEGI_PG_TEST_URL not set"); + return; + }; + let tables = ["pgmig_race_markers", "pgmig_race_a", "pgmig_race_b"]; + scrub(&pg, "pgmig_race", &tables); + + // Six independent handles = six pods booting at once. The cluster + // advisory lock must serialize them; plain CREATE TABLE (no IF NOT + // EXISTS) and unguarded marker INSERTs make any double-apply loud. + let results: Vec> = std::thread::scope(|s| { + (0..6) + .map(|_| { + s.spawn(|| { + let pod = db().unwrap(); + racing_migrator().run(&pod).unwrap() + }) + }) + .collect::>() + .into_iter() + .map(|h| h.join().unwrap()) + .collect() + }); + + let mut all: Vec = results.into_iter().flatten().collect(); + all.sort(); + assert_eq!( + all, + vec!["pgmig_race_000", "pgmig_race_001", "pgmig_race_002"] + ); + assert_eq!( + applied_versions(&pg, "pgmig_race"), + vec!["pgmig_race_000", "pgmig_race_001", "pgmig_race_002"] + ); + let markers = pg + .query("SELECT COUNT(*) AS c FROM pgmig_race_markers", &[]) + .unwrap()[0] + .get("c") + .and_then(sutegi_json::Json::as_i64); + assert_eq!(markers, Some(2), "a migration body ran more than once"); + + scrub(&pg, "pgmig_race", &tables); +} + +#[test] +fn pg_rollback_of_a_failing_down_is_atomic() { + let Some(pg) = db() else { + eprintln!("skipping: SUTEGI_PG_TEST_URL not set"); + return; + }; + scrub(&pg, "pgmig_rb", &["pgmig_rb_t"]); + + let m = Migrator::new().add(Migration::reversible( + "pgmig_rb_001", + "t", + |ops| { + ops.execute("CREATE TABLE pgmig_rb_t (id BIGINT PRIMARY KEY)", &[]) + .map(|_| ()) + }, + |ops| { + ops.execute("DROP TABLE pgmig_rb_t", &[])?; + Err("down blew up after the drop".into()) + }, + )); + m.run(&pg).unwrap(); + + let err = m.rollback(&pg, 1).unwrap_err(); + assert!(err.contains("down blew up"), "got: {err}"); + // Transactional DDL on PG: the DROP rolled back with the failure. + assert!(pg.query("SELECT 1 FROM pgmig_rb_t", &[]).is_ok()); + assert_eq!(applied_versions(&pg, "pgmig_rb"), vec!["pgmig_rb_001"]); + + scrub(&pg, "pgmig_rb", &["pgmig_rb_t"]); +} + +#[test] +fn pg_no_transaction_migration_runs_concurrent_index_ddl() { + let Some(pg) = db() else { + eprintln!("skipping: SUTEGI_PG_TEST_URL not set"); + return; + }; + scrub(&pg, "pgmig_notx", &["pgmig_notx_t"]); + + // CREATE INDEX CONCURRENTLY refuses to run inside a transaction — the + // whole reason Migration::no_transaction exists. + let m = Migrator::new() + .add(Migration::new("pgmig_notx_001", "table", |ops| { + ops.execute( + "CREATE TABLE pgmig_notx_t (id BIGINT PRIMARY KEY, k TEXT)", + &[], + ) + .map(|_| ()) + })) + .add( + Migration::new("pgmig_notx_002", "concurrent_index", |ops| { + ops.execute( + "CREATE INDEX CONCURRENTLY IF NOT EXISTS pgmig_notx_k ON pgmig_notx_t (k)", + &[], + ) + .map(|_| ()) + }) + .no_transaction(), + ); + assert_eq!( + m.run(&pg).unwrap(), + vec!["pgmig_notx_001", "pgmig_notx_002"] + ); + assert!(m.run(&pg).unwrap().is_empty()); + + scrub(&pg, "pgmig_notx", &["pgmig_notx_t"]); +} diff --git a/crates/sutegi-web/src/respond.rs b/crates/sutegi-web/src/respond.rs index 5fe3b10..e34687c 100644 --- a/crates/sutegi-web/src/respond.rs +++ b/crates/sutegi-web/src/respond.rs @@ -112,6 +112,17 @@ impl IntoResponse for Response { impl IntoResponse for Error { fn into_response(self) -> Response { + // A 5xx message is an internal detail — a SQL error, a curl stderr, a + // path on disk — that reached the client verbatim for as long as this + // rendered `self.message` unconditionally. Log it where the operator + // can see it and answer with the one thing a client can do about it. + if self.status >= 500 { + eprintln!(" ! {}: {}", self.status, self.message); + return json( + self.status, + &Json::obj(vec![("error", Json::str("internal error"))]), + ); + } let mut obj = vec![("error", Json::str(self.message))]; if let Some(fields) = self.fields { obj.push(("errors", fields)); @@ -170,3 +181,35 @@ impl IntoResponse for Result { } } } + +#[cfg(test)] +mod tests { + use super::*; + + fn body_text(resp: &Response) -> String { + match &resp.body { + sutegi_http::Body::Full(bytes) => String::from_utf8_lossy(bytes).into_owned(), + _ => String::new(), + } + } + + #[test] + fn client_errors_render_their_message() { + let resp = Error::unprocessable("email is taken").into_response(); + assert_eq!(resp.status, 422); + assert!(body_text(&resp).contains("email is taken")); + } + + #[test] + fn server_errors_never_leak_internals() { + // `?` on a database call lands here: the message is a SQLite error + // string with the schema in it. The client gets the one answer it can + // act on; the operator gets the detail on stderr. + let resp = Error::from("no such table: answer_seeds".to_string()).into_response(); + assert_eq!(resp.status, 500); + assert_eq!(body_text(&resp), r#"{"error":"internal error"}"#); + + let resp = Error::internal("curl: (6) could not resolve host").into_response(); + assert!(!body_text(&resp).contains("curl")); + } +} diff --git a/crates/sutegi/src/lib.rs b/crates/sutegi/src/lib.rs index a4cbdb1..338df42 100644 --- a/crates/sutegi/src/lib.rs +++ b/crates/sutegi/src/lib.rs @@ -152,17 +152,18 @@ pub mod config; /// ``` /// /// Then `myapp migrate` applies pending migrations, `myapp migrate:rollback [n]` -/// rolls back the last `n` batches (default 1), and `myapp migrate:status` -/// prints the ledger. Same binary serves and migrates — the Rails/Laravel shape. +/// rolls back the last `n` batches (default 1), `myapp migrate:status` prints +/// the ledger, and `myapp migrate:pending` dry-runs the pending SQL. Same +/// binary serves and migrates — the Rails/Laravel shape. #[cfg(feature = "orm")] pub mod migrate { use sutegi_json::Json; pub use sutegi_orm::migrate::{ drift, generate, generate_via, status_json, write_migration_file, DriftReport, Migration, - MigrationOps, MigrationStatus, Migrator, + MigrationOps, MigrationStatus, Migrator, PlannedMigration, }; pub use sutegi_orm::schema_diff::{self, Plan, SchemaOp}; - use sutegi_orm::{Backend, TableSchema}; + use sutegi_orm::{Backend, TableSchema, Transactional}; /// The conventional directory for generated migration files. pub const MIGRATIONS_DIR: &str = "migrations"; @@ -172,10 +173,14 @@ pub mod migrate { /// stop and exit), `false` if there was no migrate subcommand (carry on and /// serve). On a migration error it prints the error and exits the process /// with status 1, so CI and deploy scripts see a real failure. - pub fn dispatch(migrator: &Migrator, conn: &B) -> bool { + pub fn dispatch(migrator: &Migrator, conn: &B) -> bool { let args: Vec = std::env::args().collect(); match args.get(1).map(String::as_str) { Some("migrate") => finish("migrate", migrator.run(conn).map(report_applied)), + Some("migrate:pending") => finish( + "migrate:pending", + migrator.plan_run(conn).map(report_pending), + ), Some("migrate:rollback") => { let batches = args .get(2) @@ -206,7 +211,7 @@ pub mod migrate { /// return Ok(()); /// } /// ``` - pub fn dispatch_full( + pub fn dispatch_full( migrator: &Migrator, conn: &B, models: &[TableSchema], @@ -299,7 +304,10 @@ pub mod migrate { } /// Roll back every batch, then re-run — a clean rebuild for local dev. - fn run_fresh(migrator: &Migrator, conn: &B) -> Result, String> { + fn run_fresh( + migrator: &Migrator, + conn: &B, + ) -> Result, String> { // Roll back until nothing remains (each call undoes one batch). while !migrator.rollback(conn, 1)?.is_empty() {} migrator.run(conn) @@ -370,6 +378,25 @@ pub mod migrate { } } + fn report_pending(plan: Vec) { + if plan.is_empty() { + println!("migrate:pending: up to date — nothing pending"); + return; + } + println!("migrate:pending: {} migration(s) would run:", plan.len()); + for p in &plan { + println!(" {} ({})", p.version, p.name); + match &p.statements { + Some(stmts) => { + for stmt in stmts { + println!(" {}", stmt.replace('\n', "\n ")); + } + } + None => println!(" "), + } + } + } + fn report_rolled(versions: Vec) { if versions.is_empty() { println!("migrate:rollback: nothing to roll back"); diff --git a/docs/MIGRATIONS.md b/docs/MIGRATIONS.md index 4a89b64..30d8375 100644 --- a/docs/MIGRATIONS.md +++ b/docs/MIGRATIONS.md @@ -84,13 +84,14 @@ reviewed, reversible migration. | `migrate` | apply all pending migrations | | `migrate:rollback [n]` | roll back the last `n` batches (default 1) | | `migrate:status` | print the ledger (✓ applied, ? orphan) | +| `migrate:pending` | dry-run: pending migrations + the exact SQL they'd execute | | `migrate:gen ` | diff models ↔ shadow, write a migration file | | `migrate:plan` | show what `:gen` would write, without writing | | `migrate:drift` | three-way report: DB vs migrations vs models | | `migrate:fresh` | roll everything back and re-run (dev only) | `dispatch` (without `_full`) is the smaller runner: just `migrate` / -`:rollback` / `:status`, for apps that hand-write migrations. +`:rollback` / `:status` / `:pending`, for apps that hand-write migrations. ## Column attributes @@ -154,18 +155,61 @@ let app = app.get("/__migrations", "Migration status + drift.", move |_| sutegi::web::json(200, &report)); ``` -## Integrity: checksums +## Reliability: what `run`/`rollback` guarantee + +Migrations are the one moment the app rewrites its own foundations, so the +runner is built so **no outcome leaves the database in a state it can't account +for**: + +- **Each migration is atomic.** Its body and its history row commit in one real + transaction pinned to a single connection (`Transactional::run_in_tx`). A + failing migration leaves neither a half-applied schema nor a history row — + including on a pooled backend, where a statement-level `BEGIN`/`COMMIT` would + spray across connections and roll back nothing. +- **Concurrent runners serialize.** The run holds the backend's + `sutegi:migrations` advisory lock — cluster-wide on Postgres (a dedicated + session, auto-released on crash), process-wide on SQLite. The lock is waited + for by *polling*, never a server-side blocking wait (a session parked in + `pg_advisory_lock()` holds a snapshot that deadlocks against a + `CREATE INDEX CONCURRENTLY` holder). Timeout is 300 s; tune with + `Migrator::lock_timeout(...)`. +- **Lost races skip, never double-apply.** Each migration re-checks the history + table *inside its own write transaction* before running, so even where the + lock can't reach (two OS processes on one SQLite file) a racer skips what the + winner already applied. +- **The plan is validated before anything runs.** Rejected up front, with the + database untouched: duplicate or malformed versions; a pending migration that + sorts *before* one already applied (the merged-stale-branch hazard — opt in + with `Migrator::allow_out_of_order()` if your teams genuinely interleave); an + applied migration whose file was edited (checksum) or renamed. After a + deliberate edit/rename, `Migrator::repair(&db)` re-stamps the history. +- **Rollback preflights the whole batch.** If any victim is forward-only or no + longer defined in code, the rollback errors with nothing undone — never a + half-rolled-back batch. Each `down` then runs atomically with its history + delete. + +### Dry runs + +`Migrator::plan_run(&db)` returns the pending migrations in apply order with +the exact SQL each declarative migration would execute (rendered for the +backend's dialect against the live schema); closure bodies report `None`. +Nothing is executed. + +### Non-transactional migrations + +DDL that refuses to run inside a transaction — Postgres +`CREATE INDEX CONCURRENTLY` — opts out per migration: -The `_sutegi_migrations` history table stores a checksum of each applied ops -migration. If a migration file is edited after being applied, `migrate` fails -with a checksum mismatch rather than silently diverging. After a deliberate edit, -`Migrator::repair(&db)` re-stamps the stored checksums. - -## Cross-pod safety +```rust +Migration::new("20260901_000001", "index_users_email", |ops| { + ops.execute("CREATE INDEX CONCURRENTLY IF NOT EXISTS users_email ON users (email)", &[]) + .map(|_| ()) +}).no_transaction() +``` -On Postgres, `run`/`rollback` take a session advisory lock, so many pods booting -at once serialize their migrations instead of racing. SQLite is single-node, so -this is a no-op there. +The trade is explicit: if the process dies between the body finishing and the +history row committing, the next run executes the body **again** — write it +idempotently (`IF NOT EXISTS` the DDL, key the backfills). ## Honest caveats