Skip to main content

pylon_core/
migrate.rs

1//
2// This source file is part of the Pylon open source project.
3//
4// Copyright (c) 2026 Jaldis B.V.
5//
6// Licensed under the MIT OR Apache-2.0 license (the "License");
7// you may not use this file except in compliance with the License.
8// You may obtain a copy of the License at
9//
10//     https://opensource.org/licenses/MIT
11//     https://www.apache.org/licenses/LICENSE-2.0
12//
13// Unless required by applicable law or agreed to in writing, software
14// distributed under the License is distributed on an "AS IS" BASIS,
15// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
16// See the License for the specific language governing permissions and
17// limitations under the License.
18//
19
20//! Migration execution: advisory locking, tracking tables, resumable
21//! per-step progress, and dev-mode savepoint retry. A direct Rust port of
22//! `pylon.cli.commands.migrations`'s `_apply_one` (plus the tracking-table
23//! helpers `_ensure_tracking_tables`/`_read_tracking`/`_applied_tip`/
24//! `_record_applied` it and `_apply`'s outer loop share). The outer
25//! `apply` loop itself (chain resolution, `--to` targeting, the squash-
26//! backfill special case, `click.echo` progress output) stays in Python —
27//! this module is the part that's actually "migration execution."
28
29use crate::migration::{MigrationFile, parse_steps, verify_integrity};
30use pylon_pgcon::{PgPool, PgTransaction};
31use pylon_value::DecodedValue;
32
33#[derive(Debug, thiserror::Error)]
34pub enum MigrateError {
35    #[error(transparent)]
36    Integrity(#[from] crate::migration::MigrationError),
37    #[error(transparent)]
38    Db(#[from] pylon_pgcon::Error),
39    #[error(
40        "migration {id} is already recorded as applied onto {recorded_onto}, but the file \
41         being applied claims onto {new_onto} — two different migrations share one ID"
42    )]
43    IdCollision {
44        id: String,
45        recorded_onto: String,
46        new_onto: String,
47    },
48}
49
50pub type Result<T> = std::result::Result<T, MigrateError>;
51
52/// Fixed session-level advisory-lock key used by `apply` — must match the
53/// value the Python implementation used forever, since it's what makes
54/// concurrent `apply` runs (across processes, even across old/new
55/// implementations during a rollout) mutually exclusive.
56pub const ADVISORY_LOCK_KEY: i64 = 7_461_999;
57
58const DUPLICATE_OBJECT_CODES: [tokio_postgres::error::SqlState; 5] = [
59    tokio_postgres::error::SqlState::DUPLICATE_TABLE,
60    tokio_postgres::error::SqlState::DUPLICATE_COLUMN,
61    tokio_postgres::error::SqlState::DUPLICATE_SCHEMA,
62    tokio_postgres::error::SqlState::DUPLICATE_OBJECT,
63    tokio_postgres::error::SqlState::DUPLICATE_DATABASE,
64];
65
66fn is_duplicate_object_error(err: &pylon_pgcon::Error) -> bool {
67    err.sqlstate().is_some_and(|code| DUPLICATE_OBJECT_CODES.contains(code))
68}
69
70/// Brings the whole internal `_pylon` schema up to date — tracking tables,
71/// the index/signal outboxes, the cache-invalidate trigger function, and the
72/// stdlib functions.
73///
74/// Runs the *entire* `export_stdlib()` blob, not just the migration tracking
75/// subset. Those blobs carry their own upgrade statements (`ADD COLUMN IF
76/// NOT EXISTS`, `CREATE OR REPLACE FUNCTION`), so which ones a database
77/// receives decides which internal changes ever reach it. Applying only the
78/// tracking subset here meant `_pylon."Migrations".schema_state` arrived on
79/// every database while `_pylon."IndexOutbox".claimed_at` reached only
80/// databases that had been re-initialized — and the index workers on the
81/// rest failed every claim against a column that was never added.
82///
83/// Every statement is idempotent, and `batch_execute` runs them in one
84/// implicit transaction, so this is safe to call on every migration and
85/// cheap enough at that frequency — it is a few dozen statements, not a
86/// schema diff.
87pub async fn ensure_internal_schema(pool: &PgPool) -> Result<()> {
88    pool.batch_execute(&crate::stdlib::export_stdlib()).await?;
89    Ok(())
90}
91
92/// How a database's internal schema relates to the one this build expects.
93///
94/// The classification is deliberately separate from what any caller *does*
95/// about it: `pylon-server` refuses to start on `TooOld` while a client
96/// raises, and both merely note `Behind`, but the rule for which is which
97/// belongs in one place.
98#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum InternalSchemaState {
100    /// No `_pylon."Internal"` row — a database no migration has ever run
101    /// against. Legal, and not a mismatch: callers should behave exactly as
102    /// they did before this check existed rather than treat it as an error.
103    Unmigrated,
104    /// Older than `MIN_SUPPORTED_INTERNAL_VERSION`: this build cannot work
105    /// against it. The only fix is `pylon migration apply`.
106    TooOld { found: i32, required: i32 },
107    /// Behind this build but still within its supported range — an upgrade
108    /// is pending and everything works meanwhile.
109    Behind { found: i32, current: i32 },
110    /// Exactly what this build writes.
111    Current,
112    /// Written by a *newer* build. Reported, never fatal: this is what a
113    /// rollback looks like, and failing closed would turn the recovery
114    /// lever into a second outage.
115    Newer { found: i32, current: i32 },
116}
117
118impl InternalSchemaState {
119    /// Whether this build should refuse to proceed.
120    pub fn is_fatal(&self) -> bool {
121        matches!(self, InternalSchemaState::TooOld { .. })
122    }
123
124    /// A one-line explanation, or `None` when there is nothing to say.
125    pub fn message(&self) -> Option<String> {
126        match self {
127            InternalSchemaState::Unmigrated | InternalSchemaState::Current => None,
128            InternalSchemaState::TooOld { found, required } => Some(format!(
129                "this database's internal schema (version {found}) is older than this \
130                 version of Pylon supports (version {required}); run `pylon migration apply` \
131                 to bring it up to date"
132            )),
133            InternalSchemaState::Behind { found, current } => Some(format!(
134                "this database's internal schema is at version {found}, this version of \
135                 Pylon writes version {current}; `pylon migration apply` will update it"
136            )),
137            InternalSchemaState::Newer { found, current } => Some(format!(
138                "this database's internal schema (version {found}) was written by a newer \
139                 version of Pylon than this one (version {current}); continuing, but this \
140                 build may not understand everything it finds"
141            )),
142        }
143    }
144}
145
146/// Reads `_pylon."Internal".version`, or `None` for a database that has no
147/// such table — one no migration has ever run against.
148pub async fn read_internal_version(pool: &PgPool) -> Result<Option<i32>> {
149    let rows = match pool
150        .query_typed(
151            r#"SELECT (version) AS result FROM _pylon."Internal" WHERE singleton"#,
152            &[],
153            &pool.types(),
154        )
155        .await
156    {
157        Ok(rows) => rows,
158        // `42P01 undefined_table` is the never-migrated case, not a failure
159        // — same graceful degradation `read_schema_snapshot`'s callers rely
160        // on for `_pylon."Schema"`.
161        Err(e) if e.sqlstate() == Some(&tokio_postgres::error::SqlState::UNDEFINED_TABLE) => return Ok(None),
162        Err(e) => return Err(e.into()),
163    };
164    Ok(match rows.into_iter().next() {
165        Some(DecodedValue::I64(v)) => Some(v as i32),
166        _ => None,
167    })
168}
169
170/// Classifies this database against what this build expects — see
171/// `InternalSchemaState`.
172pub async fn check_internal_schema(pool: &PgPool) -> Result<InternalSchemaState> {
173    use crate::stdlib::ddl::{INTERNAL_SCHEMA_VERSION, MIN_SUPPORTED_INTERNAL_VERSION};
174
175    Ok(match read_internal_version(pool).await? {
176        None => InternalSchemaState::Unmigrated,
177        Some(found) if found < MIN_SUPPORTED_INTERNAL_VERSION => InternalSchemaState::TooOld {
178            found,
179            required: MIN_SUPPORTED_INTERNAL_VERSION,
180        },
181        Some(found) if found < INTERNAL_SCHEMA_VERSION => InternalSchemaState::Behind {
182            found,
183            current: INTERNAL_SCHEMA_VERSION,
184        },
185        Some(found) if found > INTERNAL_SCHEMA_VERSION => InternalSchemaState::Newer {
186            found,
187            current: INTERNAL_SCHEMA_VERSION,
188        },
189        Some(_) => InternalSchemaState::Current,
190    })
191}
192
193/// Upserts the process-wide schema snapshot every client fetches at
194/// startup instead of reading `.pylon/schema.json` — a single row (the
195/// `singleton` PK/CHECK forces at most one), written both by `migration
196/// apply` (the formal path) and `watch` (immediate dev-mode sync), since
197/// either is a point where the live database's actual shape just changed.
198/// A bare schema-file edit with neither applied has no effect here, by
199/// design — clients keep seeing the last-applied/synced shape until one of
200/// those two actually run.
201pub async fn write_schema_snapshot(pool: &PgPool, snapshot_json: &str) -> Result<()> {
202    pool.execute_typed(
203        r#"INSERT INTO _pylon."Schema" (singleton, snapshot, updated_at) VALUES (true, $1::jsonb, now())
204           ON CONFLICT (singleton) DO UPDATE SET snapshot = $1::jsonb, updated_at = now()"#,
205        &[DecodedValue::Str(snapshot_json.to_string())],
206    )
207    .await?;
208    Ok(())
209}
210
211/// Reads the current schema snapshot, or `None` if neither `migration
212/// apply` nor `watch` has ever run against this database.
213pub async fn read_schema_snapshot(pool: &PgPool) -> Result<Option<String>> {
214    let rows = pool
215        .query_typed(
216            r#"SELECT (snapshot::text) AS result FROM _pylon."Schema" WHERE singleton"#,
217            &[],
218            &pool.types(),
219        )
220        .await?;
221    Ok(match rows.into_iter().next() {
222        Some(DecodedValue::Str(s)) => Some(s),
223        _ => None,
224    })
225}
226
227/// One row of `_pylon."Migrations"` — covers both `apply`'s own tracking
228/// logic (id/onto/applied, for computing the applied tip) and `migration
229/// create`'s diff-baseline lookup (db_state and schema_state, the JSON
230/// snapshots recorded on the tip row by `apply`). `filename` isn't read —
231/// nothing in this codebase actually consults it once a row exists.
232#[derive(Debug, Clone, PartialEq)]
233pub struct TrackingRow {
234    pub id: String,
235    pub onto: String,
236    /// The `db_state` JSON snapshot (catalog/DDL-visible shape only),
237    /// pre-rendered as text (`db_state::text`) rather than decoded as
238    /// jsonb — callers (`db_state_from_json`) want the raw JSON string to
239    /// re-parse, not an already-decoded value tree.
240    pub db_state: Option<String>,
241    /// The full `SchemaDescriptor` JSON as of this migration (everything
242    /// `db_state` has, plus schema semantics with zero DDL footprint —
243    /// `readonly`, rewrites, `Channel`s, ...). `migration create` diffs
244    /// against *this* (the tip row's own recorded state), not against
245    /// `_pylon."Schema"`, for the same reason `db_state` already does:
246    /// `watch` may have pushed ad hoc changes straight to the live database
247    /// without ever going through `migration create`, and those must still
248    /// show up as a pending change here rather than silently being treated
249    /// as already-baselined. Pre-rendered as text for the same reason as
250    /// `db_state` — re-parse via `SchemaDescriptor.from_json`.
251    pub schema_state: Option<String>,
252    pub applied: bool,
253}
254
255pub async fn read_tracking(pool: &PgPool) -> Result<Vec<TrackingRow>> {
256    let rows = pool
257        .query_typed(
258            r#"SELECT (id, onto, (db_state::text), (schema_state::text), (applied_at IS NOT NULL)) AS result FROM _pylon."Migrations""#,
259            &[],
260            &pool.types(),
261        )
262        .await?;
263    Ok(rows
264        .into_iter()
265        .filter_map(|row| {
266            let DecodedValue::Composite(fields) = row else {
267                return None;
268            };
269            let [
270                DecodedValue::Str(id),
271                DecodedValue::Str(onto),
272                db_state,
273                schema_state,
274                DecodedValue::Bool(applied),
275            ] = <[DecodedValue; 5]>::try_from(fields).ok()?
276            else {
277                return None;
278            };
279            let db_state = match db_state {
280                DecodedValue::Str(s) => Some(s),
281                _ => None,
282            };
283            let schema_state = match schema_state {
284                DecodedValue::Str(s) => Some(s),
285                _ => None,
286            };
287            Some(TrackingRow {
288                id,
289                onto,
290                db_state,
291                schema_state,
292                applied,
293            })
294        })
295        .collect())
296}
297
298/// Computes the tip ID from applied tracking rows (the one with no
299/// descendant) — a pure function, no I/O, mirroring `_applied_tip`.
300///
301/// A healthy chain has exactly one such row, but a tracking table can end
302/// up with several orphaned single-node "tips" (e.g. leftover rows from
303/// migration files that were since deleted/regenerated without cleaning up
304/// the tracking table) — when that happens, this deterministically returns
305/// the lexicographically smallest ID among them, rather than an arbitrary
306/// one. The original Python implementation this ports picked from a
307/// `set()` difference (`next(iter(tips))`), whose iteration order is
308/// randomized per-process by Python's string hash randomization — a real,
309/// pre-existing bug (not introduced by this port) that could make
310/// `status`/`apply` disagree on which tip is current from one invocation
311/// to the next, since each CLI command is a separate process.
312pub fn applied_tip(tracking: &[TrackingRow]) -> Option<String> {
313    let applied: Vec<&TrackingRow> = tracking.iter().filter(|r| r.applied).collect();
314    if applied.is_empty() {
315        return None;
316    }
317    let onto_targets: std::collections::HashSet<&str> = applied.iter().map(|r| r.onto.as_str()).collect();
318    applied
319        .iter()
320        .filter(|r| !onto_targets.contains(r.id.as_str()))
321        .map(|r| r.id.as_str())
322        .min()
323        .map(|s| s.to_string())
324}
325
326/// Blocks until the advisory lock is acquired (`pg_advisory_lock`), on a
327/// connection the caller then holds for the lock's entire lifetime — an
328/// advisory lock is released by an explicit unlock (or the session
329/// ending), not by a transaction boundary or by returning to the pool, so
330/// acquire and release must run on the *same* connection (see
331/// `PgPool::connection`). Pass the returned handle to `advisory_unlock`
332/// once `apply` is done; dropping it without unlocking first would leave
333/// the lock held until that specific connection eventually closes.
334pub async fn advisory_lock(pool: &PgPool) -> Result<pylon_pgcon::PgConnection> {
335    let conn = pool.connection().await?;
336    conn.batch_execute(&format!("SELECT pg_advisory_lock({ADVISORY_LOCK_KEY})"))
337        .await?;
338    Ok(conn)
339}
340
341/// Attempts to acquire the advisory lock without blocking
342/// (`pg_try_advisory_lock`); `None` means another `apply` holds it —
343/// same held-connection contract as `advisory_lock`.
344pub async fn try_advisory_lock(pool: &PgPool) -> Result<Option<pylon_pgcon::PgConnection>> {
345    let conn = pool.connection().await?;
346    let rows = conn
347        .query_typed(
348            &format!("SELECT (pg_try_advisory_lock({ADVISORY_LOCK_KEY})) AS result"),
349            &[],
350            &pool.types(),
351        )
352        .await?;
353    Ok(if matches!(rows.first(), Some(DecodedValue::Bool(true))) {
354        Some(conn)
355    } else {
356        None
357    })
358}
359
360pub async fn advisory_unlock(conn: pylon_pgcon::PgConnection) -> Result<()> {
361    conn.batch_execute(&format!("SELECT pg_advisory_unlock({ADVISORY_LOCK_KEY})"))
362        .await?;
363    Ok(())
364}
365
366/// Records a migration as applied.
367///
368/// The conflict path updates `onto` and `filename` too, not just
369/// `applied_at`: re-applying the same migration rewrites them with identical
370/// values (free), while leaving them stale would let the applied-tip walk
371/// read an `onto` that doesn't match the migration the ID refers to.
372/// `check_no_id_collision` runs first and rejects the case where they would
373/// genuinely differ.
374const RECORD_APPLIED_SQL: &str = r#"
375    INSERT INTO _pylon."Migrations" (id, onto, filename, applied_at)
376    VALUES ($1, $2, $3, now())
377    ON CONFLICT (id) DO UPDATE
378        SET applied_at = now(),
379            onto = EXCLUDED.onto,
380            filename = EXCLUDED.filename
381"#;
382
383/// Rejects recording `id` when a *different* migration is already tracked
384/// under it — i.e. one whose parent isn't `onto`.
385///
386/// Under the `m2` ID format this can't arise from two same-bodied migrations
387/// at different chain positions, since `onto` is hashed in. It remains
388/// reachable for `m1`-era IDs (body-only hash), where exactly that collision
389/// silently overwrote the first migration's tracking row and left the chain
390/// walk reading a parent that no longer matched.
391async fn check_no_id_collision(pool: &PgPool, id: &str, onto: &str) -> Result<()> {
392    let rows = pool
393        .query_typed(
394            r#"SELECT (onto) AS result FROM _pylon."Migrations" WHERE id = $1"#,
395            &[DecodedValue::Str(id.to_string())],
396            &pool.types(),
397        )
398        .await?;
399    if let Some(DecodedValue::Str(existing_onto)) = rows.into_iter().next()
400        && existing_onto != onto
401    {
402        return Err(MigrateError::IdCollision {
403            id: id.to_string(),
404            recorded_onto: existing_onto,
405            new_onto: onto.to_string(),
406        });
407    }
408    Ok(())
409}
410
411fn record_applied_params(id: &str, onto: &str, filename: &str) -> Vec<DecodedValue> {
412    vec![
413        DecodedValue::Str(id.to_string()),
414        DecodedValue::Str(onto.to_string()),
415        DecodedValue::Str(filename.to_string()),
416    ]
417}
418
419/// Records a migration as applied without running its DDL — used both by
420/// `apply_one`'s last step (inside its transaction, via `record_applied_in_tx`)
421/// and by `apply`'s squash-backfill case (a migration whose squashed
422/// constituent IDs are already applied under the old chain: no DDL to run,
423/// just mark it applied so future `apply` runs see it as done).
424pub async fn record_applied(pool: &PgPool, id: &str, onto: &str, filename: &str) -> Result<()> {
425    check_no_id_collision(pool, id, onto).await?;
426    pool.execute_typed(RECORD_APPLIED_SQL, &record_applied_params(id, onto, filename))
427        .await?;
428    Ok(())
429}
430
431async fn record_applied_in_tx(tx: &PgTransaction, id: &str, onto: &str, filename: &str) -> Result<()> {
432    tx.execute_typed(RECORD_APPLIED_SQL, &record_applied_params(id, onto, filename))
433        .await?;
434    Ok(())
435}
436
437async fn read_progress(pool: &PgPool, id: &str) -> Result<Option<i64>> {
438    let rows = pool
439        .query_typed(
440            r#"SELECT (step_index) AS result FROM _pylon."Progress" WHERE id = $1"#,
441            &[DecodedValue::Str(id.to_string())],
442            &pool.types(),
443        )
444        .await?;
445    Ok(match rows.into_iter().next() {
446        Some(DecodedValue::I64(n)) => Some(n),
447        _ => None,
448    })
449}
450
451async fn record_progress(pool: &PgPool, id: &str, step_index: i64) -> Result<()> {
452    pool.execute_typed(
453        r#"INSERT INTO _pylon."Progress" (id, step_index) VALUES ($1, $2)
454           ON CONFLICT (id) DO UPDATE SET step_index = $2, updated_at = now()"#,
455        &[DecodedValue::Str(id.to_string()), DecodedValue::I64(step_index)],
456    )
457    .await?;
458    Ok(())
459}
460
461async fn delete_progress(pool: &PgPool, id: &str) -> Result<()> {
462    pool.execute_typed(
463        r#"DELETE FROM _pylon."Progress" WHERE id = $1"#,
464        &[DecodedValue::Str(id.to_string())],
465    )
466    .await?;
467    Ok(())
468}
469
470/// Same upsert as `record_progress`, but on an open transaction so a step's
471/// DDL and the progress row that claims it commit together — see
472/// `apply_one`'s own note on why that atomicity matters.
473async fn record_progress_in_tx(tx: &PgTransaction, id: &str, step_index: i64) -> Result<()> {
474    tx.execute_typed(
475        r#"INSERT INTO _pylon."Progress" (id, step_index) VALUES ($1, $2)
476           ON CONFLICT (id) DO UPDATE SET step_index = $2, updated_at = now()"#,
477        &[DecodedValue::Str(id.to_string()), DecodedValue::I64(step_index)],
478    )
479    .await?;
480    Ok(())
481}
482
483async fn delete_progress_in_tx(tx: &PgTransaction, id: &str) -> Result<()> {
484    tx.execute_typed(
485        r#"DELETE FROM _pylon."Progress" WHERE id = $1"#,
486        &[DecodedValue::Str(id.to_string())],
487    )
488    .await?;
489    Ok(())
490}
491
492/// Before retrying a `CONCURRENTLY` step, drops any invalid index it left
493/// behind from a prior failed attempt (a `CREATE INDEX CONCURRENTLY` that
494/// errors partway leaves an unusable index rather than rolling back, since
495/// it can't run inside a transaction).
496async fn drop_invalid_concurrent_index(pool: &PgPool, sql: &str) -> Result<()> {
497    let Some(index_name) = concurrent_index_name(sql) else {
498        return Ok(());
499    };
500    let rows = pool
501        .query_typed(
502            "SELECT (1) AS result FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid \
503             WHERE c.relname = $1 AND NOT i.indisvalid",
504            &[DecodedValue::Str(index_name.clone())],
505            &pool.types(),
506        )
507        .await?;
508    if !rows.is_empty() {
509        pool.batch_execute(&format!("DROP INDEX CONCURRENTLY IF EXISTS \"{index_name}\""))
510            .await?;
511    }
512    Ok(())
513}
514
515/// Extracts the index name from a `CREATE INDEX CONCURRENTLY [IF NOT
516/// EXISTS] name ...` statement, case-insensitively — the exact shape
517/// pylon-core always emits for non-transactional migration steps. `None`
518/// if `sql` doesn't start with that form.
519fn concurrent_index_name(sql: &str) -> Option<String> {
520    let mut tokens = sql.split_whitespace();
521    let matches_kw = |t: Option<&str>, expected: &str| t.is_some_and(|t| t.eq_ignore_ascii_case(expected));
522    if !matches_kw(tokens.next(), "CREATE") {
523        return None;
524    }
525    if !matches_kw(tokens.next(), "INDEX") {
526        return None;
527    }
528    if !matches_kw(tokens.next(), "CONCURRENTLY") {
529        return None;
530    }
531    let mut next = tokens.next()?;
532    if next.eq_ignore_ascii_case("IF") {
533        if !matches_kw(tokens.next(), "NOT") {
534            return None;
535        }
536        if !matches_kw(tokens.next(), "EXISTS") {
537            return None;
538        }
539        next = tokens.next()?;
540    }
541    let after_quote = next.strip_prefix('"').unwrap_or(next);
542    let name: String = after_quote
543        .chars()
544        .take_while(|c| c.is_ascii_alphanumeric() || *c == '_')
545        .collect();
546    if name.is_empty() { None } else { Some(name) }
547}
548
549const DEV_SAVEPOINT: &str = "pylon_dev";
550
551/// Applies one migration's steps in order, resuming from recorded progress
552/// if a prior run failed mid-migration. Verifies `m`'s integrity first
553/// (re-hashes the body against its header ID).
554///
555/// `_pylon."Progress".step_index` is the index of the last step that
556/// **completed**, so a resume starts at `step_index + 1`. For a transactional
557/// step the progress row is written inside that step's own transaction, so
558/// the two commit together — a crash mid-step rolls both back and the step is
559/// retried rather than skipped. A non-transactional step (`CREATE INDEX
560/// CONCURRENTLY`) has no transaction to join, so it records progress
561/// immediately afterwards; a crash in that window just re-runs the step, and
562/// `drop_invalid_concurrent_index` clears the half-built index first.
563///
564/// `dev_mode` rebases onto a database `watch` has already partly changed: it
565/// runs each *statement* in its own savepoint and skips the individual ones
566/// that fail with "already exists", leaving the rest of the step to apply
567/// normally.
568pub async fn apply_one(pool: &PgPool, m: &MigrationFile, dev_mode: bool) -> Result<()> {
569    verify_integrity(m)?;
570    // Checked up front rather than at the `record_applied` at the end, so a
571    // colliding ID is rejected before any of this migration's DDL runs.
572    check_no_id_collision(pool, &m.id, &m.onto).await?;
573
574    let steps = parse_steps(&m.body);
575    let resume_from = read_progress(pool, &m.id).await?.map(|i| i + 1).unwrap_or(0) as usize;
576    let multi_step = steps.len() > 1;
577
578    for (step_idx, (transactional, sql)) in steps.iter().enumerate() {
579        if step_idx < resume_from {
580            continue;
581        }
582        let sql = sql.trim();
583        if sql.is_empty() {
584            continue;
585        }
586
587        let is_last = step_idx == steps.len() - 1;
588
589        if *transactional {
590            let tx = pool.begin_default().await?;
591
592            let step_result = if dev_mode {
593                apply_statements_rebasing(&tx, sql).await
594            } else {
595                tx.batch_execute(sql).await
596            };
597
598            if let Err(e) = step_result {
599                let _ = tx.rollback().await;
600                return Err(e.into());
601            }
602
603            if is_last {
604                record_applied_in_tx(&tx, &m.id, &m.onto, &m.filename).await?;
605                if multi_step {
606                    delete_progress_in_tx(&tx, &m.id).await?;
607                }
608            } else if multi_step {
609                record_progress_in_tx(&tx, &m.id, step_idx as i64).await?;
610            }
611            tx.commit().await?;
612        } else {
613            // One statement per round trip: several sent together are one
614            // simple-query batch, which Postgres runs in an implicit
615            // transaction -- the very thing `CREATE INDEX CONCURRENTLY`
616            // refuses to be in.
617            for statement in crate::migration::split_statements(sql) {
618                drop_invalid_concurrent_index(pool, &statement).await?;
619                pool.batch_execute(&statement).await?;
620            }
621            if is_last {
622                record_applied(pool, &m.id, &m.onto, &m.filename).await?;
623                if multi_step {
624                    delete_progress(pool, &m.id).await?;
625                }
626            } else if multi_step {
627                record_progress(pool, &m.id, step_idx as i64).await?;
628            }
629        }
630    }
631
632    // This DDL may have created a type whose OID the pool has never seen,
633    // or invalidated a cached plan. Both outlive the migration: `apply`
634    // applies every pending migration over one pool, `watch` keeps one for
635    // the life of the process.
636    pool.refresh_types().await?;
637
638    Ok(())
639}
640
641/// Runs one step's statements inside `tx`, each wrapped in its own savepoint,
642/// skipping any that fail because the object already exists.
643///
644/// The savepoint has to be per statement rather than per step: rolling the
645/// whole step back on the first duplicate would undo the statements before it
646/// *and* skip the ones after it, while still reporting the step as applied.
647async fn apply_statements_rebasing(tx: &PgTransaction, sql: &str) -> std::result::Result<(), pylon_pgcon::Error> {
648    for stmt in crate::migration::split_statements(sql) {
649        tx.savepoint(DEV_SAVEPOINT).await?;
650        match tx.batch_execute(&stmt).await {
651            Ok(()) => tx.release_savepoint(DEV_SAVEPOINT).await?,
652            Err(e) if is_duplicate_object_error(&e) => tx.rollback_to_savepoint(DEV_SAVEPOINT).await?,
653            Err(e) => return Err(e),
654        }
655    }
656    Ok(())
657}
658
659#[cfg(test)]
660mod concurrent_index_name_tests {
661    use super::concurrent_index_name;
662
663    #[test]
664    fn extracts_a_bare_index_name() {
665        assert_eq!(
666            concurrent_index_name("CREATE INDEX CONCURRENTLY idx_person_name ON \"public\".\"Person\" (name);"),
667            Some("idx_person_name".to_string())
668        );
669    }
670
671    #[test]
672    fn extracts_a_quoted_index_name() {
673        assert_eq!(
674            concurrent_index_name("CREATE INDEX CONCURRENTLY \"idx_person_name\" ON \"public\".\"Person\" (name);"),
675            Some("idx_person_name".to_string())
676        );
677    }
678
679    #[test]
680    fn handles_if_not_exists() {
681        assert_eq!(
682            concurrent_index_name("CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_x ON t (c);"),
683            Some("idx_x".to_string())
684        );
685    }
686
687    #[test]
688    fn is_case_insensitive() {
689        assert_eq!(
690            concurrent_index_name("create index concurrently idx_x on t (c);"),
691            Some("idx_x".to_string())
692        );
693    }
694
695    #[test]
696    fn returns_none_for_unrelated_sql() {
697        assert_eq!(concurrent_index_name("CREATE TABLE foo ();"), None);
698        assert_eq!(concurrent_index_name("CREATE INDEX idx_x ON t (c);"), None); // not CONCURRENTLY
699    }
700}
701
702#[cfg(test)]
703mod tests {
704    use super::*;
705    use crate::migration::render_file;
706
707    fn test_dsn() -> String {
708        std::env::var("PYLON_PGCON_TEST_DSN").expect("PYLON_PGCON_TEST_DSN must be set to run live-Postgres tests")
709    }
710
711    async fn test_pool() -> PgPool {
712        let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
713        pool.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
714        ensure_internal_schema(&pool).await.unwrap();
715        pool
716    }
717
718    fn make_migration(onto: &str, body: &str) -> MigrationFile {
719        let content = render_file(onto, body, &[]);
720        crate::migration::parse(&content, "test").unwrap()
721    }
722
723    /// A body with the leading blank-line separator `render_file`/`parse`
724    /// expect, so the migration's computed ID matches what gets hashed.
725    fn body(sql: &str) -> String {
726        format!("\n{sql}\n")
727    }
728
729    fn unique_table_name(prefix: &str) -> String {
730        use std::time::{SystemTime, UNIX_EPOCH};
731        let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
732        format!("{prefix}_{nanos}")
733    }
734
735    /// Every test in this module runs against the same shared,
736    /// non-isolated `_pylon."Migrations"` table (unlike the rest of the
737    /// live-execution suite, which gets a fresh schema per test via
738    /// `unique_module()` — there's no equivalent scoping for this
739    /// process-wide tracking table). Any test that records a tracking row
740    /// must delete it again here, or it permanently pollutes whatever
741    /// database `PYLON_PGCON_TEST_DSN` points at (this defaults to the
742    /// same DSN pylon-demo uses, and a stray row here can shadow a real
743    /// project's actual migration tip).
744    async fn cleanup_migration_row(pool: &PgPool, id: &str) {
745        pool.execute_typed(
746            r#"DELETE FROM _pylon."Migrations" WHERE id = $1"#,
747            &[DecodedValue::Str(id.to_string())],
748        )
749        .await
750        .unwrap();
751    }
752
753    /// A migration that creates an enum must leave the pool that applied it
754    /// able to decode one.
755    #[tokio::test]
756    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
757    async fn apply_one_leaves_the_pool_able_to_decode_a_type_it_just_created() {
758        let pool = test_pool().await;
759        let enum_type = unique_table_name("migrate_apply_enum");
760        pool.batch_execute(&format!("DROP TYPE IF EXISTS {enum_type}"))
761            .await
762            .unwrap();
763        let m = make_migration(
764            "initial",
765            &body(&format!("CREATE TYPE {enum_type} AS ENUM ('ok', 'nope');")),
766        );
767
768        apply_one(&pool, &m, false).await.unwrap();
769
770        // Asserted on the registry, before anything decodes an enum: the
771        // query path heals itself too (`PgPool::heal_types`), so a decode
772        // assertion alone would pass either way.
773        let oid = match pool
774            .query_typed(
775                "SELECT (oid::int8) AS result FROM pg_type WHERE typname = $1",
776                &[DecodedValue::Str(enum_type.clone())],
777                &pylon_pgcon::ExtensionOids::default(),
778            )
779            .await
780            .unwrap()
781            .first()
782        {
783            Some(DecodedValue::I64(oid)) => *oid as u32,
784            other => panic!("expected the new enum's oid, got {other:?}"),
785        };
786        assert!(
787            pool.types().enums.contains(&oid),
788            "apply_one must leave the registry knowing the type its DDL created"
789        );
790
791        let rows = pool
792            .query_composite(&format!("SELECT ('ok'::{enum_type}) AS result"), &pool.types())
793            .await
794            .unwrap();
795        assert_eq!(rows, vec![DecodedValue::Str("ok".to_string())]);
796
797        cleanup_migration_row(&pool, &m.id).await;
798        pool.batch_execute(&format!("DROP TYPE {enum_type}")).await.unwrap();
799    }
800
801    #[tokio::test]
802    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
803    async fn ensure_internal_schema_is_idempotent() {
804        let pool = test_pool().await;
805        ensure_internal_schema(&pool).await.unwrap();
806        ensure_internal_schema(&pool).await.unwrap();
807    }
808
809    #[test]
810    fn a_database_this_build_cannot_work_against_is_fatal_and_names_the_fix() {
811        let state = InternalSchemaState::TooOld { found: 1, required: 2 };
812        assert!(state.is_fatal());
813        let message = state.message().unwrap();
814        assert!(message.contains("pylon migration apply"), "got: {message}");
815    }
816
817    #[test]
818    fn a_newer_database_is_reported_but_never_fatal() {
819        // A rollback looks exactly like this. Failing closed here would turn
820        // the recovery lever into a second outage.
821        let state = InternalSchemaState::Newer { found: 2, current: 1 };
822        assert!(!state.is_fatal());
823        assert!(state.message().is_some());
824    }
825
826    #[test]
827    fn a_pending_upgrade_is_reported_but_not_fatal() {
828        let state = InternalSchemaState::Behind { found: 1, current: 2 };
829        assert!(!state.is_fatal());
830        assert!(state.message().is_some());
831    }
832
833    #[test]
834    fn an_unmigrated_or_current_database_says_nothing() {
835        // A database no migration has run against is a legal state, not a
836        // mismatch — callers must behave as they did before this existed.
837        for state in [InternalSchemaState::Unmigrated, InternalSchemaState::Current] {
838            assert!(!state.is_fatal());
839            assert_eq!(state.message(), None, "{state:?} should be silent");
840        }
841    }
842
843    #[tokio::test]
844    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
845    async fn a_freshly_ensured_database_reads_as_current() {
846        let pool = test_pool().await;
847        ensure_internal_schema(&pool).await.unwrap();
848        assert_eq!(
849            check_internal_schema(&pool).await.unwrap(),
850            InternalSchemaState::Current
851        );
852    }
853
854    #[tokio::test]
855    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
856    async fn a_database_without_the_marker_table_reads_as_unmigrated() {
857        // Not an error: `read_internal_version` has to tell "never migrated"
858        // apart from "migrated, and old", or every brand-new database would
859        // fail the check it is supposed to pass.
860        let pool = test_pool().await;
861        ensure_internal_schema(&pool).await.unwrap();
862        pool.batch_execute(r#"ALTER TABLE _pylon."Internal" RENAME TO "Internal_hidden";"#)
863            .await
864            .unwrap();
865        let state = check_internal_schema(&pool).await;
866        pool.batch_execute(r#"ALTER TABLE _pylon."Internal_hidden" RENAME TO "Internal";"#)
867            .await
868            .unwrap();
869        assert_eq!(state.unwrap(), InternalSchemaState::Unmigrated);
870    }
871
872    #[tokio::test]
873    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
874    async fn ensure_internal_schema_repairs_a_database_missing_a_newer_column() {
875        // The failure this whole change exists for: `claimed_at` was added
876        // to `_pylon."IndexOutbox"` with an `ADD COLUMN IF NOT EXISTS`
877        // alongside the worker lease, but that statement lived in a blob
878        // only `database initialize` ever shipped. Databases upgraded via
879        // `migration apply` never received it, and every index-worker claim
880        // failed against a column that was never added.
881        //
882        // Dropping the column reproduces such a database exactly.
883        let pool = test_pool().await;
884        ensure_internal_schema(&pool).await.unwrap();
885        pool.batch_execute(r#"ALTER TABLE _pylon."IndexOutbox" DROP COLUMN IF EXISTS claimed_at;"#)
886            .await
887            .unwrap();
888
889        ensure_internal_schema(&pool).await.unwrap();
890
891        let rows = pool
892            .query_typed(
893                "SELECT (count(*)) AS result FROM information_schema.columns \
894                 WHERE table_schema = '_pylon' AND table_name = 'IndexOutbox' \
895                 AND column_name = 'claimed_at'",
896                &[],
897                &pool.types(),
898            )
899            .await
900            .unwrap();
901        assert_eq!(
902            rows.into_iter().next(),
903            Some(DecodedValue::I64(1)),
904            "claimed_at should have been restored"
905        );
906    }
907
908    #[tokio::test]
909    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
910    async fn schema_snapshot_round_trips() {
911        // `_pylon."Schema"` is a shared singleton row across this whole test
912        // module's DSN (same non-isolation concern as `_pylon."Migrations"`
913        // — see `cleanup_migration_row`'s doc comment), so save and restore
914        // whatever was there before rather than leaving test data behind.
915        let pool = test_pool().await;
916        let previous = read_schema_snapshot(&pool).await.unwrap();
917
918        write_schema_snapshot(&pool, r#"{"probe": "schema_snapshot_round_trips"}"#)
919            .await
920            .unwrap();
921        let read_back = read_schema_snapshot(&pool).await.unwrap();
922        assert_eq!(
923            read_back.as_deref(),
924            Some(r#"{"probe": "schema_snapshot_round_trips"}"#)
925        );
926
927        // Upsert overwrites in place, so a second write must still round-trip
928        // (not silently keep the first value).
929        write_schema_snapshot(&pool, r#"{"probe": "second_write"}"#)
930            .await
931            .unwrap();
932        let read_back_2 = read_schema_snapshot(&pool).await.unwrap();
933        assert_eq!(read_back_2.as_deref(), Some(r#"{"probe": "second_write"}"#));
934
935        match previous {
936            Some(prior) => write_schema_snapshot(&pool, &prior).await.unwrap(),
937            None => pool.batch_execute(r#"DELETE FROM _pylon."Schema""#).await.unwrap(),
938        }
939    }
940
941    #[tokio::test]
942    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
943    async fn applied_tip_is_none_with_no_applied_rows() {
944        assert_eq!(applied_tip(&[]), None);
945        let all_pending = vec![TrackingRow {
946            id: "m1a".into(),
947            onto: "initial".into(),
948            db_state: None,
949            schema_state: None,
950            applied: false,
951        }];
952        assert_eq!(applied_tip(&all_pending), None);
953    }
954
955    #[tokio::test]
956    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
957    async fn applied_tip_is_the_row_with_no_descendant() {
958        let tracking = vec![
959            TrackingRow {
960                id: "m1a".into(),
961                onto: "initial".into(),
962                db_state: None,
963                schema_state: None,
964                applied: true,
965            },
966            TrackingRow {
967                id: "m1b".into(),
968                onto: "m1a".into(),
969                db_state: None,
970                schema_state: None,
971                applied: true,
972            },
973            TrackingRow {
974                id: "m1c".into(),
975                onto: "m1b".into(),
976                db_state: None,
977                schema_state: None,
978                applied: false,
979            }, // not applied yet
980        ];
981        assert_eq!(applied_tip(&tracking), Some("m1b".to_string()));
982    }
983
984    #[tokio::test]
985    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
986    async fn applied_tip_is_deterministic_with_multiple_orphaned_tips() {
987        // A tracking table can end up with several unrelated single-node
988        // "tips" (orphaned rows from deleted/regenerated migration files) —
989        // must always return the same answer, not one that depends on
990        // hash-map iteration order (see the doc comment on `applied_tip`).
991        let tracking = vec![
992            TrackingRow {
993                id: "m1zzz".into(),
994                onto: "initial".into(),
995                db_state: None,
996                schema_state: None,
997                applied: true,
998            },
999            TrackingRow {
1000                id: "m1aaa".into(),
1001                onto: "initial".into(),
1002                db_state: None,
1003                schema_state: None,
1004                applied: true,
1005            },
1006            TrackingRow {
1007                id: "m1mmm".into(),
1008                onto: "initial".into(),
1009                db_state: None,
1010                schema_state: None,
1011                applied: true,
1012            },
1013        ];
1014        for _ in 0..20 {
1015            assert_eq!(applied_tip(&tracking), Some("m1aaa".to_string()));
1016        }
1017    }
1018
1019    #[tokio::test]
1020    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1021    async fn apply_one_runs_ddl_and_records_tracking_row() {
1022        let pool = test_pool().await;
1023        let table = unique_table_name("migrate_apply_test");
1024        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
1025
1026        apply_one(&pool, &m, false).await.unwrap();
1027
1028        // DDL actually ran.
1029        let rows = pool
1030            .query_typed(
1031                &format!("SELECT (1) AS result FROM {table}"),
1032                &[],
1033                &pylon_pgcon::ExtensionOids::default(),
1034            )
1035            .await;
1036        assert!(rows.is_ok(), "table should exist after apply_one");
1037
1038        // Tracking row recorded.
1039        let tracking = read_tracking(&pool).await.unwrap();
1040        let row = tracking
1041            .iter()
1042            .find(|r| r.id == m.id)
1043            .expect("tracking row for this migration");
1044        assert!(row.applied);
1045        assert_eq!(row.onto, "initial");
1046        assert_eq!(
1047            row.db_state, None,
1048            "db_state is only ever set separately, by `migration create`'s own UPDATE"
1049        );
1050        assert_eq!(
1051            row.schema_state, None,
1052            "schema_state is only ever set separately, by `apply`'s own UPDATE"
1053        );
1054
1055        cleanup_migration_row(&pool, &m.id).await;
1056    }
1057
1058    #[tokio::test]
1059    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1060    async fn read_tracking_decodes_schema_state_as_raw_json_text() {
1061        let pool = test_pool().await;
1062        let m = make_migration("initial", &body("SELECT 1;"));
1063        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1064
1065        pool.execute_typed(
1066            r#"UPDATE _pylon."Migrations" SET schema_state = $1::jsonb WHERE id = $2"#,
1067            &[
1068                DecodedValue::Str(r#"{"types":[]}"#.to_string()),
1069                DecodedValue::Str(m.id.clone()),
1070            ],
1071        )
1072        .await
1073        .unwrap();
1074
1075        let tracking = read_tracking(&pool).await.unwrap();
1076        let row = tracking.iter().find(|r| r.id == m.id).unwrap();
1077        // Raw text, not a decoded DecodedValue::Object tree — `SchemaDescriptor::from_json`
1078        // (the pyo3-exposed consumer) re-parses this string itself.
1079        assert_eq!(row.schema_state.as_deref(), Some(r#"{"types": []}"#));
1080
1081        cleanup_migration_row(&pool, &m.id).await;
1082    }
1083
1084    #[tokio::test]
1085    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1086    async fn read_tracking_decodes_db_state_as_raw_json_text() {
1087        let pool = test_pool().await;
1088        let m = make_migration("initial", &body("SELECT 1;"));
1089        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1090
1091        pool.execute_typed(
1092            r#"UPDATE _pylon."Migrations" SET db_state = $1::jsonb WHERE id = $2"#,
1093            &[
1094                DecodedValue::Str(r#"{"schemas":["default"]}"#.to_string()),
1095                DecodedValue::Str(m.id.clone()),
1096            ],
1097        )
1098        .await
1099        .unwrap();
1100
1101        let tracking = read_tracking(&pool).await.unwrap();
1102        let row = tracking.iter().find(|r| r.id == m.id).unwrap();
1103        // Raw text, not a decoded DecodedValue::Object tree — `db_state_from_json`
1104        // (the pyo3-exposed consumer) re-parses this string itself.
1105        assert_eq!(row.db_state.as_deref(), Some(r#"{"schemas": ["default"]}"#));
1106
1107        cleanup_migration_row(&pool, &m.id).await;
1108    }
1109
1110    #[tokio::test]
1111    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1112    async fn apply_one_multi_step_clears_progress_after_completion() {
1113        let pool = test_pool().await;
1114        let t1 = unique_table_name("migrate_step1");
1115        let t2 = unique_table_name("migrate_step2");
1116        let m = make_migration(
1117            "initial",
1118            &format!("\nCREATE TABLE {t1} (id int8);\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
1119        );
1120
1121        apply_one(&pool, &m, false).await.unwrap();
1122
1123        for t in [&t1, &t2] {
1124            let rows = pool
1125                .query_typed(
1126                    &format!("SELECT (1) AS result FROM {t}"),
1127                    &[],
1128                    &pylon_pgcon::ExtensionOids::default(),
1129                )
1130                .await;
1131            assert!(rows.is_ok(), "table {t} should exist after apply_one");
1132        }
1133
1134        let progress = pool
1135            .query_typed(
1136                r#"SELECT (1) AS result FROM _pylon."Progress" WHERE id = $1"#,
1137                &[DecodedValue::Str(m.id.clone())],
1138                &pylon_pgcon::ExtensionOids::default(),
1139            )
1140            .await
1141            .unwrap();
1142        assert!(
1143            progress.is_empty(),
1144            "progress row must be cleared after a successful multi-step apply"
1145        );
1146
1147        cleanup_migration_row(&pool, &m.id).await;
1148    }
1149
1150    #[tokio::test]
1151    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1152    async fn apply_one_resumes_from_recorded_progress_skipping_earlier_steps() {
1153        let pool = test_pool().await;
1154        let t2 = unique_table_name("migrate_resume_step2");
1155        // Step 0 is intentionally invalid SQL — if apply_one didn't skip
1156        // it (via the pre-recorded progress row below), this test would
1157        // fail with a Postgres syntax error instead of succeeding.
1158        let m = make_migration(
1159            "initial",
1160            &format!("\nTHIS IS NOT VALID SQL;\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
1161        );
1162
1163        // Simulate a prior run that got through step 0 already.
1164        record_progress(&pool, &m.id, 0).await.unwrap();
1165
1166        apply_one(&pool, &m, false).await.unwrap();
1167
1168        let rows = pool
1169            .query_typed(
1170                &format!("SELECT (1) AS result FROM {t2}"),
1171                &[],
1172                &pylon_pgcon::ExtensionOids::default(),
1173            )
1174            .await;
1175        assert!(rows.is_ok(), "step 1 should have run");
1176
1177        cleanup_migration_row(&pool, &m.id).await;
1178    }
1179
1180    #[tokio::test]
1181    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1182    async fn apply_one_does_not_skip_the_step_it_failed_on_when_resumed() {
1183        let pool = test_pool().await;
1184        let t0 = unique_table_name("migrate_crash_step0");
1185        let t1 = unique_table_name("migrate_crash_step1");
1186        let t2 = unique_table_name("migrate_crash_step2");
1187
1188        // Step 1 fails on the first run because `t1` doesn't exist yet. The
1189        // point of the test is what the *second* run does with step 1: it
1190        // must retry it, not treat it as already done.
1191        let m = make_migration(
1192            "initial",
1193            &format!(
1194                "\nCREATE TABLE {t0} (id int8);\n\
1195                 -- pylon:step\n\
1196                 INSERT INTO {t1} (id) VALUES (1);\n\
1197                 -- pylon:step\n\
1198                 CREATE TABLE {t2} (id int8);\n"
1199            ),
1200        );
1201
1202        let first = apply_one(&pool, &m, false).await;
1203        assert!(first.is_err(), "step 1 should have failed on the first run");
1204
1205        // Progress must name step 0 — the last step that actually committed.
1206        // Recording step 1 here is the bug: it would make the retry resume at
1207        // step 2 and skip the INSERT forever.
1208        assert_eq!(
1209            read_progress(&pool, &m.id).await.unwrap(),
1210            Some(0),
1211            "progress must record the last *completed* step, not the one being attempted"
1212        );
1213
1214        // Make step 1 able to succeed, then resume.
1215        pool.batch_execute(&format!("CREATE TABLE {t1} (id int8);"))
1216            .await
1217            .unwrap();
1218        apply_one(&pool, &m, false).await.unwrap();
1219
1220        let rows = pool
1221            .query_typed(
1222                &format!("SELECT (count(*)) AS result FROM {t1}"),
1223                &[],
1224                &pylon_pgcon::ExtensionOids::default(),
1225            )
1226            .await
1227            .unwrap();
1228        assert_eq!(
1229            rows.first(),
1230            Some(&DecodedValue::I64(1)),
1231            "step 1 must have been retried on resume, not skipped"
1232        );
1233
1234        let t2_rows = pool
1235            .query_typed(
1236                &format!("SELECT (1) AS result FROM {t2}"),
1237                &[],
1238                &pylon_pgcon::ExtensionOids::default(),
1239            )
1240            .await;
1241        assert!(t2_rows.is_ok(), "step 2 should have run after the resumed step 1");
1242
1243        for t in [&t0, &t1, &t2] {
1244            pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
1245        }
1246        cleanup_migration_row(&pool, &m.id).await;
1247    }
1248
1249    #[tokio::test]
1250    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1251    async fn apply_one_dev_mode_skips_only_the_duplicate_statement_in_a_step() {
1252        let pool = test_pool().await;
1253        let before = unique_table_name("migrate_rebase_before");
1254        let existing = unique_table_name("migrate_rebase_existing");
1255        let after = unique_table_name("migrate_rebase_after");
1256
1257        // `watch` already created the middle table out of band.
1258        pool.batch_execute(&format!("CREATE TABLE {existing} (id int8);"))
1259            .await
1260            .unwrap();
1261
1262        // One step, three statements, only the middle one already applied.
1263        let m = make_migration(
1264            "initial",
1265            &body(&format!(
1266                "CREATE TABLE {before} (id int8);\n\
1267                 CREATE TABLE {existing} (id int8);\n\
1268                 CREATE TABLE {after} (id int8);"
1269            )),
1270        );
1271
1272        apply_one(&pool, &m, true).await.unwrap();
1273
1274        // Rolling back the whole step on the duplicate would drop `before`
1275        // and never reach `after`, while still recording the migration as
1276        // applied — the failure this test exists to catch.
1277        for t in [&before, &after] {
1278            let rows = pool
1279                .query_typed(
1280                    &format!("SELECT (1) AS result FROM {t}"),
1281                    &[],
1282                    &pylon_pgcon::ExtensionOids::default(),
1283                )
1284                .await;
1285            assert!(
1286                rows.is_ok(),
1287                "table {t} should exist — only the duplicate statement may be skipped"
1288            );
1289        }
1290
1291        let tracking = read_tracking(&pool).await.unwrap();
1292        assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
1293
1294        for t in [&before, &existing, &after] {
1295            pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
1296        }
1297        cleanup_migration_row(&pool, &m.id).await;
1298    }
1299
1300    #[tokio::test]
1301    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1302    async fn apply_one_dev_mode_swallows_a_duplicate_table_error() {
1303        let pool = test_pool().await;
1304        let table = unique_table_name("migrate_dev_mode_test");
1305        // Simulate `watch` having already applied this exact DDL out of band.
1306        pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
1307            .await
1308            .unwrap();
1309
1310        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
1311        apply_one(&pool, &m, true).await.unwrap(); // dev_mode=true: must not error
1312
1313        let tracking = read_tracking(&pool).await.unwrap();
1314        assert!(
1315            tracking.iter().any(|r| r.id == m.id && r.applied),
1316            "still recorded applied despite the swallowed error"
1317        );
1318
1319        cleanup_migration_row(&pool, &m.id).await;
1320    }
1321
1322    #[tokio::test]
1323    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1324    async fn apply_one_without_dev_mode_propagates_a_duplicate_table_error() {
1325        let pool = test_pool().await;
1326        let table = unique_table_name("migrate_no_dev_mode_test");
1327        pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
1328            .await
1329            .unwrap();
1330
1331        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
1332        let result = apply_one(&pool, &m, false).await; // dev_mode=false: must error
1333        assert!(result.is_err());
1334    }
1335
1336    #[tokio::test]
1337    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1338    async fn record_applied_standalone_marks_a_migration_applied_without_running_ddl() {
1339        let pool = test_pool().await;
1340        let m = make_migration("initial", &body("SELECT 1;"));
1341
1342        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1343
1344        let tracking = read_tracking(&pool).await.unwrap();
1345        assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
1346
1347        cleanup_migration_row(&pool, &m.id).await;
1348    }
1349
1350    #[tokio::test]
1351    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1352    async fn advisory_lock_round_trips_and_blocks_a_concurrent_try_lock() {
1353        let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
1354
1355        let held = advisory_lock(&pool).await.unwrap();
1356
1357        // A second, independent connection can't acquire the same key.
1358        let blocked = try_advisory_lock(&pool).await.unwrap();
1359        assert!(blocked.is_none(), "advisory lock should still be held");
1360
1361        advisory_unlock(held).await.unwrap();
1362
1363        // Now it's free again.
1364        let reacquired = try_advisory_lock(&pool).await.unwrap();
1365        assert!(reacquired.is_some());
1366        advisory_unlock(reacquired.unwrap()).await.unwrap();
1367    }
1368}