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    Ok(())
633}
634
635/// Runs one step's statements inside `tx`, each wrapped in its own savepoint,
636/// skipping any that fail because the object already exists.
637///
638/// The savepoint has to be per statement rather than per step: rolling the
639/// whole step back on the first duplicate would undo the statements before it
640/// *and* skip the ones after it, while still reporting the step as applied.
641async fn apply_statements_rebasing(tx: &PgTransaction, sql: &str) -> std::result::Result<(), pylon_pgcon::Error> {
642    for stmt in crate::migration::split_statements(sql) {
643        tx.savepoint(DEV_SAVEPOINT).await?;
644        match tx.batch_execute(&stmt).await {
645            Ok(()) => tx.release_savepoint(DEV_SAVEPOINT).await?,
646            Err(e) if is_duplicate_object_error(&e) => tx.rollback_to_savepoint(DEV_SAVEPOINT).await?,
647            Err(e) => return Err(e),
648        }
649    }
650    Ok(())
651}
652
653#[cfg(test)]
654mod concurrent_index_name_tests {
655    use super::concurrent_index_name;
656
657    #[test]
658    fn extracts_a_bare_index_name() {
659        assert_eq!(
660            concurrent_index_name("CREATE INDEX CONCURRENTLY idx_person_name ON \"public\".\"Person\" (name);"),
661            Some("idx_person_name".to_string())
662        );
663    }
664
665    #[test]
666    fn extracts_a_quoted_index_name() {
667        assert_eq!(
668            concurrent_index_name("CREATE INDEX CONCURRENTLY \"idx_person_name\" ON \"public\".\"Person\" (name);"),
669            Some("idx_person_name".to_string())
670        );
671    }
672
673    #[test]
674    fn handles_if_not_exists() {
675        assert_eq!(
676            concurrent_index_name("CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_x ON t (c);"),
677            Some("idx_x".to_string())
678        );
679    }
680
681    #[test]
682    fn is_case_insensitive() {
683        assert_eq!(
684            concurrent_index_name("create index concurrently idx_x on t (c);"),
685            Some("idx_x".to_string())
686        );
687    }
688
689    #[test]
690    fn returns_none_for_unrelated_sql() {
691        assert_eq!(concurrent_index_name("CREATE TABLE foo ();"), None);
692        assert_eq!(concurrent_index_name("CREATE INDEX idx_x ON t (c);"), None); // not CONCURRENTLY
693    }
694}
695
696#[cfg(test)]
697mod tests {
698    use super::*;
699    use crate::migration::render_file;
700
701    fn test_dsn() -> String {
702        std::env::var("PYLON_PGCON_TEST_DSN").expect("PYLON_PGCON_TEST_DSN must be set to run live-Postgres tests")
703    }
704
705    async fn test_pool() -> PgPool {
706        let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
707        pool.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
708        ensure_internal_schema(&pool).await.unwrap();
709        pool
710    }
711
712    fn make_migration(onto: &str, body: &str) -> MigrationFile {
713        let content = render_file(onto, body, &[]);
714        crate::migration::parse(&content, "test").unwrap()
715    }
716
717    /// A body with the leading blank-line separator `render_file`/`parse`
718    /// expect, so the migration's computed ID matches what gets hashed.
719    fn body(sql: &str) -> String {
720        format!("\n{sql}\n")
721    }
722
723    fn unique_table_name(prefix: &str) -> String {
724        use std::time::{SystemTime, UNIX_EPOCH};
725        let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
726        format!("{prefix}_{nanos}")
727    }
728
729    /// Every test in this module runs against the same shared,
730    /// non-isolated `_pylon."Migrations"` table (unlike the rest of the
731    /// live-execution suite, which gets a fresh schema per test via
732    /// `unique_module()` — there's no equivalent scoping for this
733    /// process-wide tracking table). Any test that records a tracking row
734    /// must delete it again here, or it permanently pollutes whatever
735    /// database `PYLON_PGCON_TEST_DSN` points at (this defaults to the
736    /// same DSN pylon-demo uses, and a stray row here can shadow a real
737    /// project's actual migration tip).
738    async fn cleanup_migration_row(pool: &PgPool, id: &str) {
739        pool.execute_typed(
740            r#"DELETE FROM _pylon."Migrations" WHERE id = $1"#,
741            &[DecodedValue::Str(id.to_string())],
742        )
743        .await
744        .unwrap();
745    }
746
747    #[tokio::test]
748    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
749    async fn ensure_internal_schema_is_idempotent() {
750        let pool = test_pool().await;
751        ensure_internal_schema(&pool).await.unwrap();
752        ensure_internal_schema(&pool).await.unwrap();
753    }
754
755    #[test]
756    fn a_database_this_build_cannot_work_against_is_fatal_and_names_the_fix() {
757        let state = InternalSchemaState::TooOld { found: 1, required: 2 };
758        assert!(state.is_fatal());
759        let message = state.message().unwrap();
760        assert!(message.contains("pylon migration apply"), "got: {message}");
761    }
762
763    #[test]
764    fn a_newer_database_is_reported_but_never_fatal() {
765        // A rollback looks exactly like this. Failing closed here would turn
766        // the recovery lever into a second outage.
767        let state = InternalSchemaState::Newer { found: 2, current: 1 };
768        assert!(!state.is_fatal());
769        assert!(state.message().is_some());
770    }
771
772    #[test]
773    fn a_pending_upgrade_is_reported_but_not_fatal() {
774        let state = InternalSchemaState::Behind { found: 1, current: 2 };
775        assert!(!state.is_fatal());
776        assert!(state.message().is_some());
777    }
778
779    #[test]
780    fn an_unmigrated_or_current_database_says_nothing() {
781        // A database no migration has run against is a legal state, not a
782        // mismatch — callers must behave as they did before this existed.
783        for state in [InternalSchemaState::Unmigrated, InternalSchemaState::Current] {
784            assert!(!state.is_fatal());
785            assert_eq!(state.message(), None, "{state:?} should be silent");
786        }
787    }
788
789    #[tokio::test]
790    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
791    async fn a_freshly_ensured_database_reads_as_current() {
792        let pool = test_pool().await;
793        ensure_internal_schema(&pool).await.unwrap();
794        assert_eq!(
795            check_internal_schema(&pool).await.unwrap(),
796            InternalSchemaState::Current
797        );
798    }
799
800    #[tokio::test]
801    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
802    async fn a_database_without_the_marker_table_reads_as_unmigrated() {
803        // Not an error: `read_internal_version` has to tell "never migrated"
804        // apart from "migrated, and old", or every brand-new database would
805        // fail the check it is supposed to pass.
806        let pool = test_pool().await;
807        ensure_internal_schema(&pool).await.unwrap();
808        pool.batch_execute(r#"ALTER TABLE _pylon."Internal" RENAME TO "Internal_hidden";"#)
809            .await
810            .unwrap();
811        let state = check_internal_schema(&pool).await;
812        pool.batch_execute(r#"ALTER TABLE _pylon."Internal_hidden" RENAME TO "Internal";"#)
813            .await
814            .unwrap();
815        assert_eq!(state.unwrap(), InternalSchemaState::Unmigrated);
816    }
817
818    #[tokio::test]
819    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
820    async fn ensure_internal_schema_repairs_a_database_missing_a_newer_column() {
821        // The failure this whole change exists for: `claimed_at` was added
822        // to `_pylon."IndexOutbox"` with an `ADD COLUMN IF NOT EXISTS`
823        // alongside the worker lease, but that statement lived in a blob
824        // only `database initialize` ever shipped. Databases upgraded via
825        // `migration apply` never received it, and every index-worker claim
826        // failed against a column that was never added.
827        //
828        // Dropping the column reproduces such a database exactly.
829        let pool = test_pool().await;
830        ensure_internal_schema(&pool).await.unwrap();
831        pool.batch_execute(r#"ALTER TABLE _pylon."IndexOutbox" DROP COLUMN IF EXISTS claimed_at;"#)
832            .await
833            .unwrap();
834
835        ensure_internal_schema(&pool).await.unwrap();
836
837        let rows = pool
838            .query_typed(
839                "SELECT (count(*)) AS result FROM information_schema.columns \
840                 WHERE table_schema = '_pylon' AND table_name = 'IndexOutbox' \
841                 AND column_name = 'claimed_at'",
842                &[],
843                pool.types(),
844            )
845            .await
846            .unwrap();
847        assert_eq!(
848            rows.into_iter().next(),
849            Some(DecodedValue::I64(1)),
850            "claimed_at should have been restored"
851        );
852    }
853
854    #[tokio::test]
855    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
856    async fn schema_snapshot_round_trips() {
857        // `_pylon."Schema"` is a shared singleton row across this whole test
858        // module's DSN (same non-isolation concern as `_pylon."Migrations"`
859        // — see `cleanup_migration_row`'s doc comment), so save and restore
860        // whatever was there before rather than leaving test data behind.
861        let pool = test_pool().await;
862        let previous = read_schema_snapshot(&pool).await.unwrap();
863
864        write_schema_snapshot(&pool, r#"{"probe": "schema_snapshot_round_trips"}"#)
865            .await
866            .unwrap();
867        let read_back = read_schema_snapshot(&pool).await.unwrap();
868        assert_eq!(
869            read_back.as_deref(),
870            Some(r#"{"probe": "schema_snapshot_round_trips"}"#)
871        );
872
873        // Upsert overwrites in place, so a second write must still round-trip
874        // (not silently keep the first value).
875        write_schema_snapshot(&pool, r#"{"probe": "second_write"}"#)
876            .await
877            .unwrap();
878        let read_back_2 = read_schema_snapshot(&pool).await.unwrap();
879        assert_eq!(read_back_2.as_deref(), Some(r#"{"probe": "second_write"}"#));
880
881        match previous {
882            Some(prior) => write_schema_snapshot(&pool, &prior).await.unwrap(),
883            None => pool.batch_execute(r#"DELETE FROM _pylon."Schema""#).await.unwrap(),
884        }
885    }
886
887    #[tokio::test]
888    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
889    async fn applied_tip_is_none_with_no_applied_rows() {
890        assert_eq!(applied_tip(&[]), None);
891        let all_pending = vec![TrackingRow {
892            id: "m1a".into(),
893            onto: "initial".into(),
894            db_state: None,
895            schema_state: None,
896            applied: false,
897        }];
898        assert_eq!(applied_tip(&all_pending), None);
899    }
900
901    #[tokio::test]
902    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
903    async fn applied_tip_is_the_row_with_no_descendant() {
904        let tracking = vec![
905            TrackingRow {
906                id: "m1a".into(),
907                onto: "initial".into(),
908                db_state: None,
909                schema_state: None,
910                applied: true,
911            },
912            TrackingRow {
913                id: "m1b".into(),
914                onto: "m1a".into(),
915                db_state: None,
916                schema_state: None,
917                applied: true,
918            },
919            TrackingRow {
920                id: "m1c".into(),
921                onto: "m1b".into(),
922                db_state: None,
923                schema_state: None,
924                applied: false,
925            }, // not applied yet
926        ];
927        assert_eq!(applied_tip(&tracking), Some("m1b".to_string()));
928    }
929
930    #[tokio::test]
931    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
932    async fn applied_tip_is_deterministic_with_multiple_orphaned_tips() {
933        // A tracking table can end up with several unrelated single-node
934        // "tips" (orphaned rows from deleted/regenerated migration files) —
935        // must always return the same answer, not one that depends on
936        // hash-map iteration order (see the doc comment on `applied_tip`).
937        let tracking = vec![
938            TrackingRow {
939                id: "m1zzz".into(),
940                onto: "initial".into(),
941                db_state: None,
942                schema_state: None,
943                applied: true,
944            },
945            TrackingRow {
946                id: "m1aaa".into(),
947                onto: "initial".into(),
948                db_state: None,
949                schema_state: None,
950                applied: true,
951            },
952            TrackingRow {
953                id: "m1mmm".into(),
954                onto: "initial".into(),
955                db_state: None,
956                schema_state: None,
957                applied: true,
958            },
959        ];
960        for _ in 0..20 {
961            assert_eq!(applied_tip(&tracking), Some("m1aaa".to_string()));
962        }
963    }
964
965    #[tokio::test]
966    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
967    async fn apply_one_runs_ddl_and_records_tracking_row() {
968        let pool = test_pool().await;
969        let table = unique_table_name("migrate_apply_test");
970        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
971
972        apply_one(&pool, &m, false).await.unwrap();
973
974        // DDL actually ran.
975        let rows = pool
976            .query_typed(
977                &format!("SELECT (1) AS result FROM {table}"),
978                &[],
979                &pylon_pgcon::ExtensionOids::default(),
980            )
981            .await;
982        assert!(rows.is_ok(), "table should exist after apply_one");
983
984        // Tracking row recorded.
985        let tracking = read_tracking(&pool).await.unwrap();
986        let row = tracking
987            .iter()
988            .find(|r| r.id == m.id)
989            .expect("tracking row for this migration");
990        assert!(row.applied);
991        assert_eq!(row.onto, "initial");
992        assert_eq!(
993            row.db_state, None,
994            "db_state is only ever set separately, by `migration create`'s own UPDATE"
995        );
996        assert_eq!(
997            row.schema_state, None,
998            "schema_state is only ever set separately, by `apply`'s own UPDATE"
999        );
1000
1001        cleanup_migration_row(&pool, &m.id).await;
1002    }
1003
1004    #[tokio::test]
1005    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1006    async fn read_tracking_decodes_schema_state_as_raw_json_text() {
1007        let pool = test_pool().await;
1008        let m = make_migration("initial", &body("SELECT 1;"));
1009        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1010
1011        pool.execute_typed(
1012            r#"UPDATE _pylon."Migrations" SET schema_state = $1::jsonb WHERE id = $2"#,
1013            &[
1014                DecodedValue::Str(r#"{"types":[]}"#.to_string()),
1015                DecodedValue::Str(m.id.clone()),
1016            ],
1017        )
1018        .await
1019        .unwrap();
1020
1021        let tracking = read_tracking(&pool).await.unwrap();
1022        let row = tracking.iter().find(|r| r.id == m.id).unwrap();
1023        // Raw text, not a decoded DecodedValue::Object tree — `SchemaDescriptor::from_json`
1024        // (the pyo3-exposed consumer) re-parses this string itself.
1025        assert_eq!(row.schema_state.as_deref(), Some(r#"{"types": []}"#));
1026
1027        cleanup_migration_row(&pool, &m.id).await;
1028    }
1029
1030    #[tokio::test]
1031    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1032    async fn read_tracking_decodes_db_state_as_raw_json_text() {
1033        let pool = test_pool().await;
1034        let m = make_migration("initial", &body("SELECT 1;"));
1035        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1036
1037        pool.execute_typed(
1038            r#"UPDATE _pylon."Migrations" SET db_state = $1::jsonb WHERE id = $2"#,
1039            &[
1040                DecodedValue::Str(r#"{"schemas":["default"]}"#.to_string()),
1041                DecodedValue::Str(m.id.clone()),
1042            ],
1043        )
1044        .await
1045        .unwrap();
1046
1047        let tracking = read_tracking(&pool).await.unwrap();
1048        let row = tracking.iter().find(|r| r.id == m.id).unwrap();
1049        // Raw text, not a decoded DecodedValue::Object tree — `db_state_from_json`
1050        // (the pyo3-exposed consumer) re-parses this string itself.
1051        assert_eq!(row.db_state.as_deref(), Some(r#"{"schemas": ["default"]}"#));
1052
1053        cleanup_migration_row(&pool, &m.id).await;
1054    }
1055
1056    #[tokio::test]
1057    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1058    async fn apply_one_multi_step_clears_progress_after_completion() {
1059        let pool = test_pool().await;
1060        let t1 = unique_table_name("migrate_step1");
1061        let t2 = unique_table_name("migrate_step2");
1062        let m = make_migration(
1063            "initial",
1064            &format!("\nCREATE TABLE {t1} (id int8);\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
1065        );
1066
1067        apply_one(&pool, &m, false).await.unwrap();
1068
1069        for t in [&t1, &t2] {
1070            let rows = pool
1071                .query_typed(
1072                    &format!("SELECT (1) AS result FROM {t}"),
1073                    &[],
1074                    &pylon_pgcon::ExtensionOids::default(),
1075                )
1076                .await;
1077            assert!(rows.is_ok(), "table {t} should exist after apply_one");
1078        }
1079
1080        let progress = pool
1081            .query_typed(
1082                r#"SELECT (1) AS result FROM _pylon."Progress" WHERE id = $1"#,
1083                &[DecodedValue::Str(m.id.clone())],
1084                &pylon_pgcon::ExtensionOids::default(),
1085            )
1086            .await
1087            .unwrap();
1088        assert!(
1089            progress.is_empty(),
1090            "progress row must be cleared after a successful multi-step apply"
1091        );
1092
1093        cleanup_migration_row(&pool, &m.id).await;
1094    }
1095
1096    #[tokio::test]
1097    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1098    async fn apply_one_resumes_from_recorded_progress_skipping_earlier_steps() {
1099        let pool = test_pool().await;
1100        let t2 = unique_table_name("migrate_resume_step2");
1101        // Step 0 is intentionally invalid SQL — if apply_one didn't skip
1102        // it (via the pre-recorded progress row below), this test would
1103        // fail with a Postgres syntax error instead of succeeding.
1104        let m = make_migration(
1105            "initial",
1106            &format!("\nTHIS IS NOT VALID SQL;\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
1107        );
1108
1109        // Simulate a prior run that got through step 0 already.
1110        record_progress(&pool, &m.id, 0).await.unwrap();
1111
1112        apply_one(&pool, &m, false).await.unwrap();
1113
1114        let rows = pool
1115            .query_typed(
1116                &format!("SELECT (1) AS result FROM {t2}"),
1117                &[],
1118                &pylon_pgcon::ExtensionOids::default(),
1119            )
1120            .await;
1121        assert!(rows.is_ok(), "step 1 should have run");
1122
1123        cleanup_migration_row(&pool, &m.id).await;
1124    }
1125
1126    #[tokio::test]
1127    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1128    async fn apply_one_does_not_skip_the_step_it_failed_on_when_resumed() {
1129        let pool = test_pool().await;
1130        let t0 = unique_table_name("migrate_crash_step0");
1131        let t1 = unique_table_name("migrate_crash_step1");
1132        let t2 = unique_table_name("migrate_crash_step2");
1133
1134        // Step 1 fails on the first run because `t1` doesn't exist yet. The
1135        // point of the test is what the *second* run does with step 1: it
1136        // must retry it, not treat it as already done.
1137        let m = make_migration(
1138            "initial",
1139            &format!(
1140                "\nCREATE TABLE {t0} (id int8);\n\
1141                 -- pylon:step\n\
1142                 INSERT INTO {t1} (id) VALUES (1);\n\
1143                 -- pylon:step\n\
1144                 CREATE TABLE {t2} (id int8);\n"
1145            ),
1146        );
1147
1148        let first = apply_one(&pool, &m, false).await;
1149        assert!(first.is_err(), "step 1 should have failed on the first run");
1150
1151        // Progress must name step 0 — the last step that actually committed.
1152        // Recording step 1 here is the bug: it would make the retry resume at
1153        // step 2 and skip the INSERT forever.
1154        assert_eq!(
1155            read_progress(&pool, &m.id).await.unwrap(),
1156            Some(0),
1157            "progress must record the last *completed* step, not the one being attempted"
1158        );
1159
1160        // Make step 1 able to succeed, then resume.
1161        pool.batch_execute(&format!("CREATE TABLE {t1} (id int8);"))
1162            .await
1163            .unwrap();
1164        apply_one(&pool, &m, false).await.unwrap();
1165
1166        let rows = pool
1167            .query_typed(
1168                &format!("SELECT (count(*)) AS result FROM {t1}"),
1169                &[],
1170                &pylon_pgcon::ExtensionOids::default(),
1171            )
1172            .await
1173            .unwrap();
1174        assert_eq!(
1175            rows.first(),
1176            Some(&DecodedValue::I64(1)),
1177            "step 1 must have been retried on resume, not skipped"
1178        );
1179
1180        let t2_rows = pool
1181            .query_typed(
1182                &format!("SELECT (1) AS result FROM {t2}"),
1183                &[],
1184                &pylon_pgcon::ExtensionOids::default(),
1185            )
1186            .await;
1187        assert!(t2_rows.is_ok(), "step 2 should have run after the resumed step 1");
1188
1189        for t in [&t0, &t1, &t2] {
1190            pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
1191        }
1192        cleanup_migration_row(&pool, &m.id).await;
1193    }
1194
1195    #[tokio::test]
1196    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1197    async fn apply_one_dev_mode_skips_only_the_duplicate_statement_in_a_step() {
1198        let pool = test_pool().await;
1199        let before = unique_table_name("migrate_rebase_before");
1200        let existing = unique_table_name("migrate_rebase_existing");
1201        let after = unique_table_name("migrate_rebase_after");
1202
1203        // `watch` already created the middle table out of band.
1204        pool.batch_execute(&format!("CREATE TABLE {existing} (id int8);"))
1205            .await
1206            .unwrap();
1207
1208        // One step, three statements, only the middle one already applied.
1209        let m = make_migration(
1210            "initial",
1211            &body(&format!(
1212                "CREATE TABLE {before} (id int8);\n\
1213                 CREATE TABLE {existing} (id int8);\n\
1214                 CREATE TABLE {after} (id int8);"
1215            )),
1216        );
1217
1218        apply_one(&pool, &m, true).await.unwrap();
1219
1220        // Rolling back the whole step on the duplicate would drop `before`
1221        // and never reach `after`, while still recording the migration as
1222        // applied — the failure this test exists to catch.
1223        for t in [&before, &after] {
1224            let rows = pool
1225                .query_typed(
1226                    &format!("SELECT (1) AS result FROM {t}"),
1227                    &[],
1228                    &pylon_pgcon::ExtensionOids::default(),
1229                )
1230                .await;
1231            assert!(
1232                rows.is_ok(),
1233                "table {t} should exist — only the duplicate statement may be skipped"
1234            );
1235        }
1236
1237        let tracking = read_tracking(&pool).await.unwrap();
1238        assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
1239
1240        for t in [&before, &existing, &after] {
1241            pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
1242        }
1243        cleanup_migration_row(&pool, &m.id).await;
1244    }
1245
1246    #[tokio::test]
1247    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1248    async fn apply_one_dev_mode_swallows_a_duplicate_table_error() {
1249        let pool = test_pool().await;
1250        let table = unique_table_name("migrate_dev_mode_test");
1251        // Simulate `watch` having already applied this exact DDL out of band.
1252        pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
1253            .await
1254            .unwrap();
1255
1256        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
1257        apply_one(&pool, &m, true).await.unwrap(); // dev_mode=true: must not error
1258
1259        let tracking = read_tracking(&pool).await.unwrap();
1260        assert!(
1261            tracking.iter().any(|r| r.id == m.id && r.applied),
1262            "still recorded applied despite the swallowed error"
1263        );
1264
1265        cleanup_migration_row(&pool, &m.id).await;
1266    }
1267
1268    #[tokio::test]
1269    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1270    async fn apply_one_without_dev_mode_propagates_a_duplicate_table_error() {
1271        let pool = test_pool().await;
1272        let table = unique_table_name("migrate_no_dev_mode_test");
1273        pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
1274            .await
1275            .unwrap();
1276
1277        let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
1278        let result = apply_one(&pool, &m, false).await; // dev_mode=false: must error
1279        assert!(result.is_err());
1280    }
1281
1282    #[tokio::test]
1283    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1284    async fn record_applied_standalone_marks_a_migration_applied_without_running_ddl() {
1285        let pool = test_pool().await;
1286        let m = make_migration("initial", &body("SELECT 1;"));
1287
1288        record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
1289
1290        let tracking = read_tracking(&pool).await.unwrap();
1291        assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
1292
1293        cleanup_migration_row(&pool, &m.id).await;
1294    }
1295
1296    #[tokio::test]
1297    #[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
1298    async fn advisory_lock_round_trips_and_blocks_a_concurrent_try_lock() {
1299        let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
1300
1301        let held = advisory_lock(&pool).await.unwrap();
1302
1303        // A second, independent connection can't acquire the same key.
1304        let blocked = try_advisory_lock(&pool).await.unwrap();
1305        assert!(blocked.is_none(), "advisory lock should still be held");
1306
1307        advisory_unlock(held).await.unwrap();
1308
1309        // Now it's free again.
1310        let reacquired = try_advisory_lock(&pool).await.unwrap();
1311        assert!(reacquired.is_some());
1312        advisory_unlock(reacquired.unwrap()).await.unwrap();
1313    }
1314}