1use 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
52pub 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
70pub async fn ensure_internal_schema(pool: &PgPool) -> Result<()> {
88 pool.batch_execute(&crate::stdlib::export_stdlib()).await?;
89 Ok(())
90}
91
92#[derive(Debug, Clone, Copy, PartialEq, Eq)]
99pub enum InternalSchemaState {
100 Unmigrated,
104 TooOld { found: i32, required: i32 },
107 Behind { found: i32, current: i32 },
110 Current,
112 Newer { found: i32, current: i32 },
116}
117
118impl InternalSchemaState {
119 pub fn is_fatal(&self) -> bool {
121 matches!(self, InternalSchemaState::TooOld { .. })
122 }
123
124 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
146pub 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 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
170pub 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
193pub 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
211pub 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#[derive(Debug, Clone, PartialEq)]
233pub struct TrackingRow {
234 pub id: String,
235 pub onto: String,
236 pub db_state: Option<String>,
241 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
298pub 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
326pub 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
341pub 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
366const 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
383async 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
419pub 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
470async 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
492async 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
515fn 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
551pub async fn apply_one(pool: &PgPool, m: &MigrationFile, dev_mode: bool) -> Result<()> {
569 verify_integrity(m)?;
570 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 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 pool.refresh_types().await?;
637
638 Ok(())
639}
640
641async 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); }
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 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 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 #[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 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 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 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 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 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 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 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 }, ];
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 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 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 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 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 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 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 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 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 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 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 pool.batch_execute(&format!("CREATE TABLE {existing} (id int8);"))
1259 .await
1260 .unwrap();
1261
1262 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 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 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(); 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; 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 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 let reacquired = try_advisory_lock(&pool).await.unwrap();
1365 assert!(reacquired.is_some());
1366 advisory_unlock(reacquired.unwrap()).await.unwrap();
1367 }
1368}