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 Ok(())
633}
634
635async 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); }
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 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 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 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 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 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 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 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 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 }, ];
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 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 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 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 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 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 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 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 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 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 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 pool.batch_execute(&format!("CREATE TABLE {existing} (id int8);"))
1205 .await
1206 .unwrap();
1207
1208 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 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 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(); 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; 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 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 let reacquired = try_advisory_lock(&pool).await.unwrap();
1311 assert!(reacquired.is_some());
1312 advisory_unlock(reacquired.unwrap()).await.unwrap();
1313 }
1314}