use crate::migration::{MigrationFile, parse_steps, verify_integrity};
use pylon_pgcon::{PgPool, PgTransaction};
use pylon_value::DecodedValue;
#[derive(Debug, thiserror::Error)]
pub enum MigrateError {
#[error(transparent)]
Integrity(#[from] crate::migration::MigrationError),
#[error(transparent)]
Db(#[from] pylon_pgcon::Error),
#[error(
"migration {id} is already recorded as applied onto {recorded_onto}, but the file \
being applied claims onto {new_onto} — two different migrations share one ID"
)]
IdCollision {
id: String,
recorded_onto: String,
new_onto: String,
},
}
pub type Result<T> = std::result::Result<T, MigrateError>;
pub const ADVISORY_LOCK_KEY: i64 = 7_461_999;
const DUPLICATE_OBJECT_CODES: [tokio_postgres::error::SqlState; 5] = [
tokio_postgres::error::SqlState::DUPLICATE_TABLE,
tokio_postgres::error::SqlState::DUPLICATE_COLUMN,
tokio_postgres::error::SqlState::DUPLICATE_SCHEMA,
tokio_postgres::error::SqlState::DUPLICATE_OBJECT,
tokio_postgres::error::SqlState::DUPLICATE_DATABASE,
];
fn is_duplicate_object_error(err: &pylon_pgcon::Error) -> bool {
err.sqlstate().is_some_and(|code| DUPLICATE_OBJECT_CODES.contains(code))
}
pub async fn ensure_internal_schema(pool: &PgPool) -> Result<()> {
pool.batch_execute(&crate::stdlib::export_stdlib()).await?;
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum InternalSchemaState {
Unmigrated,
TooOld { found: i32, required: i32 },
Behind { found: i32, current: i32 },
Current,
Newer { found: i32, current: i32 },
}
impl InternalSchemaState {
pub fn is_fatal(&self) -> bool {
matches!(self, InternalSchemaState::TooOld { .. })
}
pub fn message(&self) -> Option<String> {
match self {
InternalSchemaState::Unmigrated | InternalSchemaState::Current => None,
InternalSchemaState::TooOld { found, required } => Some(format!(
"this database's internal schema (version {found}) is older than this \
version of Pylon supports (version {required}); run `pylon migration apply` \
to bring it up to date"
)),
InternalSchemaState::Behind { found, current } => Some(format!(
"this database's internal schema is at version {found}, this version of \
Pylon writes version {current}; `pylon migration apply` will update it"
)),
InternalSchemaState::Newer { found, current } => Some(format!(
"this database's internal schema (version {found}) was written by a newer \
version of Pylon than this one (version {current}); continuing, but this \
build may not understand everything it finds"
)),
}
}
}
pub async fn read_internal_version(pool: &PgPool) -> Result<Option<i32>> {
let rows = match pool
.query_typed(
r#"SELECT (version) AS result FROM _pylon."Internal" WHERE singleton"#,
&[],
pool.types(),
)
.await
{
Ok(rows) => rows,
Err(e) if e.sqlstate() == Some(&tokio_postgres::error::SqlState::UNDEFINED_TABLE) => return Ok(None),
Err(e) => return Err(e.into()),
};
Ok(match rows.into_iter().next() {
Some(DecodedValue::I64(v)) => Some(v as i32),
_ => None,
})
}
pub async fn check_internal_schema(pool: &PgPool) -> Result<InternalSchemaState> {
use crate::stdlib::ddl::{INTERNAL_SCHEMA_VERSION, MIN_SUPPORTED_INTERNAL_VERSION};
Ok(match read_internal_version(pool).await? {
None => InternalSchemaState::Unmigrated,
Some(found) if found < MIN_SUPPORTED_INTERNAL_VERSION => InternalSchemaState::TooOld {
found,
required: MIN_SUPPORTED_INTERNAL_VERSION,
},
Some(found) if found < INTERNAL_SCHEMA_VERSION => InternalSchemaState::Behind {
found,
current: INTERNAL_SCHEMA_VERSION,
},
Some(found) if found > INTERNAL_SCHEMA_VERSION => InternalSchemaState::Newer {
found,
current: INTERNAL_SCHEMA_VERSION,
},
Some(_) => InternalSchemaState::Current,
})
}
pub async fn write_schema_snapshot(pool: &PgPool, snapshot_json: &str) -> Result<()> {
pool.execute_typed(
r#"INSERT INTO _pylon."Schema" (singleton, snapshot, updated_at) VALUES (true, $1::jsonb, now())
ON CONFLICT (singleton) DO UPDATE SET snapshot = $1::jsonb, updated_at = now()"#,
&[DecodedValue::Str(snapshot_json.to_string())],
)
.await?;
Ok(())
}
pub async fn read_schema_snapshot(pool: &PgPool) -> Result<Option<String>> {
let rows = pool
.query_typed(
r#"SELECT (snapshot::text) AS result FROM _pylon."Schema" WHERE singleton"#,
&[],
pool.types(),
)
.await?;
Ok(match rows.into_iter().next() {
Some(DecodedValue::Str(s)) => Some(s),
_ => None,
})
}
#[derive(Debug, Clone, PartialEq)]
pub struct TrackingRow {
pub id: String,
pub onto: String,
pub db_state: Option<String>,
pub schema_state: Option<String>,
pub applied: bool,
}
pub async fn read_tracking(pool: &PgPool) -> Result<Vec<TrackingRow>> {
let rows = pool
.query_typed(
r#"SELECT (id, onto, (db_state::text), (schema_state::text), (applied_at IS NOT NULL)) AS result FROM _pylon."Migrations""#,
&[],
pool.types(),
)
.await?;
Ok(rows
.into_iter()
.filter_map(|row| {
let DecodedValue::Composite(fields) = row else {
return None;
};
let [
DecodedValue::Str(id),
DecodedValue::Str(onto),
db_state,
schema_state,
DecodedValue::Bool(applied),
] = <[DecodedValue; 5]>::try_from(fields).ok()?
else {
return None;
};
let db_state = match db_state {
DecodedValue::Str(s) => Some(s),
_ => None,
};
let schema_state = match schema_state {
DecodedValue::Str(s) => Some(s),
_ => None,
};
Some(TrackingRow {
id,
onto,
db_state,
schema_state,
applied,
})
})
.collect())
}
pub fn applied_tip(tracking: &[TrackingRow]) -> Option<String> {
let applied: Vec<&TrackingRow> = tracking.iter().filter(|r| r.applied).collect();
if applied.is_empty() {
return None;
}
let onto_targets: std::collections::HashSet<&str> = applied.iter().map(|r| r.onto.as_str()).collect();
applied
.iter()
.filter(|r| !onto_targets.contains(r.id.as_str()))
.map(|r| r.id.as_str())
.min()
.map(|s| s.to_string())
}
pub async fn advisory_lock(pool: &PgPool) -> Result<pylon_pgcon::PgConnection> {
let conn = pool.connection().await?;
conn.batch_execute(&format!("SELECT pg_advisory_lock({ADVISORY_LOCK_KEY})"))
.await?;
Ok(conn)
}
pub async fn try_advisory_lock(pool: &PgPool) -> Result<Option<pylon_pgcon::PgConnection>> {
let conn = pool.connection().await?;
let rows = conn
.query_typed(
&format!("SELECT (pg_try_advisory_lock({ADVISORY_LOCK_KEY})) AS result"),
&[],
pool.types(),
)
.await?;
Ok(if matches!(rows.first(), Some(DecodedValue::Bool(true))) {
Some(conn)
} else {
None
})
}
pub async fn advisory_unlock(conn: pylon_pgcon::PgConnection) -> Result<()> {
conn.batch_execute(&format!("SELECT pg_advisory_unlock({ADVISORY_LOCK_KEY})"))
.await?;
Ok(())
}
const RECORD_APPLIED_SQL: &str = r#"
INSERT INTO _pylon."Migrations" (id, onto, filename, applied_at)
VALUES ($1, $2, $3, now())
ON CONFLICT (id) DO UPDATE
SET applied_at = now(),
onto = EXCLUDED.onto,
filename = EXCLUDED.filename
"#;
async fn check_no_id_collision(pool: &PgPool, id: &str, onto: &str) -> Result<()> {
let rows = pool
.query_typed(
r#"SELECT (onto) AS result FROM _pylon."Migrations" WHERE id = $1"#,
&[DecodedValue::Str(id.to_string())],
pool.types(),
)
.await?;
if let Some(DecodedValue::Str(existing_onto)) = rows.into_iter().next()
&& existing_onto != onto
{
return Err(MigrateError::IdCollision {
id: id.to_string(),
recorded_onto: existing_onto,
new_onto: onto.to_string(),
});
}
Ok(())
}
fn record_applied_params(id: &str, onto: &str, filename: &str) -> Vec<DecodedValue> {
vec![
DecodedValue::Str(id.to_string()),
DecodedValue::Str(onto.to_string()),
DecodedValue::Str(filename.to_string()),
]
}
pub async fn record_applied(pool: &PgPool, id: &str, onto: &str, filename: &str) -> Result<()> {
check_no_id_collision(pool, id, onto).await?;
pool.execute_typed(RECORD_APPLIED_SQL, &record_applied_params(id, onto, filename))
.await?;
Ok(())
}
async fn record_applied_in_tx(tx: &PgTransaction, id: &str, onto: &str, filename: &str) -> Result<()> {
tx.execute_typed(RECORD_APPLIED_SQL, &record_applied_params(id, onto, filename))
.await?;
Ok(())
}
async fn read_progress(pool: &PgPool, id: &str) -> Result<Option<i64>> {
let rows = pool
.query_typed(
r#"SELECT (step_index) AS result FROM _pylon."Progress" WHERE id = $1"#,
&[DecodedValue::Str(id.to_string())],
pool.types(),
)
.await?;
Ok(match rows.into_iter().next() {
Some(DecodedValue::I64(n)) => Some(n),
_ => None,
})
}
async fn record_progress(pool: &PgPool, id: &str, step_index: i64) -> Result<()> {
pool.execute_typed(
r#"INSERT INTO _pylon."Progress" (id, step_index) VALUES ($1, $2)
ON CONFLICT (id) DO UPDATE SET step_index = $2, updated_at = now()"#,
&[DecodedValue::Str(id.to_string()), DecodedValue::I64(step_index)],
)
.await?;
Ok(())
}
async fn delete_progress(pool: &PgPool, id: &str) -> Result<()> {
pool.execute_typed(
r#"DELETE FROM _pylon."Progress" WHERE id = $1"#,
&[DecodedValue::Str(id.to_string())],
)
.await?;
Ok(())
}
async fn record_progress_in_tx(tx: &PgTransaction, id: &str, step_index: i64) -> Result<()> {
tx.execute_typed(
r#"INSERT INTO _pylon."Progress" (id, step_index) VALUES ($1, $2)
ON CONFLICT (id) DO UPDATE SET step_index = $2, updated_at = now()"#,
&[DecodedValue::Str(id.to_string()), DecodedValue::I64(step_index)],
)
.await?;
Ok(())
}
async fn delete_progress_in_tx(tx: &PgTransaction, id: &str) -> Result<()> {
tx.execute_typed(
r#"DELETE FROM _pylon."Progress" WHERE id = $1"#,
&[DecodedValue::Str(id.to_string())],
)
.await?;
Ok(())
}
async fn drop_invalid_concurrent_index(pool: &PgPool, sql: &str) -> Result<()> {
let Some(index_name) = concurrent_index_name(sql) else {
return Ok(());
};
let rows = pool
.query_typed(
"SELECT (1) AS result FROM pg_index i JOIN pg_class c ON c.oid = i.indexrelid \
WHERE c.relname = $1 AND NOT i.indisvalid",
&[DecodedValue::Str(index_name.clone())],
pool.types(),
)
.await?;
if !rows.is_empty() {
pool.batch_execute(&format!("DROP INDEX CONCURRENTLY IF EXISTS \"{index_name}\""))
.await?;
}
Ok(())
}
fn concurrent_index_name(sql: &str) -> Option<String> {
let mut tokens = sql.split_whitespace();
let matches_kw = |t: Option<&str>, expected: &str| t.is_some_and(|t| t.eq_ignore_ascii_case(expected));
if !matches_kw(tokens.next(), "CREATE") {
return None;
}
if !matches_kw(tokens.next(), "INDEX") {
return None;
}
if !matches_kw(tokens.next(), "CONCURRENTLY") {
return None;
}
let mut next = tokens.next()?;
if next.eq_ignore_ascii_case("IF") {
if !matches_kw(tokens.next(), "NOT") {
return None;
}
if !matches_kw(tokens.next(), "EXISTS") {
return None;
}
next = tokens.next()?;
}
let after_quote = next.strip_prefix('"').unwrap_or(next);
let name: String = after_quote
.chars()
.take_while(|c| c.is_ascii_alphanumeric() || *c == '_')
.collect();
if name.is_empty() { None } else { Some(name) }
}
const DEV_SAVEPOINT: &str = "pylon_dev";
pub async fn apply_one(pool: &PgPool, m: &MigrationFile, dev_mode: bool) -> Result<()> {
verify_integrity(m)?;
check_no_id_collision(pool, &m.id, &m.onto).await?;
let steps = parse_steps(&m.body);
let resume_from = read_progress(pool, &m.id).await?.map(|i| i + 1).unwrap_or(0) as usize;
let multi_step = steps.len() > 1;
for (step_idx, (transactional, sql)) in steps.iter().enumerate() {
if step_idx < resume_from {
continue;
}
let sql = sql.trim();
if sql.is_empty() {
continue;
}
let is_last = step_idx == steps.len() - 1;
if *transactional {
let tx = pool.begin_default().await?;
let step_result = if dev_mode {
apply_statements_rebasing(&tx, sql).await
} else {
tx.batch_execute(sql).await
};
if let Err(e) = step_result {
let _ = tx.rollback().await;
return Err(e.into());
}
if is_last {
record_applied_in_tx(&tx, &m.id, &m.onto, &m.filename).await?;
if multi_step {
delete_progress_in_tx(&tx, &m.id).await?;
}
} else if multi_step {
record_progress_in_tx(&tx, &m.id, step_idx as i64).await?;
}
tx.commit().await?;
} else {
for statement in crate::migration::split_statements(sql) {
drop_invalid_concurrent_index(pool, &statement).await?;
pool.batch_execute(&statement).await?;
}
if is_last {
record_applied(pool, &m.id, &m.onto, &m.filename).await?;
if multi_step {
delete_progress(pool, &m.id).await?;
}
} else if multi_step {
record_progress(pool, &m.id, step_idx as i64).await?;
}
}
}
Ok(())
}
async fn apply_statements_rebasing(tx: &PgTransaction, sql: &str) -> std::result::Result<(), pylon_pgcon::Error> {
for stmt in crate::migration::split_statements(sql) {
tx.savepoint(DEV_SAVEPOINT).await?;
match tx.batch_execute(&stmt).await {
Ok(()) => tx.release_savepoint(DEV_SAVEPOINT).await?,
Err(e) if is_duplicate_object_error(&e) => tx.rollback_to_savepoint(DEV_SAVEPOINT).await?,
Err(e) => return Err(e),
}
}
Ok(())
}
#[cfg(test)]
mod concurrent_index_name_tests {
use super::concurrent_index_name;
#[test]
fn extracts_a_bare_index_name() {
assert_eq!(
concurrent_index_name("CREATE INDEX CONCURRENTLY idx_person_name ON \"public\".\"Person\" (name);"),
Some("idx_person_name".to_string())
);
}
#[test]
fn extracts_a_quoted_index_name() {
assert_eq!(
concurrent_index_name("CREATE INDEX CONCURRENTLY \"idx_person_name\" ON \"public\".\"Person\" (name);"),
Some("idx_person_name".to_string())
);
}
#[test]
fn handles_if_not_exists() {
assert_eq!(
concurrent_index_name("CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_x ON t (c);"),
Some("idx_x".to_string())
);
}
#[test]
fn is_case_insensitive() {
assert_eq!(
concurrent_index_name("create index concurrently idx_x on t (c);"),
Some("idx_x".to_string())
);
}
#[test]
fn returns_none_for_unrelated_sql() {
assert_eq!(concurrent_index_name("CREATE TABLE foo ();"), None);
assert_eq!(concurrent_index_name("CREATE INDEX idx_x ON t (c);"), None); }
}
#[cfg(test)]
mod tests {
use super::*;
use crate::migration::render_file;
fn test_dsn() -> String {
std::env::var("PYLON_PGCON_TEST_DSN").expect("PYLON_PGCON_TEST_DSN must be set to run live-Postgres tests")
}
async fn test_pool() -> PgPool {
let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
pool.batch_execute("CREATE SCHEMA IF NOT EXISTS _pylon").await.unwrap();
ensure_internal_schema(&pool).await.unwrap();
pool
}
fn make_migration(onto: &str, body: &str) -> MigrationFile {
let content = render_file(onto, body, &[]);
crate::migration::parse(&content, "test").unwrap()
}
fn body(sql: &str) -> String {
format!("\n{sql}\n")
}
fn unique_table_name(prefix: &str) -> String {
use std::time::{SystemTime, UNIX_EPOCH};
let nanos = SystemTime::now().duration_since(UNIX_EPOCH).unwrap().as_nanos();
format!("{prefix}_{nanos}")
}
async fn cleanup_migration_row(pool: &PgPool, id: &str) {
pool.execute_typed(
r#"DELETE FROM _pylon."Migrations" WHERE id = $1"#,
&[DecodedValue::Str(id.to_string())],
)
.await
.unwrap();
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn ensure_internal_schema_is_idempotent() {
let pool = test_pool().await;
ensure_internal_schema(&pool).await.unwrap();
ensure_internal_schema(&pool).await.unwrap();
}
#[test]
fn a_database_this_build_cannot_work_against_is_fatal_and_names_the_fix() {
let state = InternalSchemaState::TooOld { found: 1, required: 2 };
assert!(state.is_fatal());
let message = state.message().unwrap();
assert!(message.contains("pylon migration apply"), "got: {message}");
}
#[test]
fn a_newer_database_is_reported_but_never_fatal() {
let state = InternalSchemaState::Newer { found: 2, current: 1 };
assert!(!state.is_fatal());
assert!(state.message().is_some());
}
#[test]
fn a_pending_upgrade_is_reported_but_not_fatal() {
let state = InternalSchemaState::Behind { found: 1, current: 2 };
assert!(!state.is_fatal());
assert!(state.message().is_some());
}
#[test]
fn an_unmigrated_or_current_database_says_nothing() {
for state in [InternalSchemaState::Unmigrated, InternalSchemaState::Current] {
assert!(!state.is_fatal());
assert_eq!(state.message(), None, "{state:?} should be silent");
}
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn a_freshly_ensured_database_reads_as_current() {
let pool = test_pool().await;
ensure_internal_schema(&pool).await.unwrap();
assert_eq!(
check_internal_schema(&pool).await.unwrap(),
InternalSchemaState::Current
);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn a_database_without_the_marker_table_reads_as_unmigrated() {
let pool = test_pool().await;
ensure_internal_schema(&pool).await.unwrap();
pool.batch_execute(r#"ALTER TABLE _pylon."Internal" RENAME TO "Internal_hidden";"#)
.await
.unwrap();
let state = check_internal_schema(&pool).await;
pool.batch_execute(r#"ALTER TABLE _pylon."Internal_hidden" RENAME TO "Internal";"#)
.await
.unwrap();
assert_eq!(state.unwrap(), InternalSchemaState::Unmigrated);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn ensure_internal_schema_repairs_a_database_missing_a_newer_column() {
let pool = test_pool().await;
ensure_internal_schema(&pool).await.unwrap();
pool.batch_execute(r#"ALTER TABLE _pylon."IndexOutbox" DROP COLUMN IF EXISTS claimed_at;"#)
.await
.unwrap();
ensure_internal_schema(&pool).await.unwrap();
let rows = pool
.query_typed(
"SELECT (count(*)) AS result FROM information_schema.columns \
WHERE table_schema = '_pylon' AND table_name = 'IndexOutbox' \
AND column_name = 'claimed_at'",
&[],
pool.types(),
)
.await
.unwrap();
assert_eq!(
rows.into_iter().next(),
Some(DecodedValue::I64(1)),
"claimed_at should have been restored"
);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn schema_snapshot_round_trips() {
let pool = test_pool().await;
let previous = read_schema_snapshot(&pool).await.unwrap();
write_schema_snapshot(&pool, r#"{"probe": "schema_snapshot_round_trips"}"#)
.await
.unwrap();
let read_back = read_schema_snapshot(&pool).await.unwrap();
assert_eq!(
read_back.as_deref(),
Some(r#"{"probe": "schema_snapshot_round_trips"}"#)
);
write_schema_snapshot(&pool, r#"{"probe": "second_write"}"#)
.await
.unwrap();
let read_back_2 = read_schema_snapshot(&pool).await.unwrap();
assert_eq!(read_back_2.as_deref(), Some(r#"{"probe": "second_write"}"#));
match previous {
Some(prior) => write_schema_snapshot(&pool, &prior).await.unwrap(),
None => pool.batch_execute(r#"DELETE FROM _pylon."Schema""#).await.unwrap(),
}
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn applied_tip_is_none_with_no_applied_rows() {
assert_eq!(applied_tip(&[]), None);
let all_pending = vec![TrackingRow {
id: "m1a".into(),
onto: "initial".into(),
db_state: None,
schema_state: None,
applied: false,
}];
assert_eq!(applied_tip(&all_pending), None);
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn applied_tip_is_the_row_with_no_descendant() {
let tracking = vec![
TrackingRow {
id: "m1a".into(),
onto: "initial".into(),
db_state: None,
schema_state: None,
applied: true,
},
TrackingRow {
id: "m1b".into(),
onto: "m1a".into(),
db_state: None,
schema_state: None,
applied: true,
},
TrackingRow {
id: "m1c".into(),
onto: "m1b".into(),
db_state: None,
schema_state: None,
applied: false,
}, ];
assert_eq!(applied_tip(&tracking), Some("m1b".to_string()));
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn applied_tip_is_deterministic_with_multiple_orphaned_tips() {
let tracking = vec![
TrackingRow {
id: "m1zzz".into(),
onto: "initial".into(),
db_state: None,
schema_state: None,
applied: true,
},
TrackingRow {
id: "m1aaa".into(),
onto: "initial".into(),
db_state: None,
schema_state: None,
applied: true,
},
TrackingRow {
id: "m1mmm".into(),
onto: "initial".into(),
db_state: None,
schema_state: None,
applied: true,
},
];
for _ in 0..20 {
assert_eq!(applied_tip(&tracking), Some("m1aaa".to_string()));
}
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_runs_ddl_and_records_tracking_row() {
let pool = test_pool().await;
let table = unique_table_name("migrate_apply_test");
let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
apply_one(&pool, &m, false).await.unwrap();
let rows = pool
.query_typed(
&format!("SELECT (1) AS result FROM {table}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await;
assert!(rows.is_ok(), "table should exist after apply_one");
let tracking = read_tracking(&pool).await.unwrap();
let row = tracking
.iter()
.find(|r| r.id == m.id)
.expect("tracking row for this migration");
assert!(row.applied);
assert_eq!(row.onto, "initial");
assert_eq!(
row.db_state, None,
"db_state is only ever set separately, by `migration create`'s own UPDATE"
);
assert_eq!(
row.schema_state, None,
"schema_state is only ever set separately, by `apply`'s own UPDATE"
);
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn read_tracking_decodes_schema_state_as_raw_json_text() {
let pool = test_pool().await;
let m = make_migration("initial", &body("SELECT 1;"));
record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
pool.execute_typed(
r#"UPDATE _pylon."Migrations" SET schema_state = $1::jsonb WHERE id = $2"#,
&[
DecodedValue::Str(r#"{"types":[]}"#.to_string()),
DecodedValue::Str(m.id.clone()),
],
)
.await
.unwrap();
let tracking = read_tracking(&pool).await.unwrap();
let row = tracking.iter().find(|r| r.id == m.id).unwrap();
assert_eq!(row.schema_state.as_deref(), Some(r#"{"types": []}"#));
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn read_tracking_decodes_db_state_as_raw_json_text() {
let pool = test_pool().await;
let m = make_migration("initial", &body("SELECT 1;"));
record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
pool.execute_typed(
r#"UPDATE _pylon."Migrations" SET db_state = $1::jsonb WHERE id = $2"#,
&[
DecodedValue::Str(r#"{"schemas":["default"]}"#.to_string()),
DecodedValue::Str(m.id.clone()),
],
)
.await
.unwrap();
let tracking = read_tracking(&pool).await.unwrap();
let row = tracking.iter().find(|r| r.id == m.id).unwrap();
assert_eq!(row.db_state.as_deref(), Some(r#"{"schemas": ["default"]}"#));
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_multi_step_clears_progress_after_completion() {
let pool = test_pool().await;
let t1 = unique_table_name("migrate_step1");
let t2 = unique_table_name("migrate_step2");
let m = make_migration(
"initial",
&format!("\nCREATE TABLE {t1} (id int8);\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
);
apply_one(&pool, &m, false).await.unwrap();
for t in [&t1, &t2] {
let rows = pool
.query_typed(
&format!("SELECT (1) AS result FROM {t}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await;
assert!(rows.is_ok(), "table {t} should exist after apply_one");
}
let progress = pool
.query_typed(
r#"SELECT (1) AS result FROM _pylon."Progress" WHERE id = $1"#,
&[DecodedValue::Str(m.id.clone())],
&pylon_pgcon::ExtensionOids::default(),
)
.await
.unwrap();
assert!(
progress.is_empty(),
"progress row must be cleared after a successful multi-step apply"
);
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_resumes_from_recorded_progress_skipping_earlier_steps() {
let pool = test_pool().await;
let t2 = unique_table_name("migrate_resume_step2");
let m = make_migration(
"initial",
&format!("\nTHIS IS NOT VALID SQL;\n-- pylon:step\nCREATE TABLE {t2} (id int8);\n"),
);
record_progress(&pool, &m.id, 0).await.unwrap();
apply_one(&pool, &m, false).await.unwrap();
let rows = pool
.query_typed(
&format!("SELECT (1) AS result FROM {t2}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await;
assert!(rows.is_ok(), "step 1 should have run");
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_does_not_skip_the_step_it_failed_on_when_resumed() {
let pool = test_pool().await;
let t0 = unique_table_name("migrate_crash_step0");
let t1 = unique_table_name("migrate_crash_step1");
let t2 = unique_table_name("migrate_crash_step2");
let m = make_migration(
"initial",
&format!(
"\nCREATE TABLE {t0} (id int8);\n\
-- pylon:step\n\
INSERT INTO {t1} (id) VALUES (1);\n\
-- pylon:step\n\
CREATE TABLE {t2} (id int8);\n"
),
);
let first = apply_one(&pool, &m, false).await;
assert!(first.is_err(), "step 1 should have failed on the first run");
assert_eq!(
read_progress(&pool, &m.id).await.unwrap(),
Some(0),
"progress must record the last *completed* step, not the one being attempted"
);
pool.batch_execute(&format!("CREATE TABLE {t1} (id int8);"))
.await
.unwrap();
apply_one(&pool, &m, false).await.unwrap();
let rows = pool
.query_typed(
&format!("SELECT (count(*)) AS result FROM {t1}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await
.unwrap();
assert_eq!(
rows.first(),
Some(&DecodedValue::I64(1)),
"step 1 must have been retried on resume, not skipped"
);
let t2_rows = pool
.query_typed(
&format!("SELECT (1) AS result FROM {t2}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await;
assert!(t2_rows.is_ok(), "step 2 should have run after the resumed step 1");
for t in [&t0, &t1, &t2] {
pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
}
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_dev_mode_skips_only_the_duplicate_statement_in_a_step() {
let pool = test_pool().await;
let before = unique_table_name("migrate_rebase_before");
let existing = unique_table_name("migrate_rebase_existing");
let after = unique_table_name("migrate_rebase_after");
pool.batch_execute(&format!("CREATE TABLE {existing} (id int8);"))
.await
.unwrap();
let m = make_migration(
"initial",
&body(&format!(
"CREATE TABLE {before} (id int8);\n\
CREATE TABLE {existing} (id int8);\n\
CREATE TABLE {after} (id int8);"
)),
);
apply_one(&pool, &m, true).await.unwrap();
for t in [&before, &after] {
let rows = pool
.query_typed(
&format!("SELECT (1) AS result FROM {t}"),
&[],
&pylon_pgcon::ExtensionOids::default(),
)
.await;
assert!(
rows.is_ok(),
"table {t} should exist — only the duplicate statement may be skipped"
);
}
let tracking = read_tracking(&pool).await.unwrap();
assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
for t in [&before, &existing, &after] {
pool.batch_execute(&format!("DROP TABLE IF EXISTS {t}")).await.unwrap();
}
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_dev_mode_swallows_a_duplicate_table_error() {
let pool = test_pool().await;
let table = unique_table_name("migrate_dev_mode_test");
pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
.await
.unwrap();
let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
apply_one(&pool, &m, true).await.unwrap();
let tracking = read_tracking(&pool).await.unwrap();
assert!(
tracking.iter().any(|r| r.id == m.id && r.applied),
"still recorded applied despite the swallowed error"
);
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn apply_one_without_dev_mode_propagates_a_duplicate_table_error() {
let pool = test_pool().await;
let table = unique_table_name("migrate_no_dev_mode_test");
pool.batch_execute(&format!("CREATE TABLE {table} (id int8);"))
.await
.unwrap();
let m = make_migration("initial", &body(&format!("CREATE TABLE {table} (id int8);")));
let result = apply_one(&pool, &m, false).await; assert!(result.is_err());
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn record_applied_standalone_marks_a_migration_applied_without_running_ddl() {
let pool = test_pool().await;
let m = make_migration("initial", &body("SELECT 1;"));
record_applied(&pool, &m.id, &m.onto, &m.filename).await.unwrap();
let tracking = read_tracking(&pool).await.unwrap();
assert!(tracking.iter().any(|r| r.id == m.id && r.applied));
cleanup_migration_row(&pool, &m.id).await;
}
#[tokio::test]
#[ignore = "requires a live Postgres via PYLON_PGCON_TEST_DSN"]
async fn advisory_lock_round_trips_and_blocks_a_concurrent_try_lock() {
let pool = PgPool::connect(&test_dsn(), 5).await.unwrap();
let held = advisory_lock(&pool).await.unwrap();
let blocked = try_advisory_lock(&pool).await.unwrap();
assert!(blocked.is_none(), "advisory lock should still be held");
advisory_unlock(held).await.unwrap();
let reacquired = try_advisory_lock(&pool).await.unwrap();
assert!(reacquired.is_some());
advisory_unlock(reacquired.unwrap()).await.unwrap();
}
}