diff --git a/src/api/router/snapshot_persisted_metadata_tests.rs b/src/api/router/snapshot_persisted_metadata_tests.rs new file mode 100644 index 00000000..028f7d35 --- /dev/null +++ b/src/api/router/snapshot_persisted_metadata_tests.rs @@ -0,0 +1,870 @@ +use std::collections::BTreeSet; + +use mst2_codec::metapage::{Page, page_id}; + +use super::*; +use crate::jupiter::storage::native_snapshot_session::MetadataRouteRequest; + +fn metadata_body(items: Value, encoding: &str) -> Vec { + serde_json::to_vec_pretty(&json!({"encoding": encoding, "items": items})).unwrap() +} + +async fn metadata_bytes(fixture: &Fixture, app: Router, body: &[u8]) -> Vec { + let response = app + .oneshot(fixture.request("POST", "metadata/pages", Body::from(body.to_vec()))) + .await + .unwrap(); + assert_eq!(response.status(), 200); + assert_eq!(response.headers()["x-mega-snapshot-id"], fixture.snapshot); + to_bytes(response.into_body(), usize::MAX) + .await + .unwrap() + .to_vec() +} + +fn assert_metadata(bytes: &[u8], body: &[u8], count: u32, expected: &[([u8; 32], Vec)]) { + let frames = parse_stream(bytes).unwrap(); + let mut actual = Vec::new(); + for frame in &frames[..frames.len() - 1] { + let Frame::Meta(meta) = frame else { + panic!("metadata route emitted a non-META frame") + }; + assert!(meta.pages.len() <= mst2_codec::treeframe::META_MAX_PAGES); + assert!( + meta.pages + .iter() + .map(|(_, page)| 36 + page.len()) + .sum::() + <= mst2_codec::treeframe::META_MAX_RAW + ); + for (id, page) in &meta.pages { + assert_eq!(*id, page_id(page)); + Page::decode(page).unwrap(); + } + actual.extend_from_slice(&meta.pages); + } + assert_eq!(actual, expected); + let Frame::End(end) = frames.last().unwrap() else { + panic!("metadata route omitted END") + }; + assert_eq!(end.request_item_count, count); + assert_eq!(end.unique_unit_count, expected.len() as u32); + assert_eq!( + end.logical_bytes, + expected + .iter() + .map(|(_, page)| page.len() as u64) + .sum::() + ); + assert_eq!( + end.request_body_sha256, + <[u8; 32]>::from(Sha256::digest(body)) + ); +} + +#[tokio::test] +async fn mst2_persisted_meta_matches_canonical_routes_after_rebuild_and_advance_without_git_reads() +{ + let fixture = Fixture::new_with_pg_config_and_directories(true, 140).await; + let context = fixture + .state + .storage + .snapshot_context(&fixture.snapshot, &fixture.lease) + .await + .unwrap(); + let handler = fixture + .state + .api_handler(std::path::Path::new("/")) + .await + .unwrap(); + let root_tree = handler + .get_tree_by_hash(&context.root_tree_oid) + .await + .unwrap(); + let built = crate::ceres::snapshot::pages::build_directory_page( + handler.as_ref(), + &root_tree, + "/project", + ) + .await + .unwrap(); + let (Page::Branch { children, .. }, _) = Page::decode(&built.page_bytes).unwrap() else { + panic!("wide fixture must use canonical radix pages") + }; + let label = children + .iter() + .find(|child| child.label == b'w') + .unwrap() + .label; + let specs = [ + ("/", vec![]), + ("/", vec![label]), + ("/wide-139", vec![]), + ("/nested", vec![]), + ("/directory", vec![]), + ("/", vec![label]), + ]; + let mut expected = Vec::new(); + let mut seen = BTreeSet::new(); + let mut items = Vec::new(); + for (path, route) in &specs { + let absolute = if *path == "/" { + "/project".to_owned() + } else { + format!("/project{path}") + }; + let directory = crate::ceres::snapshot::pages::build_directory_page( + handler.as_ref(), + &root_tree, + &absolute, + ) + .await + .unwrap(); + let pages = Page::pages_along_route(&directory.codec_entries, route).unwrap(); + let reached = page_id(pages.last().unwrap()); + items.push(json!({"directory_path": path, "route": route, + "expected_digest": format!("sha256:{}", hex::encode(reached))})); + for page in pages { + let id = page_id(&page); + if seen.insert(id) { + expected.push((id, page)); + } + } + } + advance(&fixture).await; + let state = rebuilt(&fixture).await; + assert!(state.storage.native_snapshot_sessions.get().is_none()); + let mono = fixture.state.storage.mono_storage(); + let held = mono.get_connection().begin().await.unwrap(); + held.execute_unprepared("LOCK TABLE mega_tree,mst2_verified_object IN ACCESS EXCLUSIVE MODE") + .await + .unwrap(); + fixture.counts.reset(); + for encoding in ["identity", "zstd"] { + let body = metadata_body(json!(items), encoding); + let bytes = tokio::time::timeout( + Duration::from_secs(4), + metadata_bytes(&fixture, app(&state), &body), + ) + .await + .expect("persisted META must not touch locked Git tree or verified-object tables"); + assert_metadata(&bytes, &body, specs.len() as u32, &expected); + } + let requests: Vec<_> = specs + .iter() + .map(|(path, route)| MetadataRouteRequest { + directory_path: path, + route, + expected_digest: None, + }) + .collect(); + let batch = state + .storage + .snapshot_sessions() + .await + .metadata_routes(&context, &requests) + .await + .unwrap(); + assert_eq!(batch.pages, expected); + assert_eq!(batch.work.page_queries, batch.work.pages_loaded); + assert!( + batch.work.pages_loaded < 20, + "route work must not scan the 140-directory DAG" + ); + assert!( + batch.work.walk_visits > batch.work.pages_loaded, + "duplicate/path visits reuse request pages" + ); + assert!( + batch.work.payload_bytes >= batch.pages.iter().map(|(_, page)| page.len() as u64).sum() + ); + held.rollback().await.unwrap(); + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_persisted_meta_absence_scope_digest_limits_and_release_oracles() { + let fixture = Fixture::new_with_pg_config(true).await; + for (items, status, code) in [ + ( + json!([{"directory_path":"/missing"}]), + 404, + "PATH_NOT_FOUND", + ), + (json!([{"directory_path":"/file"}]), 409, "NOT_DIRECTORY"), + ( + json!([{"directory_path":"/link/child"}]), + 409, + "NOT_DIRECTORY", + ), + ( + json!([{"directory_path":"/nested", "route":[0]}]), + 404, + "PATH_NOT_FOUND", + ), + ( + json!([{"directory_path":"/../outside"}]), + 400, + "SCOPE_INVALID", + ), + ( + json!([{"directory_path":format!("/{}", vec!["a"; 256].join("/"))}]), + 400, + "SCOPE_INVALID", + ), + ( + json!([{"directory_path":"/outside"}]), + 404, + "PATH_NOT_FOUND", + ), + ( + json!([{"directory_path":"/", "expected_digest":format!("sha256:{}", "0".repeat(64))}]), + 409, + "EXPECTED_DIGEST_MISMATCH", + ), + (json!([]), 413, "LIMIT_EXCEEDED"), + ( + json!(vec![json!({"directory_path":"/"}); 65]), + 413, + "LIMIT_EXCEEDED", + ), + ] { + error( + fixture + .send( + "POST", + "metadata/pages", + Body::from(metadata_body(items, "identity")), + ) + .await, + status, + code, + false, + ) + .await; + } + let body = metadata_body(json!([{"directory_path":"/"}]), "identity"); + let response = fixture + .send("POST", "metadata/pages", Body::from(body)) + .await; + assert_eq!(response.status(), 200); + let mut stream = response.into_body().into_data_stream(); + let first = stream.next().await.unwrap().unwrap(); + let (Frame::Meta(meta), consumed) = mst2_codec::treeframe::parse_frame(&first).unwrap() else { + panic!("first persisted frame must be META") + }; + assert_eq!(consumed, first.len()); + assert!(!meta.pages.is_empty()); + lease_control(&fixture, &fixture.lease, "DELETE", false).await; + assert!(stream.next().await.unwrap().is_err()); + assert!(stream.next().await.is_none()); + assert!(matches!( + parse_stream(&first), + Err(mst2_codec::CodecError::BadOrdering( + "stream missing END/ERROR frame" + )) + )); + fixture.counts.assert(0, 0); +} + +async fn page_row(fixture: &Fixture, path: &str) -> ([u8; 32], Vec) { + let body = metadata_body(json!([{"directory_path":path}]), "identity"); + let bytes = metadata_bytes(fixture, fixture.app.clone(), &body).await; + let Frame::Meta(meta) = &parse_stream(&bytes).unwrap()[0] else { + panic!("META required") + }; + meta.pages[0].clone() +} + +async fn damage_page(fixture: &Fixture, id: [u8; 32], remove: bool) { + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let modes = handoff_trigger_modes_for_test(db).await; + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SELECT pg_advisory_xact_lock(1296717362,hashtext(current_schema()))") + .await + .unwrap(); + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload DISABLE TRIGGER mst2_metadata_payload_fenced", + ) + .await + .unwrap(); + if remove { + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload DISABLE TRIGGER mst2_metadata_payload_removed", + ) + .await + .unwrap(); + } + let sql = if remove { + "DELETE FROM mst2_metadata_payload WHERE page_id=$1" + } else { + "UPDATE mst2_metadata_payload SET payload=set_byte(payload,0,0) WHERE page_id=$1" + }; + assert_eq!( + txn.execute_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + [id.to_vec().into()] + )) + .await + .unwrap() + .rows_affected(), + 1 + ); + if remove { + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload ENABLE TRIGGER mst2_metadata_payload_removed", + ) + .await + .unwrap(); + } + txn.execute_unprepared( + "ALTER TABLE mst2_metadata_payload ENABLE TRIGGER mst2_metadata_payload_fenced", + ) + .await + .unwrap(); + txn.commit().await.unwrap(); + assert_eq!(handoff_trigger_modes_for_test(db).await, modes); + assert!( + db.execute_unprepared("UPDATE mst2_metadata_payload SET payload=payload") + .await + .is_err() + ); +} + +#[tokio::test] +async fn mst2_persisted_meta_warm_missing_and_corrupt_pages_never_reproject_from_git() { + for remove in [true, false] { + let fixture = Fixture::new_with_pg_config(true).await; + let (id, _) = page_row(&fixture, "/nested").await; + damage_page(&fixture, id, remove).await; + let body = metadata_body(json!([{"directory_path":"/nested"}]), "identity"); + let mono = fixture.state.storage.mono_storage(); + let held = mono.get_connection().begin().await.unwrap(); + held.execute_unprepared( + "LOCK TABLE mega_tree,mst2_verified_object IN ACCESS EXCLUSIVE MODE", + ) + .await + .unwrap(); + let response = tokio::time::timeout( + Duration::from_secs(4), + fixture.send("POST", "metadata/pages", Body::from(body)), + ) + .await + .expect("damaged persisted page must fail before any source read"); + error( + response, + if remove { 503 } else { 502 }, + if remove { + "OBJECT_UNAVAILABLE" + } else { + "INTEGRITY_ERROR" + }, + false, + ) + .await; + held.rollback().await.unwrap(); + fixture.counts.assert(0, 0); + } +} + +async fn wait_retention_waiter(txn: &sea_orm::DatabaseTransaction) { + tokio::time::timeout(Duration::from_secs(4), async { + loop { + txn.execute_unprepared("SELECT pg_stat_clear_snapshot()") + .await + .unwrap(); + if scalar( + txn, + "SELECT count(*) FROM pg_locks l JOIN pg_stat_activity a ON a.pid=l.pid + WHERE l.locktype='advisory' AND NOT l.granted AND a.datname=current_database() + AND a.application_name=current_schema() AND l.classid=1296717362::oid + AND l.objid=hashtext(current_schema())::oid AND l.objsubid=2", + ) + .await + > 0 + { + break; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("META request did not wait for retention authority"); +} + +#[tokio::test] +async fn mst2_persisted_meta_rechecks_deadline_and_release_after_waiting_for_retention_lock() { + for expire in [true, false] { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let held = mono.get_connection().begin().await.unwrap(); + held.execute_unprepared( + "SELECT pg_advisory_xact_lock(1296717362,hashtext(current_schema()))", + ) + .await + .unwrap(); + let pending = { + let application = fixture.app.clone(); + let request = fixture.request( + "POST", + "metadata/pages", + Body::from(metadata_body( + json!([{"directory_path":"/nested"}]), + "identity", + )), + ); + tokio::spawn(async move { application.oneshot(request).await.unwrap() }) + }; + wait_retention_waiter(&held).await; + let sql = if expire { + "UPDATE mst2_snapshot_lease SET expires_at_unix=0 WHERE lease_id=$1" + } else { + "UPDATE mst2_snapshot_lease SET state='RELEASED' WHERE lease_id=$1" + }; + held.execute_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + [fixture.lease.clone().into()], + )) + .await + .unwrap(); + held.commit().await.unwrap(); + error(pending.await.unwrap(), 410, "LEASE_EXPIRED", false).await; + fixture.counts.assert(0, 0); + } +} + +async fn bound_fixture() -> ( + Fixture, + crate::ceres::snapshot::retention_dag::MetadataPagePayload, +) { + let mut fixture = Fixture::new_with_pg_config(true).await; + advance(&fixture).await; + 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(); + let prepared = crate::ceres::snapshot::pages::prepare_native_metadata_retention( + handler.as_ref(), + &tree, + "/project", + crate::ceres::snapshot::retention_dag::MetadataDagLimits::default(), + ) + .await + .unwrap(); + assert_eq!(prepared.dag().payloads().len(), 1); + let expected = prepared.dag().payloads()[0].clone(); + let repository = crate::jupiter::storage::native_metadata_install::generations::PostgresMetadataGenerationRepository::new( + mono.get_connection().clone()).await.unwrap(); + let intent = repository + .begin_intent("persisted-meta-bound-seed", &prepared) + .await + .unwrap(); + repository + .install_pages(&intent, prepared.dag().payloads()) + .await + .unwrap(); + repository.finalize(&intent).await.unwrap(); + let resolved = success_json( + fixture + .app + .clone() + .oneshot(resolve_request("/project")) + .await + .unwrap(), + ) + .await; + assert_ne!(resolved["descriptor"]["snapshot_id"], fixture.snapshot); + fixture.snapshot = resolved["descriptor"]["snapshot_id"] + .as_str() + .unwrap() + .to_owned(); + fixture.lease = resolved["lease_id"].as_str().unwrap().to_owned(); + let db = mono.get_connection(); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM mst2_metadata_payload WHERE generation=1" + ) + .await, + 1 + ); + db.execute_raw(Statement::from_sql_and_values(DbBackend::Postgres, + "INSERT INTO mst2_metadata_lifetime(page_id,node_id,generation,state,metadata_codec,expected_size,graph_domain) + VALUES($1,$2,2,'RESERVED',1,$3,'generic-v1')", + [expected.id.to_vec().into(), format!("page:sha256:{}", hex::encode(expected.id)).into(), + (expected.size as i32).into()])).await.unwrap(); + let membership = db.query_one_raw(Statement::from_sql_and_values(DbBackend::Postgres, + "SELECT pp.generation AS member_generation,c.generation AS current_generation + FROM mst2_snapshot_context s JOIN mst2_metadata_prepare_page pp ON pp.prepare_id=s.prepare_id + JOIN mst2_metadata_current c ON c.page_id=pp.page_id WHERE s.snapshot_id=$1", + [fixture.snapshot.clone().into()])).await.unwrap().unwrap(); + assert_eq!( + membership + .try_get::>("", "member_generation") + .unwrap(), + None + ); + assert_eq!( + membership.try_get::("", "current_generation").unwrap(), + 1 + ); + fixture.counts.reset(); + (fixture, expected) +} + +#[tokio::test] +async fn mst2_persisted_meta_serves_bound_generic_pages_with_null_members_and_extra_history() { + let (fixture, expected) = bound_fixture().await; + let body = metadata_body(json!([{"directory_path":"/"}]), "identity"); + for application in [fixture.app.clone(), app(&rebuilt(&fixture).await)] { + let bytes = metadata_bytes(&fixture, application, &body).await; + assert_metadata(&bytes, &body, 1, &[(expected.id, expected.bytes.clone())]); + } + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_persisted_meta_rejects_touched_graph_damage_and_prepare_membership_loss() { + for case in 0..5 { + let fixture = Fixture::new_with_pg_config(true).await; + let (id, _) = page_row(&fixture, "/nested").await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let modes = handoff_trigger_modes_for_test(db).await; + let txn = db.begin().await.unwrap(); + txn.execute_unprepared( + "SELECT pg_advisory_xact_lock(1296717362,hashtext(current_schema()))", + ) + .await + .unwrap(); + let node = format!("page:sha256:{}", hex::encode(id)); + let (sql, values) = match case { + 0 => ("UPDATE mst2_retention_node SET state='DELETING' WHERE node_id=$1", vec![node.into()]), + 1 => ("UPDATE mst2_retention_node SET bytes=bytes+1 WHERE node_id=$1", vec![node.into()]), + 2 => ("DELETE FROM mst2_retention_edge WHERE child_id=$1", vec![node.into()]), + 3 => ("INSERT INTO mst2_retention_gc_op(operation_id,node_id,operation,state,attempts,created_at) + VALUES('persisted-meta-tombstone',$1,'REMOVE','PENDING',0,now())", vec![node.into()]), + 4 => { + txn.execute_unprepared("ALTER TABLE mst2_metadata_prepare_page DISABLE TRIGGER mst2_install_capability_mapping_guard").await.unwrap(); + ("DELETE FROM mst2_metadata_prepare_page WHERE page_id=$1 AND prepare_id= + (SELECT prepare_id FROM mst2_snapshot_context WHERE snapshot_id=$2)", + vec![id.to_vec().into(), fixture.snapshot.clone().into()]) + } + _ => unreachable!(), + }; + assert!( + txn.execute_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + values + )) + .await + .unwrap() + .rows_affected() + > 0 + ); + if case == 4 { + txn.execute_unprepared("ALTER TABLE mst2_metadata_prepare_page ENABLE TRIGGER mst2_install_capability_mapping_guard").await.unwrap(); + } + txn.commit().await.unwrap(); + assert_eq!(handoff_trigger_modes_for_test(db).await, modes); + let response = fixture + .send( + "POST", + "metadata/pages", + Body::from(metadata_body( + json!([{"directory_path":"/nested"}]), + "identity", + )), + ) + .await; + error( + response, + if case == 2 { 502 } else { 503 }, + if case == 2 { + "INTEGRITY_ERROR" + } else { + "OBJECT_UNAVAILABLE" + }, + false, + ) + .await; + fixture.counts.assert(0, 0); + } +} + +#[tokio::test] +async fn mst2_persisted_meta_reader_holds_protection_until_all_route_bytes_are_owned() { + use crate::jupiter::storage::native_snapshot_session::with_metadata_read_barriers; + + for release_lease in [false, true] { + let fixture = Fixture::new_with_pg_config(true).await; + let expected = vec![ + page_row(&fixture, "/nested").await, + page_row(&fixture, "/directory").await, + ]; + let context = fixture + .state + .storage + .snapshot_context(&fixture.snapshot, &fixture.lease) + .await + .unwrap(); + let captured = Arc::new(Barrier::new(2)); + let resume = Arc::new(Barrier::new(2)); + let reading = { + let state = fixture.state.clone(); + let context = context.clone(); + let captured = captured.clone(); + let resume = resume.clone(); + tokio::spawn(async move { + let requests = [ + MetadataRouteRequest { + directory_path: "/nested", + route: &[], + expected_digest: None, + }, + MetadataRouteRequest { + directory_path: "/directory", + route: &[], + expected_digest: None, + }, + ]; + with_metadata_read_barriers( + captured, + resume, + state + .storage + .snapshot_sessions() + .await + .metadata_routes(&context, &requests), + ) + .await + }) + }; + tokio::time::timeout(Duration::from_secs(4), captured.wait()) + .await + .unwrap(); + let root = format!("page:{}", context.built.metadata_root); + let competing = { + let state = fixture.state.clone(); + let lease = fixture.lease.clone(); + let root = root.clone(); + tokio::spawn(async move { + if release_lease { + assert!(state.storage.snapshot_release(&lease).await.unwrap()); + None + } else { + Some( + PostgresRetentionRepository::new( + state.storage.mono_storage().get_connection().clone(), + ) + .mark_deleting("persisted-meta-reader-race", &root) + .await + .unwrap(), + ) + } + }) + }; + let observer = fixture + .state + .storage + .mono_storage() + .get_connection() + .begin() + .await + .unwrap(); + wait_retention_waiter(&observer).await; + assert!(!reading.is_finished()); + assert!(!competing.is_finished()); + resume.wait().await; + let batch = reading.await.unwrap().unwrap(); + assert_eq!(batch.pages, expected); + assert_eq!( + batch.work.pages_loaded, 3, + "root and both directories are read before unlock" + ); + let outcome = competing.await.unwrap(); + observer.rollback().await.unwrap(); + if release_lease { + assert!(outcome.is_none()); + } else { + assert_eq!(outcome, Some(GcClaim::Unavailable)); + lease_control(&fixture, &fixture.lease, "DELETE", false).await; + } + let gc = PostgresRetentionRepository::new( + fixture + .state + .storage + .mono_storage() + .get_connection() + .clone(), + ); + assert_eq!( + gc.mark_deleting("persisted-meta-after-read", &root) + .await + .unwrap(), + GcClaim::Marked + ); + gc.complete_gc("persisted-meta-after-read").await.unwrap(); + let requests = [MetadataRouteRequest { + directory_path: "/nested", + route: &[], + expected_digest: None, + }]; + let rejected = fixture + .state + .storage + .snapshot_sessions() + .await + .metadata_routes(&context, &requests) + .await + .err() + .expect("released and collected root must fail before touching pages"); + assert_eq!(rejected.code, SnapshotErrorCode::LeaseExpired); + assert_eq!( + batch.pages, expected, + "collected graph does not invalidate owned route bytes" + ); + fixture.counts.assert(0, 0); + } +} + +async fn fault_binding(fixture: &Fixture, id: [u8; 32], case: usize, restore: bool) { + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let modes = handoff_trigger_modes_for_test(db).await; + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SELECT pg_advisory_xact_lock(1296717362,hashtext(current_schema()))") + .await + .unwrap(); + let (table, guard, sql) = match case { + 0 => ( + "mst2_metadata_payload", + "mst2_metadata_payload_fenced", + if restore { + "UPDATE mst2_metadata_payload SET generation=1 WHERE page_id=$1" + } else { + "UPDATE mst2_metadata_payload SET generation=NULL WHERE page_id=$1" + }, + ), + 1 => ( + "mst2_metadata_current", + "mst2_metadata_current_guard", + if restore { + "UPDATE mst2_metadata_current SET generation=1 WHERE page_id=$1" + } else { + "UPDATE mst2_metadata_current SET generation=2 WHERE page_id=$1" + }, + ), + 2 => ( + "mst2_metadata_lifetime", + "mst2_metadata_lifetime_guard", + if restore { + "UPDATE mst2_metadata_lifetime SET state='LIVE' WHERE page_id=$1 AND generation=1" + } else { + "UPDATE mst2_metadata_lifetime SET state='DELETING' WHERE page_id=$1 AND generation=1" + }, + ), + 3 => ( + "mst2_metadata_lifetime", + "mst2_metadata_lifetime_guard", + if restore { + "UPDATE mst2_metadata_lifetime SET graph_domain='generic-v1' WHERE page_id=$1 AND generation=1" + } else { + "UPDATE mst2_metadata_lifetime SET graph_domain='qualified-v1' WHERE page_id=$1 AND generation=1" + }, + ), + 4 => ( + "mst2_metadata_lifetime", + "mst2_metadata_lifetime_guard", + if restore { + "UPDATE mst2_metadata_lifetime SET expected_size=expected_size-1 WHERE page_id=$1 AND generation=1" + } else { + "UPDATE mst2_metadata_lifetime SET expected_size=expected_size+1 WHERE page_id=$1 AND generation=1" + }, + ), + 5 => ( + "mst2_metadata_prepare_page", + "mst2_install_capability_mapping_guard", + if restore { + "UPDATE mst2_metadata_prepare_page SET generation=NULL WHERE page_id=$1 AND prepare_id= + (SELECT prepare_id FROM mst2_snapshot_context WHERE snapshot_id=$2)" + } else { + "UPDATE mst2_metadata_prepare_page SET generation=2 WHERE page_id=$1 AND prepare_id= + (SELECT prepare_id FROM mst2_snapshot_context WHERE snapshot_id=$2)" + }, + ), + _ => unreachable!(), + }; + txn.execute_unprepared(&format!("ALTER TABLE {table} DISABLE TRIGGER {guard}")) + .await + .unwrap(); + let values = if case == 5 { + vec![id.to_vec().into(), fixture.snapshot.clone().into()] + } else { + vec![id.to_vec().into()] + }; + assert_eq!( + txn.execute_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + values + )) + .await + .unwrap() + .rows_affected(), + 1 + ); + txn.execute_unprepared(&format!("ALTER TABLE {table} ENABLE TRIGGER {guard}")) + .await + .unwrap(); + txn.commit().await.unwrap(); + assert_eq!(handoff_trigger_modes_for_test(db).await, modes); +} + +#[tokio::test] +async fn mst2_persisted_meta_warm_reads_require_exact_generic_current_lifetime_and_membership() { + let (fixture, expected) = bound_fixture().await; + let body = metadata_body(json!([{"directory_path":"/"}]), "identity"); + let bytes = metadata_bytes(&fixture, fixture.app.clone(), &body).await; + assert_metadata(&bytes, &body, 1, &[(expected.id, expected.bytes.clone())]); + for case in 0..6 { + fault_binding(&fixture, expected.id, case, false).await; + error( + fixture + .send("POST", "metadata/pages", Body::from(body.clone())) + .await, + if matches!(case, 0 | 1 | 5) { 502 } else { 503 }, + if matches!(case, 0 | 1 | 5) { + "INTEGRITY_ERROR" + } else { + "OBJECT_UNAVAILABLE" + }, + false, + ) + .await; + fault_binding(&fixture, expected.id, case, true).await; + let bytes = metadata_bytes(&fixture, fixture.app.clone(), &body).await; + assert_metadata(&bytes, &body, 1, &[(expected.id, expected.bytes.clone())]); + } + fixture.counts.assert(0, 0); +} diff --git a/src/api/router/snapshot_router.rs b/src/api/router/snapshot_router.rs index eed58484..f75b7279 100644 --- a/src/api/router/snapshot_router.rs +++ b/src/api/router/snapshot_router.rs @@ -1476,62 +1476,90 @@ async fn metadata_pages( } let ctx = request_context(&state, &snapshot_id).map_err(mst2_error_response)?; - let handler = state - .api_handler(std::path::Path::new("/")) - .await - .map_err(internal)?; - let root_tree = handler - .get_tree_by_hash(&ctx.root_tree_oid) - .await - .map_err(internal)?; - let scope = &ctx.built.descriptor.scope; - - // Unique pages across all items, in first-seen order. - let mut unique: Vec<([u8; 32], Vec)> = Vec::new(); - let mut seen: Vec<[u8; 32]> = Vec::new(); - let mut logical_bytes: u64 = 0; - for item in &req.items { - let abs_path = abs_view_path(scope, &item.directory_path); - let built = build_directory_page(handler.as_ref(), &root_tree, &abs_path) + let unique = if state.storage.config().mst2.publication_enabled { + use crate::jupiter::storage::native_snapshot_session::MetadataRouteRequest; + let items: Vec<_> = req + .items + .iter() + .map(|item| MetadataRouteRequest { + directory_path: &item.directory_path, + route: &item.route, + expected_digest: item.expected_digest.as_deref(), + }) + .collect(); + let batch = state + .storage + .snapshot_sessions() + .await + .metadata_routes(&ctx, &items) .await .map_err(mst2_error_response)?; - let pages = - mst2_codec::metapage::Page::pages_along_route(&built.codec_entries, &item.route) - .map_err(|e| match e { - // A label the fixed view does not have is proven absence. - mst2_codec::CodecError::BadOrdering(m) => SnapshotError::new( - SnapshotErrorCode::PathNotFound, - format!("{}: route does not resolve ({m})", item.directory_path), - ), - other => SnapshotError::new( - SnapshotErrorCode::Internal, - format!("{}: route walk failed ({other})", item.directory_path), + tracing::debug!( + page_queries = batch.work.page_queries, + pages_loaded = batch.work.pages_loaded, + payload_bytes = batch.work.payload_bytes, + walk_visits = batch.work.walk_visits, + edge_references_checked = batch.work.edge_references_checked, + "served persisted generic metadata routes" + ); + batch.pages + } else { + let handler = state + .api_handler(std::path::Path::new("/")) + .await + .map_err(internal)?; + let root_tree = handler + .get_tree_by_hash(&ctx.root_tree_oid) + .await + .map_err(internal)?; + let scope = &ctx.built.descriptor.scope; + let mut unique: Vec<([u8; 32], Vec)> = Vec::new(); + let mut seen: Vec<[u8; 32]> = Vec::new(); + for item in &req.items { + let abs_path = abs_view_path(scope, &item.directory_path); + let built = build_directory_page(handler.as_ref(), &root_tree, &abs_path) + .await + .map_err(mst2_error_response)?; + let pages = + mst2_codec::metapage::Page::pages_along_route(&built.codec_entries, &item.route) + .map_err(|e| match e { + // A label the fixed view does not have is proven absence. + mst2_codec::CodecError::BadOrdering(m) => SnapshotError::new( + SnapshotErrorCode::PathNotFound, + format!("{}: route does not resolve ({m})", item.directory_path), + ), + other => SnapshotError::new( + SnapshotErrorCode::Internal, + format!("{}: route walk failed ({other})", item.directory_path), + ), + })?; + let last = pages + .last() + .expect("pages_along_route returns at least the root page"); + let last_id = mst2_codec::metapage::page_id(last); + if let Some(expected) = &item.expected_digest + && expected != &format!("sha256:{}", hex_of(&last_id)) + { + return Err(mst2_error_response(SnapshotError::new( + SnapshotErrorCode::DigestMismatch, + format!( + "{}: route does not reach expected_digest", + item.directory_path ), - })?; - let last = pages - .last() - .expect("pages_along_route returns at least the root page"); - let last_id = mst2_codec::metapage::page_id(last); - if let Some(expected) = &item.expected_digest - && expected != &format!("sha256:{}", hex_of(&last_id)) - { - return Err(mst2_error_response(SnapshotError::new( - SnapshotErrorCode::DigestMismatch, - format!( - "{}: route does not reach expected_digest", - item.directory_path - ), - ))); - } - for page in pages { - let id = mst2_codec::metapage::page_id(&page); - if !seen.contains(&id) { - logical_bytes += page.len() as u64; - seen.push(id); - unique.push((id, page)); + ))); + } + for page in pages { + let id = mst2_codec::metapage::page_id(&page); + if !seen.contains(&id) { + seen.push(id); + unique.push((id, page)); + } } } - } + unique + }; + let page_count = unique.len(); + let logical_bytes = unique.iter().map(|(_, page)| page.len() as u64).sum(); // Frames hold at most 64 pages and at most 1 MiB of raw payload (spec 06), // so a wide route set becomes several META frames rather than one @@ -1578,7 +1606,7 @@ async fn metadata_pages( let request_body_sha256: [u8; 32] = sha2::Digest::finalize(hasher).into(); let end = stream.end( req.items.len() as u32, - u32::try_from(seen.len()).unwrap_or(u32::MAX), + u32::try_from(page_count).unwrap_or(u32::MAX), logical_bytes, request_body_sha256, ); diff --git a/src/api/router/snapshot_session_tests.rs b/src/api/router/snapshot_session_tests.rs index 17dfcb21..6c49108d 100644 --- a/src/api/router/snapshot_session_tests.rs +++ b/src/api/router/snapshot_session_tests.rs @@ -1213,3 +1213,6 @@ async fn mst2_durable_http_warm_reads_renew_and_frame_delivery_reject_primary_sc #[path = "snapshot_lookup_metadata_tests.rs"] mod metadata_lookup; + +#[path = "snapshot_persisted_metadata_tests.rs"] +mod persisted_metadata; diff --git a/src/jupiter/storage/native_metadata_install.rs b/src/jupiter/storage/native_metadata_install.rs index aeddd315..5b97f663 100644 --- a/src/jupiter/storage/native_metadata_install.rs +++ b/src/jupiter/storage/native_metadata_install.rs @@ -643,6 +643,13 @@ impl PostgresMetadataInstallRepository { .map_err(internal) } + pub(super) async fn metadata_read_barrier( + &self, + txn: &DatabaseTransaction, + ) -> Result<(), SnapshotError> { + self.capability_barrier(txn).await + } + async fn barrier(&self, txn: &DatabaseTransaction) -> Result<(), SnapshotError> { if txn.get_database_backend() != DbBackend::Postgres { return Err(internal("metadata installation requires PostgreSQL")); diff --git a/src/jupiter/storage/native_snapshot_metadata_routes.rs b/src/jupiter/storage/native_snapshot_metadata_routes.rs new file mode 100644 index 00000000..8326fca6 --- /dev/null +++ b/src/jupiter/storage/native_snapshot_metadata_routes.rs @@ -0,0 +1,492 @@ +//! Fixed-root persisted META reads under the original generic lease authority. + +use std::collections::{BTreeSet, HashMap}; + +use mst2_codec::metapage::{HEADER_LEN, PAGE_MAX_BYTES, Page, page_id}; +use sea_orm::{ConnectionTrait, DatabaseTransaction, QueryResult}; + +use super::{ + PostgresNativeSessionRepository, SnapshotContext, SnapshotError, SnapshotErrorCode, + decode_context, expired, finish, integrity, internal, routes, statement, +}; +use crate::ceres::snapshot::{ + retention_dag::MetadataDagLimits, view::validate_scope_relative_path, +}; + +pub(crate) struct MetadataRouteRequest<'a> { + pub directory_path: &'a str, + pub route: &'a [u8], + pub expected_digest: Option<&'a str>, +} + +#[derive(Debug, Default, Clone, Copy)] +pub(crate) struct PersistedMetadataReadWork { + pub page_queries: u64, + pub pages_loaded: u64, + pub payload_bytes: u64, + pub walk_visits: u64, + pub edge_references_checked: u64, +} + +pub(crate) struct PersistedMetadataRouteBatch { + pub pages: Vec<([u8; 32], Vec)>, + pub work: PersistedMetadataReadWork, +} + +#[cfg(test)] +tokio::task_local! { + static METADATA_READ_BARRIERS: (std::sync::Arc, std::sync::Arc); +} + +#[cfg(test)] +pub(crate) async fn with_metadata_read_barriers( + captured: std::sync::Arc, + release: std::sync::Arc, + future: F, +) -> F::Output { + METADATA_READ_BARRIERS + .scope((captured, release), future) + .await +} + +impl PostgresNativeSessionRepository { + pub(crate) async fn metadata_routes( + &self, + context: &SnapshotContext, + requests: &[MetadataRouteRequest<'_>], + ) -> Result { + if requests.is_empty() || requests.len() > 64 { + return Err(limit("items must hold 1..64 entries")); + } + for request in requests { + validate_scope_relative_path(request.directory_path)?; + let scope = &context.built.descriptor.scope; + let absolute = if scope == "/" { + request.directory_path.to_owned() + } else if request.directory_path == "/" { + scope.clone() + } else { + format!("{scope}{}", request.directory_path) + }; + validate_scope_relative_path(&absolute)?; + } + let installer = self.installer().await?; + let txn = self.transaction().await?; + let result = async { + installer.verify_primary_connection(&txn).await?; + if routes::lease(&txn, installer, &context.lease_id) + .await? + .as_deref() + != Some(context.built.snapshot_id.as_str()) + { + return Err(expired()); + } + // This barrier captures the namespace and uses bounded lock acquisition. + // Do not acquire the publication/route writer lock after retention. + installer.metadata_read_barrier(&txn).await?; + if routes::lease(&txn, installer, &context.lease_id) + .await? + .as_deref() + != Some(context.built.snapshot_id.as_str()) + { + return Err(expired()); + } + let row = txn + .query_one_raw(statement( + self.session_sql(installer).await, + [ + context.built.snapshot_id.clone().into(), + context.lease_id.clone().into(), + ], + )) + .await + .map_err(internal)? + .ok_or_else(expired)?; + installer.verify_primary_scope_row(&row)?; + let current = decode_context( + &row, + &context.built.snapshot_id, + &context.lease_id, + &context.built.instance_id, + )?; + if current.built.descriptor != context.built.descriptor + || current.commit_oid != context.commit_oid + || current.root_tree_oid != context.root_tree_oid + || current.authorization_epoch != context.authorization_epoch + { + return Err(integrity("fixed session changed during metadata read")); + } + let prepare_id: String = row.try_get("", "prepare_id").map_err(internal)?; + let mut reader = Reader { + txn: &txn, + prepare_id, + codec: context.built.descriptor.metadata_codec() as i16, + cache: HashMap::new(), + work: PersistedMetadataReadWork::default(), + limits: MetadataDagLimits::default(), + }; + let mut seen = BTreeSet::new(); + let mut pages = Vec::new(); + for request in requests { + let root = reader + .directory( + context.built.descriptor.metadata_root, + request.directory_path, + ) + .await?; + let route = reader.route(root, request.route).await?; + let reached = route + .last() + .ok_or_else(|| integrity("metadata route is empty"))?; + if let Some(expected) = request.expected_digest + && expected != format!("sha256:{}", hex::encode(reached)) + { + return Err(SnapshotError::new( + SnapshotErrorCode::DigestMismatch, + format!( + "{}: route does not reach expected_digest", + request.directory_path + ), + )); + } + for id in route { + if seen.insert(id) { + let page = reader + .cache + .get(&id) + .ok_or_else(|| integrity("read metadata page is missing"))?; + pages.push((id, page.bytes.clone())); + } + } + } + Ok(PersistedMetadataRouteBatch { + pages, + work: reader.work, + }) + } + .await; + // All bytes are owned before releasing protection. Frame delivery keeps + // the existing per-frame authentication and lease revalidation. + finish(txn, result).await + } +} + +struct StoredPage { + bytes: Vec, + page: Page, + entries: u64, +} + +struct Reader<'a> { + txn: &'a DatabaseTransaction, + prepare_id: String, + codec: i16, + cache: HashMap<[u8; 32], StoredPage>, + work: PersistedMetadataReadWork, + limits: MetadataDagLimits, +} + +impl Reader<'_> { + async fn load(&mut self, id: [u8; 32]) -> Result<(), SnapshotError> { + self.work.walk_visits += 1; + if self.work.walk_visits > self.limits.prepare_entry_visits as u64 { + return Err(limit("metadata route work budget exceeded")); + } + if self.cache.contains_key(&id) { + return Ok(()); + } + if self.cache.len() >= self.limits.nodes { + return Err(limit("metadata route page budget exceeded")); + } + self.work.page_queries += 1; + let row = self + .txn + .query_one_raw(statement( + PAGE_SQL, + [self.prepare_id.clone().into(), id.to_vec().into()], + )) + .await + .map_err(internal)? + .ok_or_else(|| unavailable("metadata page is not a prepared member"))?; + let page = decode_page(&row, id, self.codec)?; + self.work.payload_bytes = self + .work + .payload_bytes + .checked_add(page.bytes.len() as u64) + .filter(|bytes| *bytes <= self.limits.payload_bytes) + .ok_or_else(|| limit("metadata route payload budget exceeded"))?; + let expected = page_references(&page.page); + self.work.edge_references_checked += expected.len() as u64; + if self.work.edge_references_checked > self.limits.edges as u64 { + return Err(limit("metadata route edge budget exceeded")); + } + let edges: String = row.try_get("", "outgoing_edges").map_err(internal)?; + let edges: Vec = serde_json::from_str(&edges).map_err(internal)?; + let actual: BTreeSet<_> = edges.iter().cloned().collect(); + if edges.len() != actual.len() || actual != expected { + return Err(integrity( + "metadata graph edges disagree with fixed page references", + )); + } + self.work.pages_loaded += 1; + self.cache.insert(id, page); + #[cfg(test)] + if self.cache.len() == 1 + && let Ok((captured, release)) = METADATA_READ_BARRIERS.try_with(Clone::clone) + { + captured.wait().await; + release.wait().await; + } + Ok(()) + } + + async fn descend(&mut self, parent: [u8; 32], label: u8) -> Result<[u8; 32], SnapshotError> { + self.load(parent).await?; + let (prefix, child) = match &self.cache[&parent].page { + Page::Branch { + prefix, children, .. + } => { + let child = children + .iter() + .find(|child| child.label == label) + .ok_or_else(|| absent("route label is absent in fixed metadata page"))?; + (prefix.clone(), child.clone()) + } + Page::Leaf { .. } => return Err(absent("route descends past a leaf page")), + }; + let id = child.child_page_id; + self.load(id).await?; + let received = &self.cache[&id]; + let mut partition = prefix; + partition.push(label); + let valid_prefix = match &received.page { + Page::Leaf { entries } => entries + .iter() + .all(|entry| entry.name.starts_with(&partition)), + Page::Branch { prefix, .. } => prefix.starts_with(&partition), + }; + if received.entries != child.subtree_entries || !valid_prefix { + return Err(integrity( + "metadata child differs from fixed radix partition", + )); + } + Ok(id) + } + + async fn directory( + &mut self, + mut root: [u8; 32], + path: &str, + ) -> Result<[u8; 32], SnapshotError> { + if path == "/" { + return Ok(root); + } + for name in path[1..].split('/') { + let mut id = root; + loop { + self.load(id).await?; + let entry = match &self.cache[&id].page { + Page::Leaf { entries } => entries + .iter() + .find(|entry| entry.name == name.as_bytes()) + .cloned(), + Page::Branch { + prefix, terminal, .. + } => { + if name.as_bytes() == prefix { + terminal.clone() + } else { + if !name.as_bytes().starts_with(prefix) { + return Err(absent("name is absent in fixed directory")); + } + let label = name + .as_bytes() + .get(prefix.len()) + .copied() + .ok_or_else(|| absent("name is absent in fixed directory"))?; + id = self.descend(id, label).await?; + continue; + } + } + } + .ok_or_else(|| absent("name is absent in fixed directory"))?; + if !entry.is_dir() { + return Err(SnapshotError::new( + SnapshotErrorCode::NotDirectory, + format!("{path} is not a directory"), + )); + } + root = entry.child_root; + break; + } + } + Ok(root) + } + + async fn route( + &mut self, + root: [u8; 32], + labels: &[u8], + ) -> Result, SnapshotError> { + self.load(root).await?; + let mut pages = vec![root]; + let mut current = root; + for label in labels { + current = self.descend(current, *label).await?; + pages.push(current); + } + Ok(pages) + } +} + +const PAGE_SQL: &str = "SELECT m.expected_size,m.generation AS member_generation,p.graph_domain AS prepare_graph_domain, + b.metadata_codec AS payload_codec,b.byte_size AS payload_size,octet_length(b.payload) AS actual_size, + CASE WHEN octet_length(b.payload) BETWEEN 20 AND 16384 THEN b.payload END AS payload, + b.generation AS payload_generation,c.generation AS current_generation, + l.generation AS lifetime_generation,l.graph_domain,l.state AS lifetime_state, + l.metadata_codec AS lifetime_codec,l.expected_size AS lifetime_size, + n.state AS graph_state,n.kind AS graph_kind,n.bytes AS graph_bytes, + EXISTS(SELECT 1 FROM mst2_retention_gc_op g WHERE g.node_id='page:sha256:'||encode(m.page_id,'hex') + AND g.operation='REMOVE' AND g.state IN ('PENDING','APPLIED')) AS tombstone, + COALESCE((SELECT jsonb_agg(e.child_id) FROM + (SELECT child_id FROM mst2_retention_edge WHERE parent_id='page:sha256:'||encode(m.page_id,'hex') + ORDER BY child_id LIMIT 258) e),'[]'::jsonb)::text AS outgoing_edges + FROM mst2_metadata_prepare_page m + JOIN mst2_metadata_prepare p ON p.prepare_id=m.prepare_id + LEFT JOIN mst2_metadata_payload b ON b.page_id=m.page_id + LEFT JOIN mst2_metadata_current c ON c.page_id=m.page_id + LEFT JOIN mst2_metadata_lifetime l ON l.page_id=c.page_id AND l.generation=c.generation + LEFT JOIN mst2_retention_node n ON n.node_id='page:sha256:'||encode(m.page_id,'hex') + WHERE m.prepare_id=$1 AND m.page_id=$2"; + +fn decode_page(row: &QueryResult, id: [u8; 32], codec: i16) -> Result { + if row + .try_get::>("", "prepare_graph_domain") + .map_err(internal)? + .as_deref() + .is_some_and(|domain| domain != "generic-v1") + { + return Err(unavailable( + "fixed metadata preparation is outside the generic namespace", + )); + } + let size: i32 = row.try_get("", "expected_size").map_err(internal)?; + if !(HEADER_LEN..=PAGE_MAX_BYTES).contains(&(size as usize)) + || row + .try_get::>("", "payload_codec") + .map_err(internal)? + != Some(codec) + || row + .try_get::>("", "payload_size") + .map_err(internal)? + != Some(size) + || row + .try_get::>("", "actual_size") + .map_err(internal)? + != Some(size) + { + return Err(unavailable( + "fixed metadata payload profile or size is unavailable", + )); + } + if row.try_get::("", "tombstone").map_err(internal)? + || row + .try_get::>("", "graph_state") + .map_err(internal)? + .as_deref() + != Some("LIVE") + || row + .try_get::>("", "graph_kind") + .map_err(internal)? + .as_deref() + != Some("page") + || row + .try_get::>("", "graph_bytes") + .map_err(internal)? + != Some(size as i64) + { + return Err(unavailable("fixed generic metadata graph is unavailable")); + } + let member: Option = row.try_get("", "member_generation").map_err(internal)?; + let payload: Option = row.try_get("", "payload_generation").map_err(internal)?; + let current: Option = row.try_get("", "current_generation").map_err(internal)?; + let lifetime: Option = row.try_get("", "lifetime_generation").map_err(internal)?; + if payload != current || payload != lifetime || member.is_some() && member != payload { + return Err(integrity( + "fixed generic page differs from its exact physical lifetime", + )); + } + if let Some(generation) = payload + && (generation <= 0 + || row + .try_get::>("", "graph_domain") + .map_err(internal)? + .as_deref() + != Some("generic-v1") + || row + .try_get::>("", "lifetime_state") + .map_err(internal)? + .as_deref() + != Some("LIVE") + || row + .try_get::>("", "lifetime_codec") + .map_err(internal)? + != Some(codec) + || row + .try_get::>("", "lifetime_size") + .map_err(internal)? + != Some(size)) + { + return Err(unavailable( + "fixed metadata lifetime cannot be served in the generic namespace", + )); + } + let bytes: Vec = row + .try_get::>>("", "payload") + .map_err(internal)? + .ok_or_else(|| unavailable("fixed metadata payload is missing"))?; + if page_id(&bytes) != id { + return Err(integrity("fixed metadata payload digest mismatch")); + } + let (page, entries) = Page::decode(&bytes) + .map_err(|_| integrity("fixed metadata payload is not a canonical page"))?; + Ok(StoredPage { + bytes, + page, + entries, + }) +} + +fn page_references(page: &Page) -> BTreeSet { + let mut edges = BTreeSet::new(); + let entries = match page { + Page::Leaf { entries } => entries.as_slice(), + Page::Branch { + terminal, children, .. + } => { + edges.extend( + children + .iter() + .map(|child| format!("page:sha256:{}", hex::encode(child.child_page_id))), + ); + terminal.as_slice() + } + }; + edges.extend( + entries + .iter() + .filter(|entry| entry.is_dir()) + .map(|entry| format!("page:sha256:{}", hex::encode(entry.child_root))), + ); + edges +} + +fn absent(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::PathNotFound, message) +} +fn unavailable(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::ObjectUnavailable, message) +} +fn limit(message: &str) -> SnapshotError { + SnapshotError::new(SnapshotErrorCode::LimitExceeded, message) +} diff --git a/src/jupiter/storage/native_snapshot_session.rs b/src/jupiter/storage/native_snapshot_session.rs index f4abe782..fb8d9134 100644 --- a/src/jupiter/storage/native_snapshot_session.rs +++ b/src/jupiter/storage/native_snapshot_session.rs @@ -363,6 +363,13 @@ impl PostgresNativeSessionRepository { #[path = "native_snapshot_routes.rs"] mod routes; +#[path = "native_snapshot_metadata_routes.rs"] +mod metadata_routes; + +pub(crate) use metadata_routes::MetadataRouteRequest; +#[cfg(test)] +pub(crate) use metadata_routes::with_metadata_read_barriers; + const SESSION_SQL: &str = "SELECT s.snapshot_id,s.canonical_descriptor,s.commit_oid,s.root_tree_oid, (SELECT storage_uuid FROM mst2_metadata_storage_scope WHERE singleton=1) AS authority_storage_uuid, current_database() AS authority_database,