diff --git a/src/api/router/snapshot_session_tests.rs b/src/api/router/snapshot_session_tests.rs index 6c49108d..822890d0 100644 --- a/src/api/router/snapshot_session_tests.rs +++ b/src/api/router/snapshot_session_tests.rs @@ -31,6 +31,54 @@ async fn scalar(db: &C, sql: &str) -> i64 { .unwrap() } +async fn root_metadata( + fixture: &Fixture, +) -> crate::ceres::snapshot::pages::PreparedNativeMetadataRetention { + let mono = fixture.state.storage.mono_storage(); + let head = mono + .read_native_publication_head( + fixture + .state + .storage + .config() + .mst2 + .instance_uuid + .as_deref() + .unwrap(), + ) + .await + .unwrap(); + let handler = fixture + .state + .api_handler(std::path::Path::new("/")) + .await + .unwrap(); + let tree = handler.get_tree_by_hash(&head.root.tree).await.unwrap(); + crate::ceres::snapshot::pages::prepare_native_metadata_retention( + handler.as_ref(), + &tree, + "/", + crate::ceres::snapshot::retention_dag::MetadataDagLimits::default(), + ) + .await + .unwrap() +} + +async fn stored_metadata_ids(db: &DatabaseConnection) -> std::collections::BTreeSet<[u8; 32]> { + db.query_all_raw(Statement::from_string( + DbBackend::Postgres, + "SELECT page_id FROM mst2_metadata_payload", + )) + .await + .unwrap() + .into_iter() + .map(|row| { + let id: Vec = row.try_get_by_index(0).unwrap(); + id.try_into().unwrap() + }) + .collect() +} + async fn rebuilt(fixture: &Fixture) -> MonoApiServiceState { let config = fixture.state.storage.config(); let connection = crate::jupiter::storage::init::postgres_connection(&config.database) @@ -679,6 +727,16 @@ async fn mst2_durable_http_install_fault_never_hands_off_or_returns_success() { let fixture = Fixture::new_with_pg_config(true).await; let mono = fixture.state.storage.mono_storage(); let db = mono.get_connection(); + let prepared = root_metadata(&fixture).await; + let existing = stored_metadata_ids(db).await; + let missing: Vec<_> = prepared + .dag() + .payloads() + .iter() + .filter(|page| !existing.contains(&page.id)) + .map(|page| page.id) + .collect(); + assert_eq!(missing, [prepared.dag().root()]); db.execute_unprepared( "CREATE FUNCTION reject_http_page() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'forced HTTP install interruption'; END $$; @@ -691,7 +749,7 @@ async fn mst2_durable_http_install_fault_never_hands_off_or_returns_success() { fixture .app .clone() - .oneshot(resolve_request("/project/nested")) + .oneshot(resolve_request("/")) .await .unwrap(), 503, @@ -699,6 +757,7 @@ async fn mst2_durable_http_install_fault_never_hands_off_or_returns_success() { true, ) .await; + assert_eq!(stored_metadata_ids(db).await, existing); assert_eq!( scalar(db, "SELECT count(*) FROM mst2_snapshot_context").await, 1 @@ -739,12 +798,21 @@ async fn mst2_durable_http_install_fault_never_hands_off_or_returns_success() { fixture .app .clone() - .oneshot(resolve_request("/project/nested")) + .oneshot(resolve_request("/")) .await .unwrap(), ) .await; - assert_eq!(success["descriptor"]["scope"], "/project/nested"); + assert_eq!(success["descriptor"]["scope"], "/"); + assert_eq!( + stored_metadata_ids(db).await, + prepared + .dag() + .payloads() + .iter() + .map(|page| page.id) + .collect() + ); assert_eq!( scalar( db, @@ -869,14 +937,35 @@ async fn mst2_durable_http_large_dag_batches_handoff_without_scanning_pages_or_e let fixture = Fixture::new_with_pg_config_and_directories(true, 80).await; let mono = fixture.state.storage.mono_storage(); let db = mono.get_connection(); - assert!(scalar(db, "SELECT count(*) FROM mst2_metadata_payload").await > 64); + let projected = root_metadata(&fixture).await; + let existing = stored_metadata_ids(db).await; + assert!(existing.len() > 64); + let missing: Vec<_> = projected + .dag() + .payloads() + .iter() + .filter(|page| !existing.contains(&page.id)) + .map(|page| page.id) + .collect(); + assert_eq!(missing, [projected.dag().root()]); + let missing_batches = projected + .dag() + .payloads() + .chunks(64) + .filter(|batch| batch.iter().any(|page| !existing.contains(&page.id))) + .count() as i64; + assert_eq!(missing_batches, 1); db.execute_unprepared( - "CREATE TABLE http_install_batch_count(singleton integer PRIMARY KEY,batches bigint NOT NULL); - INSERT INTO http_install_batch_count VALUES(1,0); + "CREATE TABLE http_install_batch_count(singleton integer PRIMARY KEY,batches bigint NOT NULL,rows bigint NOT NULL); + INSERT INTO http_install_batch_count VALUES(1,0,0); CREATE FUNCTION count_http_install_batch() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN UPDATE http_install_batch_count SET batches=batches+1 WHERE singleton=1; RETURN NULL; END $$; CREATE TRIGGER count_http_install_batch AFTER INSERT ON mst2_metadata_payload - FOR EACH STATEMENT EXECUTE FUNCTION count_http_install_batch()", + FOR EACH STATEMENT EXECUTE FUNCTION count_http_install_batch(); + CREATE FUNCTION count_http_install_row() RETURNS trigger LANGUAGE plpgsql AS $$ + BEGIN UPDATE http_install_batch_count SET rows=rows+1 WHERE singleton=1; RETURN NULL; END $$; + CREATE TRIGGER count_http_install_row AFTER INSERT ON mst2_metadata_payload + FOR EACH ROW EXECUTE FUNCTION count_http_install_row()", ).await.unwrap(); let prepared = Arc::new(Barrier::new(2)); let release = Arc::new(Barrier::new(2)); @@ -898,9 +987,23 @@ async fn mst2_durable_http_large_dag_batches_handoff_without_scanning_pages_or_e ) .await; assert!(members > 64); + assert_eq!(members as usize, projected.dag().payloads().len()); + assert_eq!( + scalar(db, "SELECT rows FROM http_install_batch_count").await, + missing.len() as i64 + ); + assert_eq!( + stored_metadata_ids(db).await, + projected + .dag() + .payloads() + .iter() + .map(|page| page.id) + .collect() + ); assert_eq!( scalar(db, "SELECT batches FROM http_install_batch_count").await, - (members + 63) / 64 + missing_batches ); let writer = db.begin().await.unwrap(); assert!( @@ -931,7 +1034,7 @@ async fn mst2_durable_http_large_dag_batches_handoff_without_scanning_pages_or_e assert_eq!(resolved["descriptor"]["scope"], "/"); assert_eq!( scalar(db, "SELECT batches FROM http_install_batch_count").await, - (members + 63) / 64 + missing_batches ); fixture.counts.assert(0, 0); let warm = success_json( @@ -949,7 +1052,7 @@ async fn mst2_durable_http_large_dag_batches_handoff_without_scanning_pages_or_e ); assert_eq!( scalar(db, "SELECT batches FROM http_install_batch_count").await, - (members + 63) / 64 + missing_batches ); advance(&fixture).await; let state = rebuilt(&fixture).await; diff --git a/src/jupiter/service/native_publication_push_tests.rs b/src/jupiter/service/native_publication_push_tests.rs index 4fb89aa1..ce3febbb 100644 --- a/src/jupiter/service/native_publication_push_tests.rs +++ b/src/jupiter/service/native_publication_push_tests.rs @@ -106,7 +106,10 @@ async fn maintenance_cannot_reinitialize_published_history_after_head_loss() { mono.get_connection().execute_unprepared("DELETE FROM mst2_native_head").await.unwrap(); for instance in [NATIVE_INSTANCE, "11111111-2222-4333-8444-555555555555"] { let error = mono.initialize_native_publication_for_maintenance(instance, &head.root).await.unwrap_err(); - assert!(error.to_string().contains("native publication history exists")); + assert!(error.to_string().contains("native publication history exists"), "rejected initialization for {instance}: {error}"); + let released = mono.get_connection().begin().await.unwrap(); + assert!(PushQueueStorage::try_mono_write_lock(&released).await.unwrap(), "rejected initialization must release the mono write lock before returning"); + released.rollback().await.unwrap(); } assert_eq!(mst2_native_head::Entity::find().count(mono.get_connection()).await.unwrap(), 0); assert_eq!(mst2_native_publication::Entity::find().all(mono.get_connection()).await.unwrap(), certificates); @@ -116,7 +119,10 @@ async fn maintenance_cannot_reinitialize_published_history_after_head_loss() { mono.get_connection().execute_unprepared("DELETE FROM mst2_native_publication").await.unwrap(); for instance in [NATIVE_INSTANCE, "11111111-2222-4333-8444-555555555555"] { let error = mono.initialize_native_publication_for_maintenance(instance, &head.root).await.unwrap_err(); - assert!(error.to_string().contains("native publication history exists")); + assert!(error.to_string().contains("native publication history exists"), "rejected initialization for {instance}: {error}"); + let released = mono.get_connection().begin().await.unwrap(); + assert!(PushQueueStorage::try_mono_write_lock(&released).await.unwrap(), "rejected initialization must release the mono write lock before returning"); + released.rollback().await.unwrap(); } assert_eq!(mst2_native_head::Entity::find().count(mono.get_connection()).await.unwrap(), 0); assert_eq!(mst2_native_publication::Entity::find().count(mono.get_connection()).await.unwrap(), 0); diff --git a/src/jupiter/storage/native_publication_storage.rs b/src/jupiter/storage/native_publication_storage.rs index 2dbff900..979c7314 100644 --- a/src/jupiter/storage/native_publication_storage.rs +++ b/src/jupiter/storage/native_publication_storage.rs @@ -351,6 +351,7 @@ impl MonoStorage { .await? .ok_or_else(|| integrity("native history observation missing"))?; if history.try_get::("", "present")? { + txn.rollback().await?; return Err(integrity( "native publication history exists; initialization cannot repair or replace it", ));