From b77341643e3525c85ec3268ca095d24c47939571 Mon Sep 17 00:00:00 2001 From: Xiaoyang Han Date: Wed, 7 Oct 2026 23:53:02 +0800 Subject: [PATCH] feat(mst2): bind generic sessions to immutable storage routes Derive permanent semantic SID routes, separate generic physical bindings and immutable original lease routes on actual v3 open/read/renew/release paths. Serialize context/lease insert before row locks via captured core mono, route and namespace retention barriers; backfill complete inventory with fixed-scope proof. Authenticate generic prepare domain at fixed-core registration, preserve original domain immutability and old upgrade reads, and add ten actual PostgreSQL/HTTP regressions. Qualified serving and collector remain closed; native and performance evidence remain separate. --- src/api/router/snapshot_content_tests.rs | 6 + .../router/snapshot_storage_route_fixture.rs | 165 +++ .../router/snapshot_storage_route_tests.rs | 1183 +++++++++++++++++ ...20261007_000600_add_mst2_storage_routes.rs | 42 + .../m20261007_000600_storage_routes.sql | 329 +++++ src/jupiter/migration/mod.rs | 16 +- .../storage/native_metadata_install.rs | 4 + src/jupiter/storage/native_snapshot_routes.rs | 122 ++ .../storage/native_snapshot_session.rs | 44 +- 9 files changed, 1900 insertions(+), 11 deletions(-) create mode 100644 src/api/router/snapshot_storage_route_fixture.rs create mode 100644 src/api/router/snapshot_storage_route_tests.rs create mode 100644 src/jupiter/migration/m20261007_000600_add_mst2_storage_routes.rs create mode 100644 src/jupiter/migration/m20261007_000600_storage_routes.sql create mode 100644 src/jupiter/storage/native_snapshot_routes.rs diff --git a/src/api/router/snapshot_content_tests.rs b/src/api/router/snapshot_content_tests.rs index 1e0dcd55..add40d69 100644 --- a/src/api/router/snapshot_content_tests.rs +++ b/src/api/router/snapshot_content_tests.rs @@ -1346,3 +1346,9 @@ async fn mst2_fixed_metadata_walker_rejects_real_tree_gitlinks_without_body_read )); fixture.counts.assert(0, 0); } + +#[path = "snapshot_storage_route_tests.rs"] +mod storage_routes; + +#[path = "snapshot_storage_route_fixture.rs"] +mod storage_route_fixture; diff --git a/src/api/router/snapshot_storage_route_fixture.rs b/src/api/router/snapshot_storage_route_fixture.rs new file mode 100644 index 00000000..751502b1 --- /dev/null +++ b/src/api/router/snapshot_storage_route_fixture.rs @@ -0,0 +1,165 @@ +use sea_orm::{ConnectionTrait, DatabaseConnection, DbBackend, Statement, TransactionTrait}; + +async fn trigger_modes(db: &C) -> String { + db.query_one_raw(Statement::from_string( + DbBackend::Postgres, + "SELECT coalesce(string_agg(c.relname||':'||t.tgname||':'||t.tgenabled::text,',' + ORDER BY c.relname,t.tgname),'') AS modes FROM pg_trigger t + JOIN pg_class c ON c.oid=t.tgrelid JOIN pg_namespace n ON n.oid=c.relnamespace + WHERE n.nspname=current_schema() AND NOT t.tgisinternal", + )) + .await + .unwrap() + .unwrap() + .try_get_by_index(0) + .unwrap() +} + +pub(super) async fn restore_pre_route_schema(db: &DatabaseConnection) { + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SELECT mst2_route_enter(current_schema())") + .await + .unwrap(); + let untouched = txn.query_one_raw(Statement::from_string(DbBackend::Postgres, + "SELECT coalesce(string_agg(c.relname||':'||t.tgname||':'||t.tgenabled::text,',' ORDER BY c.relname,t.tgname),'') + FROM pg_trigger t JOIN pg_class c ON c.oid=t.tgrelid + JOIN pg_namespace n ON n.oid=c.relnamespace JOIN pg_proc p ON p.oid=t.tgfoid + WHERE n.nspname=current_schema() AND NOT t.tgisinternal AND left(p.proname,11)<>'mst2_route_'", + )).await.unwrap().unwrap().try_get_by_index::(0).unwrap(); + // Remove only this additive layer in an isolated deployment-upgrade fixture. + // The original contexts, leases, plans, payloads and protection stay intact. + txn.execute_unprepared( + "DO $$ DECLARE r record; functions text; BEGIN + FOR r IN SELECT c.relname,t.tgname FROM pg_trigger t + JOIN pg_class c ON c.oid=t.tgrelid JOIN pg_namespace n ON n.oid=c.relnamespace + JOIN pg_proc p ON p.oid=t.tgfoid + WHERE n.nspname=current_schema() AND NOT t.tgisinternal AND left(p.proname,11)='mst2_route_' + LOOP EXECUTE format('DROP TRIGGER %I ON %I.%I',r.tgname,current_schema(),r.relname); END LOOP; + DROP TABLE mst2_lease_storage_route,mst2_generic_session_storage_binding,mst2_snapshot_storage_route,mst2_metadata_namespace; + FOR r IN SELECT conname FROM pg_constraint + WHERE conrelid='mst2_snapshot_context'::regclass AND contype='u' + AND pg_get_constraintdef(oid)='UNIQUE (snapshot_id, prepare_id, metadata_root)' + LOOP EXECUTE format('ALTER TABLE mst2_snapshot_context DROP CONSTRAINT %I',r.conname); END LOOP; + SELECT string_agg(format('%I.%I(%s)',n.nspname,p.proname,pg_get_function_identity_arguments(p.oid)),',') + INTO functions FROM pg_proc p JOIN pg_namespace n ON n.oid=p.pronamespace + WHERE n.nspname=current_schema() AND left(p.proname,11)='mst2_route_'; + IF functions IS NULL THEN RAISE EXCEPTION 'route upgrade fixture has no route functions'; END IF; + EXECUTE 'DROP FUNCTION '||functions; + END $$; + DELETE FROM seaql_migrations WHERE version='m20261007_000600_add_mst2_storage_routes'", + ) + .await + .unwrap(); + assert_eq!(trigger_modes(&txn).await, untouched); + txn.commit().await.unwrap(); +} + +pub(super) async fn reject_half_lease_commit_for_test( + db: &DatabaseConnection, + lease_id: &str, + insert_sql: &str, +) { + let before = trigger_modes(db).await; + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SELECT mst2_route_enter(current_schema())") + .await + .unwrap(); + // Simulate an omitted derived write; all statement, identity and deferred + // completeness guards remain active, and the deriving trigger is restored + // before the actual commit attempt. + txn.execute_unprepared( + "ALTER TABLE mst2_snapshot_lease DISABLE TRIGGER mst2_route_lease_insert", + ) + .await + .unwrap(); + assert_eq!( + txn.execute_unprepared(insert_sql) + .await + .unwrap() + .rows_affected(), + 1 + ); + let half = txn + .query_one_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + "SELECT (SELECT count(*) FROM mst2_snapshot_lease WHERE lease_id=$1)::bigint AS actual, + (SELECT count(*) FROM mst2_lease_storage_route WHERE lease_id=$1)::bigint AS routes", + [lease_id.into()], + )) + .await + .unwrap() + .unwrap(); + assert_eq!(half.try_get::("", "actual").unwrap(), 1); + assert_eq!(half.try_get::("", "routes").unwrap(), 0); + txn.execute_unprepared( + "ALTER TABLE mst2_snapshot_lease ENABLE TRIGGER mst2_route_lease_insert", + ) + .await + .unwrap(); + assert_eq!(trigger_modes(&txn).await, before); + let rejected = txn.commit().await.unwrap_err(); + assert!( + rejected + .to_string() + .contains("storage route lease committed without exact routing"), + "{rejected}" + ); + assert_eq!(trigger_modes(db).await, before); +} + +pub(super) async fn overwrite_lease_route_incarnation_for_test( + db: &DatabaseConnection, + lease_id: &str, +) { + let sql = "UPDATE mst2_lease_storage_route SET session_incarnation=$2::uuid WHERE lease_id=$1"; + let incarnation = uuid::Uuid::new_v4().to_string(); + let mutation = || { + Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + [lease_id.into(), incarnation.clone().into()], + ) + }; + assert!( + db.execute_raw(mutation()) + .await + .unwrap_err() + .to_string() + .contains("immutable") + ); + let before = trigger_modes(db).await; + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SELECT mst2_route_enter(current_schema())") + .await + .unwrap(); + txn.execute_unprepared( + "ALTER TABLE mst2_lease_storage_route DISABLE TRIGGER mst2_route_immutable", + ) + .await + .unwrap(); + assert_eq!( + txn.execute_raw(mutation()).await.unwrap().rows_affected(), + 1 + ); + txn.execute_unprepared( + "ALTER TABLE mst2_lease_storage_route ENABLE TRIGGER mst2_route_immutable", + ) + .await + .unwrap(); + assert_eq!(trigger_modes(&txn).await, before); + txn.commit().await.unwrap(); + assert_eq!(trigger_modes(db).await, before); + let forbidden = uuid::Uuid::new_v4().to_string(); + assert_ne!(forbidden, incarnation); + assert!( + db.execute_raw(Statement::from_sql_and_values( + DbBackend::Postgres, + sql, + [lease_id.into(), forbidden.into()], + )) + .await + .unwrap_err() + .to_string() + .contains("immutable") + ); +} diff --git a/src/api/router/snapshot_storage_route_tests.rs b/src/api/router/snapshot_storage_route_tests.rs new file mode 100644 index 00000000..bd8df6f9 --- /dev/null +++ b/src/api/router/snapshot_storage_route_tests.rs @@ -0,0 +1,1183 @@ +use mst2_codec::{ + descriptor::ServingDescriptor, + metapage::{Entry, EntryKind, Page, page_id}, +}; +use sea_orm::{ + ConnectionTrait, DatabaseConnection, DatabaseTransaction, DbBackend, IsolationLevel, Statement, + TransactionTrait, +}; +use sea_orm_migration::MigratorTrait; + +use super::*; +use crate::{ + ceres::snapshot::{ + pages::PreparedNativeMetadataRetention, + retention::RetentionRoot, + retention_dag::{MetadataDagBuilder, MetadataDagLimits}, + }, + jupiter::{ + migration::Migrator, + storage::{ + mst2_retention::{PostgresRetentionRepository, RETENTION_LOCK_KEY}, + native_metadata_install::generations::qualified::PostgresQualifiedMetadataRepository, + push_queue_storage::MONO_WRITE_LOCK_KEY1, + }, + }, +}; + +const ROUTE_LOCK_KEY: i32 = 1_296_718_001; +const ROUTE_TABLES: [&str; 4] = [ + "mst2_metadata_namespace", + "mst2_snapshot_storage_route", + "mst2_generic_session_storage_binding", + "mst2_lease_storage_route", +]; + +fn statement(sql: &str, values: impl IntoIterator) -> Statement { + Statement::from_sql_and_values(DbBackend::Postgres, sql, values) +} + +async fn scalar(db: &C, sql: &str) -> i64 { + db.query_one_raw(Statement::from_string(DbBackend::Postgres, sql)) + .await + .unwrap() + .unwrap() + .try_get_by_index(0) + .unwrap() +} + +async fn json_sql(db: &C, sql: &str) -> Value { + let value: String = db + .query_one_raw(Statement::from_string(DbBackend::Postgres, sql)) + .await + .unwrap() + .unwrap() + .try_get_by_index(0) + .unwrap(); + serde_json::from_str(&value).unwrap() +} + +async fn routes(db: &C) -> Value { + json_sql( + db, + "SELECT jsonb_build_object( + 'namespace',(SELECT jsonb_agg(to_jsonb(n) ORDER BY namespace_uuid) FROM mst2_metadata_namespace n), + 'snapshot',(SELECT jsonb_agg(to_jsonb(r) ORDER BY snapshot_id) FROM mst2_snapshot_storage_route r), + 'binding',(SELECT jsonb_agg(to_jsonb(b) ORDER BY snapshot_id) FROM mst2_generic_session_storage_binding b), + 'lease',(SELECT jsonb_agg(to_jsonb(l) ORDER BY lease_id) FROM mst2_lease_storage_route l))::text", + ) + .await +} + +async fn sources(db: &C) -> Value { + json_sql( + db, + "SELECT jsonb_build_object( + 'context',(SELECT jsonb_agg(to_jsonb(s) ORDER BY snapshot_id) FROM mst2_snapshot_context s), + 'lease',(SELECT jsonb_agg(to_jsonb(l) ORDER BY lease_id) FROM mst2_snapshot_lease l), + 'prepare',(SELECT jsonb_agg(to_jsonb(p) ORDER BY prepare_id) FROM mst2_metadata_prepare p), + 'member',(SELECT jsonb_agg(to_jsonb(m) ORDER BY prepare_id,page_id) FROM mst2_metadata_prepare_page m), + 'payload',(SELECT jsonb_agg(to_jsonb(p) ORDER BY page_id) FROM mst2_metadata_payload p), + 'root',(SELECT jsonb_agg(to_jsonb(r) ORDER BY root_key,node_id) FROM mst2_retention_root r))::text", + ) + .await +} + +async fn roots(db: &C) -> Value { + json_sql( + db, + "SELECT coalesce(jsonb_agg(to_jsonb(r) ORDER BY root_key,node_id),'[]'::jsonb)::text FROM mst2_retention_root r", + ) + .await +} + +async fn domain_boundary_digests(db: &DatabaseConnection) -> Value { + let mut inventory = serde_json::Map::new(); + for (table, keys) in [ + ("mst2_snapshot_context", "r.snapshot_id"), + ("mst2_snapshot_lease", "r.lease_id"), + ("mst2_metadata_prepare", "r.prepare_id"), + ("mst2_metadata_prepare_page", "r.prepare_id,r.page_id"), + ("mst2_metadata_payload", "r.page_id"), + ("mst2_metadata_lifetime", "r.page_id,r.generation"), + ("mst2_metadata_current", "r.page_id"), + ("mst2_metadata_graph_node", "r.page_id,r.generation"), + ( + "mst2_metadata_graph_edge", + "r.parent_page,r.parent_generation,r.child_page,r.child_generation", + ), + ( + "mst2_metadata_graph_root", + "r.prepare_id,r.page_id,r.generation", + ), + ("mst2_metadata_gc_op", "r.operation_id"), + ("mst2_retention_node", "r.node_id"), + ("mst2_retention_edge", "r.parent_id,r.child_id"), + ("mst2_retention_root", "r.node_id,r.root_key"), + ("mst2_retention_gc_op", "r.operation_id"), + ("mst2_metadata_install_seal", "r.prepare_id"), + ("mst2_metadata_storage_scope", "r.singleton"), + ("mst2_metadata_namespace", "r.namespace_uuid"), + ("mst2_snapshot_storage_route", "r.snapshot_id"), + ( + "mst2_generic_session_storage_binding", + "r.session_incarnation", + ), + ("mst2_lease_storage_route", "r.lease_id"), + ] { + // Keep plans, binding bytes and payloads in PostgreSQL; only their + // keyed row digests leave the database for the rollback oracle. + inventory.insert( + table.to_owned(), + json_sql( + db, + &format!( + "SELECT coalesce(jsonb_agg(jsonb_build_object('key',jsonb_build_array({keys}), + 'sha256',encode(sha256(convert_to(to_jsonb(r)::text,'UTF8')),'hex')) + ORDER BY {keys}),'[]'::jsonb)::text FROM {table} r", + ), + ) + .await, + ); + } + Value::Object(inventory) +} + +fn resolve_request() -> Request { + Request::builder() + .method("POST") + .uri("/api/v2/snapshots/resolve") + .header("authorization", format!("Bearer {TOKEN}")) + .header("content-type", "application/json") + .body(Body::from( + json!({"target":{"kind":"latest"},"scope":"/project"}).to_string(), + )) + .unwrap() +} + +async fn lease_control(fixture: &Fixture, lease: &str, renew: bool) -> Response { + fixture + .app + .clone() + .oneshot( + Request::builder() + .method(if renew { "POST" } else { "DELETE" }) + .uri(format!( + "/api/v2/snapshots/leases/{lease}{}", + if renew { "/renew" } else { "" } + )) + .header("authorization", format!("Bearer {TOKEN}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap() +} + +async fn rebuilt(fixture: &Fixture) -> Router { + let state = rebuilt_state(fixture).await; + Router::new().nest("/api/v2", routers(state.clone()).with_state(state)) +} + +async fn rebuilt_state(fixture: &Fixture) -> MonoApiServiceState { + let config = fixture.state.storage.config(); + let mut database = config.database.clone(); + database.max_connection = 1; + database.min_connection = 1; + let connection = crate::jupiter::storage::init::postgres_connection(&database) + .await + .unwrap(); + let storage = crate::jupiter::storage::Storage::new_with_connection( + config, + Arc::new(connection), + fixture.state.storage.git_service.obj_storage.clone(), + ) + .await + .unwrap(); + MonoApiServiceState { + storage, + git_object_cache: Arc::new(GitObjectCache { + connection: fixture.state.git_object_cache.connection.clone(), + prefix: uuid::Uuid::new_v4().to_string(), + }), + ..fixture.state.clone() + } +} + +async fn assert_exact_bindings(db: &C, leases: i64) { + assert_eq!( + scalar(db, "SELECT count(*) FROM mst2_metadata_namespace").await, + 1 + ); + assert_eq!( + scalar(db, "SELECT count(*) FROM mst2_snapshot_storage_route").await, + 1 + ); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM mst2_generic_session_storage_binding" + ) + .await, + 1 + ); + assert_eq!( + scalar(db, "SELECT count(*) FROM mst2_lease_storage_route").await, + leases + ); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM mst2_snapshot_context s + JOIN mst2_metadata_prepare p ON p.prepare_id=s.prepare_id + JOIN mst2_snapshot_storage_route r ON r.snapshot_id=s.snapshot_id + JOIN mst2_generic_session_storage_binding b ON b.snapshot_id=s.snapshot_id + JOIN mst2_metadata_namespace n ON n.namespace_uuid=r.namespace_uuid + WHERE r.canonical_descriptor=s.canonical_descriptor + AND r.instance_id=s.instance_id AND r.commit_oid=s.commit_oid + AND r.root_tree_oid=s.root_tree_oid AND r.metadata_root=s.metadata_root + AND b.namespace_uuid=r.namespace_uuid AND b.prepare_id=s.prepare_id + AND b.metadata_root=s.metadata_root AND p.metadata_root=s.metadata_root + AND r.source_profile=jsonb_build_object('source_domain',p.source_domain, + 'tagged_root_tree_oid',p.tagged_root_tree_oid,'scope',p.scope, + 'schema_version',p.schema_version,'metadata_codec',p.metadata_codec, + 'materialization_policy',p.materialization_policy,'fs_semantics',p.fs_semantics, + 'access_projection',p.access_projection,'verification_revision',p.verification_revision, + 'projection_revision',p.projection_revision)", + ) + .await, + 1, + ); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM mst2_snapshot_lease l + JOIN mst2_lease_storage_route r ON r.lease_id=l.lease_id + JOIN mst2_generic_session_storage_binding b ON b.snapshot_id=l.snapshot_id + WHERE r.snapshot_id=l.snapshot_id AND r.namespace_uuid=b.namespace_uuid + AND r.session_incarnation=b.session_incarnation AND r.prepare_id=b.prepare_id + AND r.metadata_root=b.metadata_root AND r.authorization_epoch=l.authorization_epoch + AND r.publication_sequence=l.publication_sequence AND r.writer_epoch=l.writer_epoch + AND r.certificate_receipt_id=l.certificate_receipt_id", + ) + .await, + leases, + ); + let stored = routes(db).await; + uuid::Uuid::parse_str( + stored["binding"][0]["session_incarnation"] + .as_str() + .unwrap(), + ) + .unwrap(); + let namespace = &stored["namespace"][0]; + assert_eq!(namespace["graph_domain"], "generic-v1"); + assert_eq!(namespace["family_identity"], "v3-generic-session-1"); + assert_eq!(namespace["admission_state"], "G_ADMITTED_Q_CLOSED"); + assert_eq!(namespace["collector_state"], "CLOSED"); + assert_eq!(scalar(db, + "SELECT count(*) FROM mst2_metadata_namespace n JOIN mst2_metadata_storage_scope s ON s.singleton=1 + JOIN pg_namespace c ON c.nspname=current_schema() JOIN pg_database d ON d.datname=current_database() + WHERE n.core_schema=c.nspname AND n.core_schema_oid=c.oid + AND n.metadata_schema=c.nspname AND n.metadata_schema_oid=c.oid + AND n.database_name=d.datname AND n.database_oid=d.oid AND n.storage_uuid=s.storage_uuid + AND n.mono_lock_key2=hashtext(current_schema()) + AND n.server_address IS NOT DISTINCT FROM inet_server_addr()::text + AND n.server_port IS NOT DISTINCT FROM inet_server_port()", + ).await, 1); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM information_schema.columns + WHERE table_schema=current_schema() AND table_name='mst2_snapshot_storage_route' + AND column_name IN ('prepare_id','generation','root_generation','session_incarnation')", + ) + .await, + 0, + "the permanent SID route must not freeze a physical incarnation", + ); + assert_eq!( + scalar( + db, + "SELECT count(*) FROM pg_constraint WHERE conrelid='mst2_lease_storage_route'::regclass + AND contype='f' AND confrelid='mst2_generic_session_storage_binding'::regclass", + ) + .await, + 0, + "future namespace-specific incarnations must not have a universal G ledger FK" + ); +} + +async fn rejected(db: &DatabaseConnection, sql: &str) { + let txn = db.begin().await.unwrap(); + match txn.execute_unprepared(sql).await { + Ok(_) => { + assert!(txn.commit().await.is_err(), "mutation committed: {sql}"); + } + Err(error) => { + assert!( + error.to_string().contains("storage route"), + "wrong rejection for {sql}: {error}" + ); + txn.rollback().await.unwrap(); + } + } +} + +#[tokio::test] +async fn mst2_generic_storage_routes_actual_resolve_warm_and_fresh_service_keep_exact_tuple() { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + assert_exact_bindings(db, 1).await; + let initial = routes(db).await; + assert_eq!(initial["snapshot"][0]["snapshot_id"], fixture.snapshot); + assert_eq!(initial["lease"][0]["lease_id"], fixture.lease); + let original = success_json(fixture.send("GET", "descriptor", Body::empty()).await).await; + let warm = success_json( + fixture + .app + .clone() + .oneshot(resolve_request()) + .await + .unwrap(), + ) + .await; + assert_eq!(warm["descriptor"]["snapshot_id"], fixture.snapshot); + assert_ne!(warm["lease_id"], fixture.lease); + assert_exact_bindings(db, 2).await; + let after = routes(db).await; + for key in ["namespace", "snapshot", "binding"] { + assert_eq!(after[key], initial[key], "warm resolve changed {key}"); + } + let cold = rebuilt(&fixture).await; + assert_eq!( + success_json( + cold.clone() + .oneshot(fixture.request("GET", "descriptor", Body::empty())) + .await + .unwrap() + ) + .await, + original, + ); + assert_eq!( + cold.oneshot(fixture.request("HEAD", "blob?path=/file", Body::empty())) + .await + .unwrap() + .status(), + 200, + ); + assert_eq!(routes(db).await, after); + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_every_member_delete_and_truncate_are_immutable() { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let original = routes(db).await; + let source = sources(db).await; + for table in ROUTE_TABLES { + let columns = db + .query_all_raw(statement( + "SELECT column_name,udt_name FROM information_schema.columns + WHERE table_schema=current_schema() AND table_name=$1 ORDER BY ordinal_position", + [table.into()], + )) + .await + .unwrap(); + assert!(!columns.is_empty(), "missing route relation: {table}"); + for column in columns { + let name: String = column.try_get("", "column_name").unwrap(); + let kind: String = column.try_get("", "udt_name").unwrap(); + let quoted = format!("\"{}\"", name.replace('"', "\"\"")); + let mutation = match kind.as_str() { + "uuid" => format!("'{}'::uuid", uuid::Uuid::new_v4()), + "text" | "varchar" | "bpchar" => format!("coalesce({quoted},'')||'-mutated'"), + "int2" | "int4" | "int8" => format!("coalesce({quoted},0)+1"), + "oid" => format!("({quoted}::bigint+1)::oid"), + "bytea" => format!("set_byte({quoted},0,(get_byte({quoted},0)+1)%256)"), + "jsonb" => format!("{quoted}||jsonb_build_object('_mutation',true)"), + "timestamptz" => format!("{quoted}+interval '1 second'"), + "bool" => format!("NOT {quoted}"), + other => panic!("uncovered immutable route member {table}.{name}: {other}"), + }; + rejected(db, &format!("UPDATE {table} SET {quoted}={mutation}")).await; + assert_eq!(routes(db).await, original, "{table}.{name}"); + } + rejected(db, &format!("DELETE FROM {table}")).await; + rejected(db, &format!("TRUNCATE {table} CASCADE")).await; + assert_eq!(routes(db).await, original, "{table}"); + } + for table in ["mst2_snapshot_context", "mst2_snapshot_lease"] { + rejected(db, &format!("DELETE FROM {table}")).await; + rejected(db, &format!("TRUNCATE {table} CASCADE")).await; + } + assert_eq!(sources(db).await, source); + fixture.counts.assert(0, 0); +} + +fn lease_insert(lease: &str) -> String { + format!( + "INSERT INTO mst2_snapshot_lease(lease_id,snapshot_id,authorization_epoch, + publication_sequence,writer_epoch,certificate_receipt_id,expires_at_unix,state) + SELECT '{lease}',snapshot_id,authorization_epoch,publication_sequence,writer_epoch, + certificate_receipt_id,floor(extract(epoch FROM clock_timestamp()))::bigint+3600,'ACTIVE' + FROM mst2_snapshot_context", + ) +} + +#[tokio::test] +async fn mst2_generic_storage_routes_raw_lease_derives_atomically_and_half_rows_roll_back() { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let original = routes(db).await; + let source = sources(db).await; + let lease = uuid::Uuid::new_v4().to_string(); + let txn = db.begin().await.unwrap(); + assert_eq!( + txn.execute_unprepared(&lease_insert(&lease)) + .await + .unwrap() + .rows_affected(), + 1 + ); + assert_exact_bindings(&txn, 2).await; + txn.rollback().await.unwrap(); + assert_eq!(routes(db).await, original); + assert_eq!(sources(db).await, source); + super::storage_route_fixture::reject_half_lease_commit_for_test( + db, + &lease, + &lease_insert(&lease), + ) + .await; + assert_eq!(routes(db).await, original); + assert_eq!(sources(db).await, source); + for table in [ + "mst2_snapshot_storage_route", + "mst2_generic_session_storage_binding", + ] { + assert_eq!(routes(db).await, original); + rejected( + db, + &format!( + "INSERT INTO {table} SELECT (jsonb_populate_record(NULL::{table}, + to_jsonb(r)||jsonb_build_object('snapshot_id','sha256:{}'))).* FROM {table} r", + hex::encode([7; 32]), + ), + ) + .await; + } + rejected(db, &format!( + "INSERT INTO mst2_lease_storage_route + SELECT (jsonb_populate_record(NULL::mst2_lease_storage_route, + to_jsonb(r)||jsonb_build_object('lease_id','{lease}'))).* FROM mst2_lease_storage_route r", + )).await; + rejected(db, + "INSERT INTO mst2_metadata_namespace SELECT (jsonb_populate_record(NULL::mst2_metadata_namespace, + to_jsonb(n)||jsonb_build_object('namespace_uuid','00000000-0000-4000-8000-000000000001', + 'graph_domain','qualified-v1'))).* FROM mst2_metadata_namespace n", + ).await; + assert_eq!(routes(db).await, original); + assert_eq!(sources(db).await, source); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_renew_terminal_and_unknown_release_preserve_route_and_roots() { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let original = routes(db).await; + let unknown = uuid::Uuid::new_v4().to_string(); + let root: Vec = db + .query_one_raw(Statement::from_string( + DbBackend::Postgres, + "SELECT metadata_root FROM mst2_snapshot_context", + )) + .await + .unwrap() + .unwrap() + .try_get_by_index(0) + .unwrap(); + let txn = db.begin().await.unwrap(); + PostgresRetentionRepository::acquire_existing_roots_in_txn( + &txn, + &format!("page:sha256:{}", hex::encode(root)), + &[RetentionRoot::Lease(unknown.clone())], + ) + .await + .unwrap(); + txn.commit().await.unwrap(); + let protected = roots(db).await; + assert_eq!( + success_json(lease_control(&fixture, &unknown, false).await).await["released"], + false + ); + assert_eq!( + roots(db).await, + protected, + "unknown release must not remove an unrelated protection root" + ); + assert_eq!(routes(db).await, original); + let renewed = success_json(lease_control(&fixture, &fixture.lease, true).await).await; + assert_eq!(renewed["lease_id"], fixture.lease); + assert_eq!(renewed["snapshot_id"], fixture.snapshot); + assert_eq!(routes(db).await, original); + assert_eq!( + success_json(lease_control(&fixture, &fixture.lease, false).await).await["released"], + true + ); + let terminal_roots = roots(db).await; + assert_eq!(routes(db).await, original); + assert_eq!( + success_json(lease_control(&fixture, &fixture.lease, false).await).await["released"], + false + ); + assert_eq!(roots(db).await, terminal_roots); + error( + lease_control(&fixture, &fixture.lease, true).await, + 410, + "LEASE_EXPIRED", + false, + ) + .await; + error( + fixture.send("GET", "descriptor", Body::empty()).await, + 410, + "LEASE_EXPIRED", + false, + ) + .await; + assert_eq!( + fixture + .send("HEAD", "blob?path=/file", Body::empty()) + .await + .status(), + 410 + ); + assert_eq!(routes(db).await, original); + assert_eq!(roots(db).await, terminal_roots); + assert_exact_bindings(db, 1).await; + fixture.counts.assert(0, 0); +} + +async fn wait_for_advisory_waiter(held: &DatabaseTransaction, key: i32) -> i64 { + tokio::time::timeout(Duration::from_secs(30), async { + loop { + held.execute_unprepared("SELECT pg_stat_clear_snapshot()").await.unwrap(); + let row = held.query_one_raw(statement( + "SELECT l.pid::bigint AS pid FROM pg_locks l JOIN pg_stat_activity a ON a.pid=l.pid + WHERE l.locktype='advisory' AND l.classid=$1::bigint::oid AND l.objid=hashtext(current_schema())::oid + AND l.objsubid=2 AND NOT l.granted AND a.datname=current_database() + AND a.application_name=current_schema() LIMIT 1", + [i64::from(key).into()], + )).await.unwrap(); + if let Some(row) = row { break row.try_get("", "pid").unwrap(); } + tokio::task::yield_now().await; + } + }).await.expect("raw lease statement did not wait on the expected advisory barrier") +} + +async fn lock_count(held: &DatabaseTransaction, pid: i64, key: i32) -> i64 { + held.query_one_raw(statement( + "SELECT count(*)::bigint FROM pg_locks WHERE pid::bigint=$1 AND locktype='advisory' + AND classid=$2::bigint::oid AND objid=hashtext(current_schema())::oid AND objsubid=2 AND granted", + [pid.into(), i64::from(key).into()], + )).await.unwrap().unwrap().try_get_by_index(0).unwrap() +} + +#[tokio::test] +async fn mst2_generic_storage_routes_raw_statement_waits_mono_then_route_then_retention() { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + for held_key in [MONO_WRITE_LOCK_KEY1, ROUTE_LOCK_KEY, RETENTION_LOCK_KEY] { + let held = db.begin().await.unwrap(); + held.execute_raw(statement( + "SELECT pg_advisory_xact_lock($1,hashtext(current_schema()))", + [held_key.into()], + )) + .await + .unwrap(); + let writing = { + let db = db.clone(); + let lease = fixture.lease.clone(); + tokio::spawn(async move { + let txn = db.begin().await.unwrap(); + let result = txn.execute_unprepared(&format!( + "{} ON CONFLICT(lease_id) DO UPDATE SET expires_at_unix=EXCLUDED.expires_at_unix", + lease_insert(&lease), + )).await; + txn.rollback().await.unwrap(); + result + }) + }; + let pid = wait_for_advisory_waiter(&held, held_key).await; + assert_eq!( + lock_count(&held, pid, MONO_WRITE_LOCK_KEY1).await, + if held_key == MONO_WRITE_LOCK_KEY1 { + 0 + } else { + 1 + } + ); + assert_eq!( + lock_count(&held, pid, ROUTE_LOCK_KEY).await, + if held_key == RETENTION_LOCK_KEY { 1 } else { 0 } + ); + assert_eq!(lock_count(&held, pid, RETENTION_LOCK_KEY).await, 0); + assert!( + held.query_one_raw(statement( + "SELECT lease_id FROM mst2_snapshot_lease WHERE lease_id=$1 FOR UPDATE NOWAIT", + [fixture.lease.clone().into()], + )) + .await + .unwrap() + .is_some(), + "the blocked statement locked its existing lease before the barrier" + ); + assert_eq!( + held.query_one_raw(statement( + "SELECT count(*)::bigint FROM pg_locks WHERE pid::bigint=$1 AND locktype='tuple' + AND relation IN ('mst2_snapshot_context'::regclass,'mst2_snapshot_lease'::regclass)", + [pid.into()], + )) + .await + .unwrap() + .unwrap() + .try_get_by_index::(0) + .unwrap(), + 0 + ); + held.rollback().await.unwrap(); + assert_eq!(writing.await.unwrap().unwrap().rows_affected(), 1); + assert_exact_bindings(db, 1).await; + } + for (key, message) in [ + ( + RETENTION_LOCK_KEY, + "cannot acquire core locks after retention", + ), + (ROUTE_LOCK_KEY, "lock was acquired before core mono"), + ] { + for lock in ["pg_advisory_xact_lock", "pg_advisory_xact_lock_shared"] { + let held = db.begin().await.unwrap(); + held.execute_raw(statement( + &format!("SELECT {lock}($1,hashtext(current_schema()))"), + [key.into()], + )) + .await + .unwrap(); + let result = tokio::time::timeout( + Duration::from_secs(5), + held.execute_unprepared(&lease_insert(&uuid::Uuid::new_v4().to_string())), + ) + .await + .expect("reverse-order statement waited instead of rejecting"); + let error = result.unwrap_err(); + assert!( + error.to_string().contains(message), + "wrong reverse-order {lock} rejection: {error}" + ); + held.rollback().await.unwrap(); + } + } + assert_exact_bindings(db, 1).await; +} + +#[tokio::test] +async fn mst2_generic_storage_routes_schema_isolation_wrong_caller_rr_and_temp_shadow_fail_closed() +{ + let fixture = Fixture::new_with_pg_config(true).await; + let alien = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let original = routes(db).await; + let other = routes(alien.state.storage.mono_storage().get_connection()).await; + assert_ne!( + original["namespace"][0]["namespace_uuid"], + other["namespace"][0]["namespace_uuid"] + ); + let transplanted = original["snapshot"][0].clone(); + let alien_db = alien.state.storage.mono_storage(); + let txn = alien_db.get_connection().begin().await.unwrap(); + let transplant = txn.execute_raw(statement( + "INSERT INTO mst2_snapshot_storage_route SELECT (jsonb_populate_record(NULL::mst2_snapshot_storage_route, + $1::jsonb||jsonb_build_object('namespace_uuid',(SELECT namespace_uuid FROM mst2_metadata_namespace)))).*", + [transplanted.to_string().into()], + )).await.unwrap_err(); + assert!(transplant.to_string().contains("actual generic context")); + txn.rollback().await.unwrap(); + let schema = fixture._schema.as_ref().unwrap().schema(); + let qualified = format!("\"{}\".mst2_route_enter", schema.replace('"', "\"\"")); + let txn = db + .begin_with_config(Some(IsolationLevel::RepeatableRead), None) + .await + .unwrap(); + assert!( + txn.execute_unprepared(&lease_insert(&uuid::Uuid::new_v4().to_string())) + .await + .is_err() + ); + txn.rollback().await.unwrap(); + let txn = db.begin().await.unwrap(); + txn.execute_unprepared("SET LOCAL search_path=pg_catalog") + .await + .unwrap(); + assert!( + txn.execute_raw(statement( + &format!("SELECT {qualified}($1)"), + ["pg_catalog".into()] + )) + .await + .is_err() + ); + txn.rollback().await.unwrap(); + let txn = db.begin().await.unwrap(); + assert!( + txn.execute_raw(statement( + &format!("SELECT {qualified}($1)"), + [alien._schema.as_ref().unwrap().schema().into()] + )) + .await + .is_err() + ); + txn.rollback().await.unwrap(); + let txn = db.begin().await.unwrap(); + txn.execute_unprepared( + "CREATE TEMP TABLE mst2_metadata_namespace(namespace_uuid uuid); + CREATE TEMP TABLE mst2_snapshot_storage_route(snapshot_id text); + CREATE TEMP TABLE mst2_generic_session_storage_binding(snapshot_id text); + CREATE TEMP TABLE mst2_lease_storage_route(lease_id text)", + ) + .await + .unwrap(); + txn.execute_raw(statement( + &format!("SELECT {qualified}($1)"), + [schema.into()], + )) + .await + .unwrap(); + assert_eq!( + scalar( + &txn, + &format!( + "SELECT count(*) FROM \"{}\".mst2_metadata_namespace", + schema.replace('"', "\"\"") + ) + ) + .await, + 1 + ); + assert_eq!( + scalar(&txn, "SELECT count(*) FROM pg_temp.mst2_metadata_namespace").await, + 0 + ); + txn.rollback().await.unwrap(); + error( + rebuilt(&alien) + .await + .oneshot(fixture.request("GET", "descriptor", Body::empty())) + .await + .unwrap(), + 410, + "LEASE_EXPIRED", + false, + ) + .await; + assert_eq!(routes(db).await, original); + assert_eq!( + routes(alien.state.storage.mono_storage().get_connection()).await, + other + ); + fixture.counts.assert(0, 0); + alien.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_additive_backfill_keeps_active_terminal_sources_bytes_and_roots() + { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let warm = success_json( + fixture + .app + .clone() + .oneshot(resolve_request()) + .await + .unwrap(), + ) + .await; + let second = warm["lease_id"].as_str().unwrap(); + assert_eq!( + success_json(lease_control(&fixture, second, false).await).await["released"], + true + ); + assert_exact_bindings(db, 2).await; + let source = sources(db).await; + let original = routes(db).await; + let descriptor = success_json(fixture.send("GET", "descriptor", Body::empty()).await).await; + super::storage_route_fixture::restore_pre_route_schema(db).await; + assert_eq!(sources(db).await, source); + Migrator::up(db, None).await.unwrap(); + assert_eq!(sources(db).await, source); + assert_exact_bindings(db, 2).await; + let restored = routes(db).await; + for key in [ + "canonical_descriptor", + "instance_id", + "commit_oid", + "root_tree_oid", + "metadata_root", + "source_profile", + ] { + assert_eq!( + restored["snapshot"][0][key], original["snapshot"][0][key], + "backfill changed {key}" + ); + } + assert_eq!( + success_json( + rebuilt(&fixture) + .await + .oneshot(fixture.request("GET", "descriptor", Body::empty()),) + .await + .unwrap() + ) + .await, + descriptor + ); + assert_eq!( + success_json(lease_control(&fixture, second, false).await).await["released"], + false + ); + assert_eq!(sources(db).await, source); + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_actual_http_ignores_temp_source_and_ledger_shadows() { + let fixture = Fixture::new_with_pg_config(true).await; + let state = rebuilt_state(&fixture).await; + let app = Router::new().nest("/api/v2", routers(state.clone()).with_state(state.clone())); + let mono = state.storage.mono_storage(); + let db = mono.get_connection(); + let original = routes(db).await; + assert_eq!( + app.clone() + .oneshot(fixture.request("HEAD", "blob?path=/file", Body::empty())) + .await + .unwrap() + .status(), + 200 + ); + let schema = fixture + ._schema + .as_ref() + .unwrap() + .schema() + .replace('"', "\"\""); + db.execute_unprepared(&format!( + "CREATE TEMP TABLE mst2_snapshot_context(LIKE \"{schema}\".mst2_snapshot_context); + CREATE TEMP TABLE mst2_snapshot_lease(LIKE \"{schema}\".mst2_snapshot_lease); + CREATE TEMP TABLE mst2_metadata_namespace(namespace_uuid uuid); + CREATE TEMP TABLE mst2_snapshot_storage_route(snapshot_id text); + CREATE TEMP TABLE mst2_generic_session_storage_binding(snapshot_id text); + CREATE TEMP TABLE mst2_lease_storage_route(lease_id text)", + )) + .await + .unwrap(); + assert_eq!( + app.clone() + .oneshot(fixture.request("HEAD", "blob?path=/file", Body::empty())) + .await + .unwrap() + .status(), + 200 + ); + let renewed = success_json( + app.oneshot( + Request::builder() + .method("POST") + .uri(format!("/api/v2/snapshots/leases/{}/renew", fixture.lease)) + .header("authorization", format!("Bearer {TOKEN}")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(), + ) + .await; + assert_eq!(renewed["lease_id"], fixture.lease); + assert_eq!(renewed["snapshot_id"], fixture.snapshot); + for table in [ + "mst2_snapshot_context", + "mst2_snapshot_lease", + "mst2_metadata_namespace", + "mst2_snapshot_storage_route", + "mst2_generic_session_storage_binding", + "mst2_lease_storage_route", + ] { + assert_eq!( + scalar(db, &format!("SELECT count(*) FROM pg_temp.{table}")).await, + 0 + ); + } + let real = fixture.state.storage.mono_storage(); + assert_eq!(routes(real.get_connection()).await, original); + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_temp_prepare_shadow_rejects_qualified_context_and_registered_generic_rebind() + { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let original_routes = routes(db).await; + let unique = uuid::Uuid::new_v4().to_string(); + let entries = [Entry::file( + EntryKind::Regular, + unique.as_bytes(), + 3, + digest(unique.as_bytes()), + )]; + let child = Page::build(&entries).unwrap(); + let roots = [Entry::dir(b"qualified-route-boundary", page_id(&child))]; + let root = Page::build(&roots).unwrap(); + let mut builder = MetadataDagBuilder::new(MetadataDagLimits::default()); + builder.add_directory(&child, &entries).unwrap(); + builder.add_directory(&root, &roots).unwrap(); + let prepared = PreparedNativeMetadataRetention::test_installation( + Arc::new(builder.finish(page_id(&root)).unwrap()), + "/", + ); + let tree = "a".repeat(40); + assert_eq!(prepared.fixed_root_tree_oid(), format!("sha1:{tree}")); + let qualified = PostgresQualifiedMetadataRepository::new(db.clone()) + .await + .unwrap(); + let intent = qualified + .begin_intent(&format!("qualified-route-boundary:{unique}"), &prepared) + .await + .unwrap(); + qualified + .install_pages(&intent, prepared.dag().payloads()) + .await + .unwrap(); + let receipt = qualified.finalize(&intent).await.unwrap(); + assert_eq!(receipt.metadata_root(), prepared.dag().root()); + let committed = db.query_one_raw(statement( + "SELECT state,graph_domain,metadata_root,tagged_root_tree_oid FROM mst2_metadata_prepare WHERE prepare_id=$1", + [intent.prepare_id().into()], + )).await.unwrap().unwrap(); + assert_eq!( + committed.try_get::("", "state").unwrap(), + "COMMITTED" + ); + assert_eq!( + committed.try_get::("", "graph_domain").unwrap(), + "qualified-v1" + ); + assert_eq!( + committed.try_get::>("", "metadata_root").unwrap(), + receipt.metadata_root().to_vec() + ); + assert_eq!( + committed + .try_get::("", "tagged_root_tree_oid") + .unwrap(), + format!("sha1:{tree}") + ); + assert_eq!(scalar(db, + "SELECT count(*) FROM mst2_metadata_lifetime WHERE graph_domain='qualified-v1' AND state='LIVE'", + ).await, prepared.dag().payloads().len() as i64); + assert_eq!(routes(db).await, original_routes); + let before = domain_boundary_digests(db).await; + let descriptor = ServingDescriptor { + instance_uuid: *uuid::Uuid::parse_str( + fixture + .state + .storage + .config() + .mst2 + .instance_uuid + .as_deref() + .unwrap(), + ) + .unwrap() + .as_bytes(), + namespace_view_id: digest(unique.as_bytes()), + scope: prepared.scope().to_owned(), + metadata_root: receipt.metadata_root(), + }; + let sid = format!("sha256:{}", hex::encode(descriptor.snapshot_id().unwrap())); + assert_ne!(sid, fixture.snapshot); + let schema = format!( + "\"{}\"", + fixture + ._schema + .as_ref() + .unwrap() + .schema() + .replace('"', "\"\"") + ); + let shadow = + format!("CREATE TEMP TABLE mst2_metadata_prepare(LIKE {schema}.mst2_metadata_prepare)"); + let txn = db.begin().await.unwrap(); + txn.execute_unprepared(&shadow).await.unwrap(); + txn.execute_unprepared(&format!("SET LOCAL search_path={schema},pg_catalog")) + .await + .unwrap(); + assert_eq!( + scalar(&txn, "SELECT count(*) FROM mst2_metadata_prepare").await, + 0 + ); + assert_eq!(scalar(&txn, "SELECT (to_regclass('mst2_metadata_prepare')='pg_temp.mst2_metadata_prepare'::regclass)::bigint").await, 1); + let rejected = txn.execute_raw(statement(&format!( + "INSERT INTO {schema}.mst2_snapshot_context(snapshot_id,canonical_descriptor,instance_id,commit_oid, + root_tree_oid,metadata_root,prepare_id,publication_sequence,writer_epoch,certificate_receipt_id,authorization_epoch,state) + SELECT $1,$2,instance_id,commit_oid,$3,$4,$5,publication_sequence,writer_epoch,certificate_receipt_id,authorization_epoch,'READY' + FROM {schema}.mst2_snapshot_context WHERE snapshot_id=$6", + ), [sid.into(),descriptor.encode().unwrap().into(),tree.into(),receipt.metadata_root().to_vec().into(), + intent.prepare_id().into(),fixture.snapshot.clone().into()], + )).await.unwrap_err(); + assert!( + rejected + .to_string() + .contains("storage route is not derived from its actual generic context"), + "{rejected}" + ); + txn.rollback().await.unwrap(); + assert_eq!( + domain_boundary_digests(db).await, + before, + "rejected Q adoption changed source, bytes, graph roots or route inventory" + ); + let generic = db + .query_one_raw(statement( + &format!( + "SELECT p.prepare_id,p.graph_domain FROM {schema}.mst2_metadata_prepare p + JOIN {schema}.mst2_metadata_install_seal i ON i.prepare_id=p.prepare_id + JOIN {schema}.mst2_snapshot_context s ON s.prepare_id=p.prepare_id WHERE s.snapshot_id=$1", + ), + [fixture.snapshot.clone().into()], + )) + .await + .unwrap() + .unwrap(); + let prepare_id: String = generic.try_get("", "prepare_id").unwrap(); + let domain: Option = generic.try_get("", "graph_domain").unwrap(); + assert!( + domain + .as_deref() + .is_none_or(|domain| domain == "generic-v1") + ); + let txn = db.begin().await.unwrap(); + txn.execute_unprepared(&shadow).await.unwrap(); + assert_eq!( + scalar(&txn, "SELECT count(*) FROM mst2_metadata_prepare").await, + 0 + ); + let rejected = txn.execute_raw(statement(&format!( + "UPDATE {schema}.mst2_metadata_prepare SET graph_domain='qualified-v1' WHERE prepare_id=$1", + ), [prepare_id.clone().into()])).await.unwrap_err(); + let error = rejected.to_string(); + assert!( + error.contains("registered metadata preparation identity is immutable") + || error.contains("metadata generation seal cannot be rebound"), + "{error}" + ); + txn.rollback().await.unwrap(); + let after_domain: Option = db.query_one_raw(statement(&format!( + "SELECT graph_domain FROM {schema}.mst2_metadata_prepare WHERE prepare_id=$1", + ), [prepare_id.into()])).await.unwrap().unwrap().try_get_by_index(0).unwrap(); + assert_eq!(after_domain, domain); + assert_eq!( + domain_boundary_digests(db).await, + before, + "registered G rebind changed source, bytes, graph roots or route inventory" + ); + assert_eq!(routes(db).await, original_routes); + assert_exact_bindings(db, 1).await; + assert_eq!( + fixture + .send("HEAD", "blob?path=/file", Body::empty()) + .await + .status(), + 200 + ); + fixture.counts.assert(0, 0); +} + +#[tokio::test] +async fn mst2_generic_storage_routes_corrupt_original_incarnation_rejects_reads_renew_release_without_repair() + { + let fixture = Fixture::new_with_pg_config(true).await; + let mono = fixture.state.storage.mono_storage(); + let db = mono.get_connection(); + let source = sources(db).await; + let original = routes(db).await; + super::storage_route_fixture::overwrite_lease_route_incarnation_for_test(db, &fixture.lease) + .await; + let corrupted = routes(db).await; + assert_ne!(corrupted["lease"], original["lease"]); + for key in ["namespace", "snapshot", "binding"] { + assert_eq!(corrupted[key], original[key]); + } + error( + fixture.send("GET", "descriptor", Body::empty()).await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + assert_eq!( + fixture + .send("HEAD", "blob?path=/file", Body::empty()) + .await + .status(), + 502 + ); + error( + lease_control(&fixture, &fixture.lease, true).await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + error( + lease_control(&fixture, &fixture.lease, false).await, + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + error( + rebuilt(&fixture) + .await + .oneshot(fixture.request("GET", "descriptor", Body::empty())) + .await + .unwrap(), + 502, + "INTEGRITY_ERROR", + false, + ) + .await; + assert_eq!( + routes(db).await, + corrupted, + "route corruption must not be silently repaired" + ); + assert_eq!( + sources(db).await, + source, + "failed access must not change leases, plans, bytes or roots" + ); + fixture.counts.assert(0, 0); +} diff --git a/src/jupiter/migration/m20261007_000600_add_mst2_storage_routes.rs b/src/jupiter/migration/m20261007_000600_add_mst2_storage_routes.rs new file mode 100644 index 00000000..541ce531 --- /dev/null +++ b/src/jupiter/migration/m20261007_000600_add_mst2_storage_routes.rs @@ -0,0 +1,42 @@ +//! Permanent semantic routes, with a separate existing generic session ledger. + +use sea_orm::{ConnectionTrait, DbBackend, Statement}; +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + let connection = manager.get_connection(); + let row = connection + .query_one_raw(Statement::from_string( + DbBackend::Postgres, + "SELECT current_schema() AS schema,n.oid::bigint AS oid + FROM pg_catalog.pg_namespace n WHERE n.nspname=current_schema()", + )) + .await? + .ok_or_else(|| DbErr::Custom("storage route core schema is missing".into()))?; + let schema: String = row.try_get("", "schema")?; + let oid: i64 = row.try_get("", "oid")?; + let quoted = format!("\"{}\"", schema.replace('"', "\"\"")); + let literal = format!("'{}'", schema.replace('\'', "''")); + #[cfg(test)] + let mono_key2 = format!("pg_catalog.hashtext({literal})"); + #[cfg(not(test))] + let mono_key2 = super::super::storage::push_queue_storage::MONO_WRITE_LOCK_KEY2.to_string(); + let sql = include_str!("m20261007_000600_storage_routes.sql") + .replace("$CORE_SCHEMA$", "ed) + .replace("$CORE_LITERAL$", &literal) + .replace("$CORE_OID$", &oid.to_string()) + .replace("$MONO_KEY2$", &mono_key2) + .replace("$NAMESPACE_UUID$", &uuid::Uuid::new_v4().to_string()); + connection.execute_unprepared(&sql).await?; + Ok(()) + } + + async fn down(&self, _manager: &SchemaManager) -> Result<(), DbErr> { + Ok(()) + } +} diff --git a/src/jupiter/migration/m20261007_000600_storage_routes.sql b/src/jupiter/migration/m20261007_000600_storage_routes.sql new file mode 100644 index 00000000..13e8b18c --- /dev/null +++ b/src/jupiter/migration/m20261007_000600_storage_routes.sql @@ -0,0 +1,329 @@ +DO $$ BEGIN + IF pg_catalog.pg_is_in_recovery() OR pg_catalog.current_setting('transaction_isolation')<>'read committed' + OR pg_catalog.current_schema()<>$CORE_LITERAL$ + OR NOT EXISTS(SELECT 1 FROM pg_catalog.pg_namespace WHERE oid=$CORE_OID$ AND nspname=$CORE_LITERAL$) THEN + RAISE EXCEPTION 'storage route migration requires its captured primary core schema'; + END IF; + IF EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND classid=1296718001::oid AND objid=pg_catalog.hashtext($CORE_LITERAL$)::oid AND objsubid=2 AND granted) + AND NOT EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND classid=1297043024::oid AND objid=($MONO_KEY2$)::oid AND objsubid=2 AND granted AND mode='ExclusiveLock') THEN + RAISE EXCEPTION 'storage route migration cannot acquire mono after route'; + END IF; + IF EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND database=(SELECT oid FROM pg_catalog.pg_database WHERE datname=pg_catalog.current_database()) + AND classid=1296717362::oid AND objid=pg_catalog.hashtext($CORE_LITERAL$)::oid AND objsubid=2 AND granted) + AND NOT (EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND classid=1297043024::oid AND objid=($MONO_KEY2$)::oid AND objsubid=2 AND granted AND mode='ExclusiveLock') + AND EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND classid=1296718001::oid AND objid=pg_catalog.hashtext($CORE_LITERAL$)::oid + AND objsubid=2 AND granted AND mode='ExclusiveLock')) THEN + RAISE EXCEPTION 'storage route migration cannot acquire core locks after retention'; + END IF; +END $$; +SELECT pg_catalog.set_config('lock_timeout','5000ms',true); +SELECT pg_catalog.pg_advisory_xact_lock(1297043024,$MONO_KEY2$); +SELECT pg_catalog.pg_advisory_xact_lock(1296718001,pg_catalog.hashtext($CORE_LITERAL$)); +SELECT pg_catalog.pg_advisory_xact_lock(1296717362,pg_catalog.hashtext($CORE_LITERAL$)); +SET LOCAL search_path=$CORE_SCHEMA$,pg_catalog,pg_temp; +LOCK TABLE mst2_snapshot_context,mst2_snapshot_lease,mst2_metadata_prepare, + mst2_metadata_storage_scope IN ACCESS EXCLUSIVE MODE; + +CREATE TABLE mst2_metadata_namespace ( + singleton integer PRIMARY KEY CHECK (singleton=1), + namespace_uuid uuid NOT NULL UNIQUE, + core_schema text NOT NULL, + core_schema_oid oid NOT NULL, + database_name text NOT NULL, + database_oid oid NOT NULL, + storage_uuid text NOT NULL, + server_address text, + server_port integer, + mono_lock_key2 integer NOT NULL, + metadata_schema text NOT NULL, + metadata_schema_oid oid NOT NULL, + family_identity text NOT NULL CHECK (family_identity='v3-generic-session-1'), + graph_domain text NOT NULL CHECK (graph_domain='generic-v1'), + admission_state text NOT NULL CHECK (admission_state='G_ADMITTED_Q_CLOSED'), + collector_state text NOT NULL CHECK (collector_state='CLOSED'), + CHECK (core_schema=metadata_schema AND core_schema_oid=metadata_schema_oid) +); +INSERT INTO mst2_metadata_namespace +SELECT 1,'$NAMESPACE_UUID$'::uuid,$CORE_LITERAL$,$CORE_OID$,pg_catalog.current_database(),d.oid,s.storage_uuid, + pg_catalog.inet_server_addr()::text,pg_catalog.inet_server_port(),$MONO_KEY2$,$CORE_LITERAL$,$CORE_OID$, + 'v3-generic-session-1','generic-v1','G_ADMITTED_Q_CLOSED','CLOSED' +FROM mst2_metadata_storage_scope s JOIN pg_catalog.pg_database d ON d.datname=pg_catalog.current_database() +WHERE s.singleton=1; + +CREATE TABLE mst2_snapshot_storage_route ( + snapshot_id text PRIMARY KEY, + namespace_uuid uuid NOT NULL REFERENCES mst2_metadata_namespace(namespace_uuid), + canonical_descriptor bytea NOT NULL, + instance_id text NOT NULL, + commit_oid text NOT NULL, + root_tree_oid text NOT NULL, + metadata_root bytea NOT NULL CHECK (octet_length(metadata_root)=32), + source_profile jsonb NOT NULL, + UNIQUE(snapshot_id,namespace_uuid), + CHECK (snapshot_id='sha256:'||pg_catalog.encode(pg_catalog.sha256( + pg_catalog.convert_to('mega.mst2.descriptor','UTF8')||pg_catalog.decode('00','hex')||canonical_descriptor),'hex')) +); +ALTER TABLE mst2_snapshot_context ADD UNIQUE(snapshot_id,prepare_id,metadata_root); +CREATE TABLE mst2_generic_session_storage_binding ( + session_incarnation uuid PRIMARY KEY, + snapshot_id text NOT NULL UNIQUE, + namespace_uuid uuid NOT NULL, + prepare_id text NOT NULL, + metadata_root bytea NOT NULL, + FOREIGN KEY(snapshot_id,namespace_uuid) REFERENCES mst2_snapshot_storage_route(snapshot_id,namespace_uuid), + FOREIGN KEY(snapshot_id,prepare_id,metadata_root) REFERENCES mst2_snapshot_context(snapshot_id,prepare_id,metadata_root), + UNIQUE(snapshot_id,namespace_uuid,session_incarnation,prepare_id,metadata_root) +); +CREATE TABLE mst2_lease_storage_route ( + lease_id text PRIMARY KEY, + snapshot_id text NOT NULL, + namespace_uuid uuid NOT NULL, + session_incarnation uuid NOT NULL, + prepare_id text NOT NULL, + metadata_root bytea NOT NULL CHECK (octet_length(metadata_root)=32), + authorization_epoch bigint NOT NULL, + publication_sequence bigint NOT NULL, + writer_epoch bigint NOT NULL, + certificate_receipt_id bigint NOT NULL, + FOREIGN KEY(snapshot_id,namespace_uuid) REFERENCES mst2_snapshot_storage_route(snapshot_id,namespace_uuid) +); +CREATE INDEX mst2_lease_storage_route_incarnation ON mst2_lease_storage_route(namespace_uuid,session_incarnation,lease_id); + +CREATE FUNCTION mst2_route_scope_valid(caller_schema text) RETURNS boolean LANGUAGE sql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ + SELECT NOT pg_catalog.pg_is_in_recovery() AND pg_catalog.current_setting('transaction_isolation')='read committed' + AND EXISTS(SELECT 1 FROM mst2_metadata_namespace n JOIN mst2_metadata_storage_scope s ON s.singleton=1 + JOIN pg_catalog.pg_database d ON d.datname=pg_catalog.current_database() + JOIN pg_catalog.pg_namespace c ON c.oid=n.core_schema_oid AND c.nspname=n.core_schema + WHERE n.singleton=1 AND caller_schema=n.core_schema AND n.core_schema=$CORE_LITERAL$ AND n.core_schema_oid=$CORE_OID$ + AND n.database_name=d.datname AND n.database_oid=d.oid AND n.storage_uuid=s.storage_uuid + AND n.server_address IS NOT DISTINCT FROM pg_catalog.inet_server_addr()::text + AND n.server_port IS NOT DISTINCT FROM pg_catalog.inet_server_port() + AND n.metadata_schema=n.core_schema AND n.metadata_schema_oid=n.core_schema_oid + AND n.graph_domain='generic-v1' AND n.family_identity='v3-generic-session-1' + AND n.admission_state='G_ADMITTED_Q_CLOSED' AND n.collector_state='CLOSED') +$$; + +CREATE FUNCTION mst2_route_lock_held(k1 integer,k2 integer,exclusive_only boolean DEFAULT true) RETURNS boolean LANGUAGE sql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ + SELECT EXISTS(SELECT 1 FROM pg_catalog.pg_locks WHERE locktype='advisory' AND pid=pg_catalog.pg_backend_pid() + AND database=(SELECT oid FROM pg_catalog.pg_database WHERE datname=pg_catalog.current_database()) + AND classid=k1::oid AND objid=k2::oid AND objsubid=2 AND granted AND (NOT exclusive_only OR mode='ExclusiveLock')) +$$; + +CREATE FUNCTION mst2_route_enter(caller_schema text) RETURNS uuid LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +DECLARE n mst2_metadata_namespace%ROWTYPE; mono_held boolean; route_held boolean; +BEGIN + IF NOT mst2_route_scope_valid(caller_schema) THEN RAISE EXCEPTION 'storage route captured primary scope is unavailable'; END IF; + SELECT * INTO STRICT n FROM mst2_metadata_namespace WHERE singleton=1; + mono_held:=mst2_route_lock_held(1297043024,n.mono_lock_key2); + route_held:=mst2_route_lock_held(1296718001,pg_catalog.hashtext(n.core_schema)); + IF mst2_route_lock_held(1296718001,pg_catalog.hashtext(n.core_schema),false) AND NOT mono_held THEN + RAISE EXCEPTION 'storage route lock was acquired before core mono'; + END IF; + IF EXISTS(SELECT 1 FROM mst2_metadata_namespace x + WHERE mst2_route_lock_held(1296717362,pg_catalog.hashtext(x.metadata_schema),false)) + AND NOT (mono_held AND route_held) THEN + RAISE EXCEPTION 'storage route cannot acquire core locks after retention'; + END IF; + PERFORM pg_catalog.set_config('lock_timeout','5000ms',true); + PERFORM pg_catalog.pg_advisory_xact_lock(1297043024,n.mono_lock_key2); + PERFORM pg_catalog.pg_advisory_xact_lock(1296718001,pg_catalog.hashtext(n.core_schema)); + FOR n IN SELECT * FROM mst2_metadata_namespace ORDER BY namespace_uuid LOOP + PERFORM pg_catalog.pg_advisory_xact_lock(1296717362,pg_catalog.hashtext(n.metadata_schema)); + END LOOP; + RETURN n.namespace_uuid; +END $$; + +-- This wrapper observes the caller before entering a fixed trusted path. It +-- uses only pg_catalog and the explicitly captured core function/relation. +CREATE FUNCTION mst2_route_statement_barrier() RETURNS trigger LANGUAGE plpgsql VOLATILE AS $$ +BEGIN + IF TG_TABLE_SCHEMA<>$CORE_LITERAL$ OR NOT EXISTS(SELECT 1 FROM pg_catalog.pg_class + WHERE oid=TG_RELID AND relnamespace=$CORE_OID$) THEN + RAISE EXCEPTION 'storage route mutation is outside its captured core schema'; + END IF; + PERFORM $CORE_SCHEMA$.mst2_route_enter(pg_catalog.current_schema()); + RETURN NULL; +END $$; + +CREATE FUNCTION mst2_route_profile(source_domain text,tagged_root_tree_oid text,scope text,schema_version smallint, + metadata_codec smallint,materialization_policy smallint,fs_semantics smallint,access_projection smallint, + verification_revision integer,projection_revision smallint) RETURNS jsonb LANGUAGE sql IMMUTABLE STRICT +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ + SELECT pg_catalog.jsonb_build_object('source_domain',source_domain,'tagged_root_tree_oid',tagged_root_tree_oid, + 'scope',scope,'schema_version',schema_version,'metadata_codec',metadata_codec, + 'materialization_policy',materialization_policy,'fs_semantics',fs_semantics, + 'access_projection',access_projection,'verification_revision',verification_revision, + 'projection_revision',projection_revision) +$$; + +CREATE FUNCTION mst2_route_snapshot_proof(sid text) RETURNS boolean LANGUAGE sql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ + SELECT EXISTS(SELECT 1 FROM mst2_snapshot_storage_route r JOIN mst2_snapshot_context s USING(snapshot_id) + JOIN mst2_generic_session_storage_binding b USING(snapshot_id,namespace_uuid) + JOIN mst2_metadata_namespace n USING(namespace_uuid) + JOIN mst2_metadata_prepare p ON p.prepare_id=s.prepare_id + WHERE r.snapshot_id=sid AND r.canonical_descriptor=s.canonical_descriptor + AND r.instance_id=s.instance_id AND r.commit_oid=s.commit_oid AND r.root_tree_oid=s.root_tree_oid + AND r.metadata_root=s.metadata_root AND r.source_profile=mst2_route_profile(p.source_domain,p.tagged_root_tree_oid, + p.scope,p.schema_version,p.metadata_codec,p.materialization_policy,p.fs_semantics,p.access_projection, + p.verification_revision,p.projection_revision) + AND b.prepare_id=s.prepare_id AND b.metadata_root=s.metadata_root + AND p.state='COMMITTED' AND p.metadata_root=s.metadata_root AND p.source_domain='native-git' + AND p.tagged_root_tree_oid IN ('sha1:'||s.root_tree_oid,'sha256:'||s.root_tree_oid) + AND n.graph_domain='generic-v1') +$$; + +CREATE FUNCTION mst2_route_lease_proof(lid text) RETURNS boolean LANGUAGE sql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ + SELECT EXISTS(SELECT 1 FROM mst2_lease_storage_route r JOIN mst2_snapshot_lease l USING(lease_id) + JOIN mst2_generic_session_storage_binding b ON b.snapshot_id=r.snapshot_id AND b.namespace_uuid=r.namespace_uuid + AND b.session_incarnation=r.session_incarnation AND b.prepare_id=r.prepare_id AND b.metadata_root=r.metadata_root + WHERE r.lease_id=lid AND r.snapshot_id=l.snapshot_id AND r.authorization_epoch=l.authorization_epoch + AND r.publication_sequence=l.publication_sequence AND r.writer_epoch=l.writer_epoch + AND r.certificate_receipt_id=l.certificate_receipt_id AND mst2_route_snapshot_proof(r.snapshot_id)) +$$; + +CREATE FUNCTION mst2_route_select_snapshot(sid text,caller_schema text) +RETURNS TABLE(context_present boolean,route_present boolean,valid boolean) LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + IF NOT mst2_route_scope_valid(caller_schema) THEN RAISE EXCEPTION 'storage route captured primary scope is unavailable'; END IF; + RETURN QUERY SELECT EXISTS(SELECT 1 FROM mst2_snapshot_context WHERE snapshot_id=sid), + EXISTS(SELECT 1 FROM mst2_snapshot_storage_route WHERE snapshot_id=sid),mst2_route_snapshot_proof(sid); +END $$; + +CREATE FUNCTION mst2_route_select_lease(lid text,caller_schema text) +RETURNS TABLE(actual_sid text,route_present boolean,valid boolean) LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + IF NOT mst2_route_scope_valid(caller_schema) THEN RAISE EXCEPTION 'storage route captured primary scope is unavailable'; END IF; + RETURN QUERY SELECT (SELECT snapshot_id FROM mst2_snapshot_lease WHERE lease_id=lid), + EXISTS(SELECT 1 FROM mst2_lease_storage_route WHERE lease_id=lid),mst2_route_lease_proof(lid); +END $$; + +CREATE FUNCTION mst2_route_immutable() RETURNS trigger LANGUAGE plpgsql +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN RAISE EXCEPTION 'storage route identity and historical ledger are immutable'; END $$; + +CREATE FUNCTION mst2_route_insert_guard() RETURNS trigger LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + IF TG_TABLE_NAME='mst2_snapshot_storage_route' THEN + IF NOT EXISTS(SELECT 1 FROM mst2_snapshot_context s JOIN mst2_metadata_prepare p ON p.prepare_id=s.prepare_id + JOIN mst2_metadata_namespace n ON n.namespace_uuid=NEW.namespace_uuid + WHERE s.snapshot_id=NEW.snapshot_id AND NEW.canonical_descriptor=s.canonical_descriptor + AND NEW.instance_id=s.instance_id AND NEW.commit_oid=s.commit_oid AND NEW.root_tree_oid=s.root_tree_oid + AND NEW.metadata_root=s.metadata_root AND NEW.source_profile=mst2_route_profile(p.source_domain,p.tagged_root_tree_oid, + p.scope,p.schema_version,p.metadata_codec,p.materialization_policy,p.fs_semantics,p.access_projection, + p.verification_revision,p.projection_revision) + AND p.state='COMMITTED' AND p.metadata_root=s.metadata_root AND p.source_domain='native-git' + AND (p.graph_domain IS NULL OR p.graph_domain='generic-v1') + AND p.tagged_root_tree_oid IN ('sha1:'||s.root_tree_oid,'sha256:'||s.root_tree_oid)) THEN + RAISE EXCEPTION 'storage route is not derived from its actual generic context'; + END IF; + ELSIF TG_TABLE_NAME='mst2_generic_session_storage_binding' THEN + IF NOT EXISTS(SELECT 1 FROM mst2_snapshot_context s JOIN mst2_snapshot_storage_route r USING(snapshot_id) + WHERE s.snapshot_id=NEW.snapshot_id AND r.namespace_uuid=NEW.namespace_uuid + AND s.prepare_id=NEW.prepare_id AND s.metadata_root=NEW.metadata_root) THEN + RAISE EXCEPTION 'storage route generic incarnation is not its actual context'; + END IF; + ELSIF TG_TABLE_NAME='mst2_lease_storage_route' THEN + IF NOT EXISTS(SELECT 1 FROM mst2_snapshot_lease l JOIN mst2_generic_session_storage_binding b USING(snapshot_id) + WHERE l.lease_id=NEW.lease_id AND l.snapshot_id=NEW.snapshot_id AND b.namespace_uuid=NEW.namespace_uuid + AND b.session_incarnation=NEW.session_incarnation AND b.prepare_id=NEW.prepare_id AND b.metadata_root=NEW.metadata_root + AND l.authorization_epoch=NEW.authorization_epoch AND l.publication_sequence=NEW.publication_sequence + AND l.writer_epoch=NEW.writer_epoch AND l.certificate_receipt_id=NEW.certificate_receipt_id + AND mst2_route_snapshot_proof(l.snapshot_id)) THEN + RAISE EXCEPTION 'storage route lease is not its exact generic incarnation and source'; + END IF; + ELSE RAISE EXCEPTION 'storage route insertion target is not registered'; END IF; + RETURN NEW; +END $$; + +CREATE FUNCTION mst2_route_context_insert() RETURNS trigger LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + INSERT INTO mst2_snapshot_storage_route + SELECT s.snapshot_id,n.namespace_uuid,s.canonical_descriptor,s.instance_id,s.commit_oid,s.root_tree_oid, + s.metadata_root,mst2_route_profile(p.source_domain,p.tagged_root_tree_oid,p.scope,p.schema_version,p.metadata_codec, + p.materialization_policy,p.fs_semantics,p.access_projection,p.verification_revision,p.projection_revision) + FROM mst2_snapshot_context s JOIN mst2_metadata_prepare p ON p.prepare_id=s.prepare_id + CROSS JOIN mst2_metadata_namespace n WHERE s.snapshot_id=NEW.snapshot_id AND n.singleton=1 + ON CONFLICT(snapshot_id) DO NOTHING; + INSERT INTO mst2_generic_session_storage_binding + SELECT pg_catalog.gen_random_uuid(),s.snapshot_id,r.namespace_uuid,s.prepare_id,s.metadata_root + FROM mst2_snapshot_context s JOIN mst2_snapshot_storage_route r USING(snapshot_id) WHERE s.snapshot_id=NEW.snapshot_id + ON CONFLICT(snapshot_id) DO NOTHING; + RETURN NULL; +END $$; +CREATE FUNCTION mst2_route_lease_insert() RETURNS trigger LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + INSERT INTO mst2_lease_storage_route + SELECT l.lease_id,l.snapshot_id,b.namespace_uuid,b.session_incarnation,b.prepare_id,b.metadata_root, + l.authorization_epoch,l.publication_sequence,l.writer_epoch,l.certificate_receipt_id + FROM mst2_snapshot_lease l JOIN mst2_generic_session_storage_binding b USING(snapshot_id) WHERE l.lease_id=NEW.lease_id + ON CONFLICT(lease_id) DO NOTHING; + RETURN NULL; +END $$; +CREATE FUNCTION mst2_route_complete() RETURNS trigger LANGUAGE plpgsql VOLATILE +SET search_path=$CORE_SCHEMA$,pg_catalog,pg_temp AS $$ +BEGIN + IF TG_TABLE_NAME='mst2_snapshot_context' THEN + IF NOT mst2_route_snapshot_proof(NEW.snapshot_id) THEN RAISE EXCEPTION 'storage route context committed without exact routing'; END IF; + ELSIF NOT mst2_route_lease_proof(NEW.lease_id) THEN RAISE EXCEPTION 'storage route lease committed without exact routing'; END IF; + RETURN NULL; +END $$; + +DO $$ DECLARE t text; BEGIN + FOREACH t IN ARRAY ARRAY['mst2_metadata_namespace','mst2_snapshot_storage_route', + 'mst2_generic_session_storage_binding','mst2_lease_storage_route'] LOOP + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_00_route_statement_barrier BEFORE INSERT OR UPDATE OR DELETE ON %I FOR EACH STATEMENT EXECUTE FUNCTION mst2_route_statement_barrier()',t); + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_route_immutable BEFORE UPDATE OR DELETE ON %I FOR EACH ROW EXECUTE FUNCTION mst2_route_immutable()',t); + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_route_truncate_guard BEFORE TRUNCATE ON %I FOR EACH STATEMENT EXECUTE FUNCTION mst2_route_immutable()',t); + IF t<>'mst2_metadata_namespace' THEN + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_route_insert_guard BEFORE INSERT ON %I FOR EACH ROW EXECUTE FUNCTION mst2_route_insert_guard()',t); + END IF; + END LOOP; + FOREACH t IN ARRAY ARRAY['mst2_snapshot_context','mst2_snapshot_lease'] LOOP + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_00_route_statement_barrier BEFORE INSERT ON %I FOR EACH STATEMENT EXECUTE FUNCTION mst2_route_statement_barrier()',t); + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_route_source_delete_guard BEFORE DELETE ON %I FOR EACH ROW EXECUTE FUNCTION mst2_route_immutable()',t); + EXECUTE pg_catalog.format('CREATE TRIGGER mst2_route_source_truncate_guard BEFORE TRUNCATE ON %I FOR EACH STATEMENT EXECUTE FUNCTION mst2_route_immutable()',t); + EXECUTE pg_catalog.format('CREATE CONSTRAINT TRIGGER mst2_route_complete AFTER INSERT ON %I DEFERRABLE INITIALLY DEFERRED FOR EACH ROW EXECUTE FUNCTION mst2_route_complete()',t); + END LOOP; +END $$; +CREATE TRIGGER mst2_route_namespace_registration_closed BEFORE INSERT ON mst2_metadata_namespace + FOR EACH ROW EXECUTE FUNCTION mst2_route_immutable(); +CREATE TRIGGER mst2_route_context_insert AFTER INSERT ON mst2_snapshot_context + FOR EACH ROW EXECUTE FUNCTION mst2_route_context_insert(); +CREATE TRIGGER mst2_route_lease_insert AFTER INSERT ON mst2_snapshot_lease + FOR EACH ROW EXECUTE FUNCTION mst2_route_lease_insert(); + +INSERT INTO mst2_snapshot_storage_route +SELECT s.snapshot_id,n.namespace_uuid,s.canonical_descriptor,s.instance_id,s.commit_oid,s.root_tree_oid, + s.metadata_root,mst2_route_profile(p.source_domain,p.tagged_root_tree_oid,p.scope,p.schema_version,p.metadata_codec, + p.materialization_policy,p.fs_semantics,p.access_projection,p.verification_revision,p.projection_revision) +FROM mst2_snapshot_context s JOIN mst2_metadata_prepare p ON p.prepare_id=s.prepare_id +CROSS JOIN mst2_metadata_namespace n; +INSERT INTO mst2_generic_session_storage_binding +SELECT pg_catalog.gen_random_uuid(),s.snapshot_id,r.namespace_uuid,s.prepare_id,s.metadata_root +FROM mst2_snapshot_context s JOIN mst2_snapshot_storage_route r USING(snapshot_id); +INSERT INTO mst2_lease_storage_route +SELECT l.lease_id,l.snapshot_id,b.namespace_uuid,b.session_incarnation,b.prepare_id,b.metadata_root, + l.authorization_epoch,l.publication_sequence,l.writer_epoch,l.certificate_receipt_id +FROM mst2_snapshot_lease l JOIN mst2_generic_session_storage_binding b USING(snapshot_id); +DO $$ BEGIN + IF NOT mst2_route_scope_valid($CORE_LITERAL$) + OR EXISTS(SELECT 1 FROM mst2_snapshot_context s WHERE NOT mst2_route_snapshot_proof(s.snapshot_id)) + OR EXISTS(SELECT 1 FROM mst2_snapshot_lease l WHERE NOT mst2_route_lease_proof(l.lease_id)) THEN + RAISE EXCEPTION 'storage route backfill did not cover the complete durable inventory'; + END IF; +END $$; diff --git a/src/jupiter/migration/mod.rs b/src/jupiter/migration/mod.rs index fa611fce..2d667a39 100644 --- a/src/jupiter/migration/mod.rs +++ b/src/jupiter/migration/mod.rs @@ -150,6 +150,7 @@ mod m20261007_000200_add_mst2_metadata_generations; mod m20261007_000300_add_mst2_metadata_lifetime_history; mod m20261007_000400_add_mst2_qualified_metadata_gc; mod m20261007_000500_add_mst2_install_capability; +mod m20261007_000600_add_mst2_storage_routes; mod runner; pub use m20260905_000100_add_push_queue::ensure_queue_control_seed; pub use runner::apply_migrations; @@ -288,6 +289,7 @@ impl MigratorTrait for Migrator { Box::new(m20261007_000300_add_mst2_metadata_lifetime_history::Migration), Box::new(m20261007_000400_add_mst2_qualified_metadata_gc::Migration), Box::new(m20261007_000500_add_mst2_install_capability::Migration), + Box::new(m20261007_000600_add_mst2_storage_routes::Migration), ] } } @@ -1198,7 +1200,7 @@ mod tests { async fn import_repo_alias_rows_canonicalized() { let names = migration_names(); assert_eq!( - &names[names.len() - 13..names.len() - 4], + &names[names.len() - 14..names.len() - 5], &[ "m20260923_000200_canonicalize_import_repo_paths".to_string(), "m20260925_000100_media_paging".to_string(), @@ -1213,21 +1215,25 @@ mod tests { "native retention, metadata installation and view tables follow media paging" ); assert_eq!( - &names[names.len() - 4], + &names[names.len() - 5], "m20261007_000200_add_mst2_metadata_generations" ); assert_eq!( - &names[names.len() - 3], + &names[names.len() - 4], "m20261007_000300_add_mst2_metadata_lifetime_history" ); assert_eq!( - &names[names.len() - 2], + &names[names.len() - 3], "m20261007_000400_add_mst2_qualified_metadata_gc" ); assert_eq!( - names.last().unwrap(), + &names[names.len() - 2], "m20261007_000500_add_mst2_install_capability" ); + assert_eq!( + names.last().unwrap(), + "m20261007_000600_add_mst2_storage_routes" + ); let db = alias_db().await; insert_repo(&db, 1, "/third-party//a").await; diff --git a/src/jupiter/storage/native_metadata_install.rs b/src/jupiter/storage/native_metadata_install.rs index 735dcc54..6a1c8500 100644 --- a/src/jupiter/storage/native_metadata_install.rs +++ b/src/jupiter/storage/native_metadata_install.rs @@ -163,6 +163,10 @@ impl PostgresMetadataInstallRepository { }) } + pub(crate) fn captured_schema(&self) -> &str { + &self.storage_scope.schema + } + /// These columns are returned by the same query that authorizes a lease, /// so a warm cache cannot authorize a stale replica or another schema. pub(crate) fn verify_primary_scope_row(&self, row: &QueryResult) -> Result<(), SnapshotError> { diff --git a/src/jupiter/storage/native_snapshot_routes.rs b/src/jupiter/storage/native_snapshot_routes.rs new file mode 100644 index 00000000..ac818034 --- /dev/null +++ b/src/jupiter/storage/native_snapshot_routes.rs @@ -0,0 +1,122 @@ +//! Permanent namespace selection precedes the generic physical session proof. + +use super::{ + ConnectionTrait, DatabaseTransaction, PostgresMetadataInstallRepository, SESSION_SQL, + SnapshotError, integrity, internal, statement, +}; + +fn function(installer: &PostgresMetadataInstallRepository, name: &str) -> String { + format!( + "\"{}\".{name}", + installer.captured_schema().replace('"', "\"\"") + ) +} + +pub(super) async fn enter( + txn: &DatabaseTransaction, + installer: &PostgresMetadataInstallRepository, +) -> Result<(), SnapshotError> { + installer.verify_primary_connection(txn).await?; + txn.execute_raw(statement( + &format!( + "SELECT {}(current_schema())", + function(installer, "mst2_route_enter") + ), + [], + )) + .await + .map_err(internal)?; + generic_path(txn, installer).await?; + Ok(()) +} + +pub(super) async fn generic_path( + txn: &DatabaseTransaction, + installer: &PostgresMetadataInstallRepository, +) -> Result<(), SnapshotError> { + let schema = installer.captured_schema().replace('"', "\"\""); + txn.execute_unprepared(&format!( + "SET LOCAL search_path=\"{schema}\",pg_catalog,pg_temp" + )) + .await + .map_err(internal)?; + Ok(()) +} + +pub(super) fn session_sql(installer: &PostgresMetadataInstallRepository) -> String { + let schema = installer.captured_schema().replace('"', "\"\""); + let mut sql = SESSION_SQL.to_owned(); + for table in [ + "mst2_metadata_storage_scope", + "mst2_snapshot_context", + "mst2_snapshot_lease", + "mst2_retention_node", + "mst2_retention_root", + "mst2_metadata_prepare", + "mst2_native_publication", + "mst2_publication", + "mst2_publication_outbox", + ] { + for join in ["FROM", "JOIN"] { + sql = sql.replace( + &format!("{join} {table} "), + &format!("{join} \"{schema}\".{table} "), + ); + } + } + sql +} + +pub(super) async fn snapshot( + connection: &C, + installer: &PostgresMetadataInstallRepository, + sid: &str, +) -> Result<(), SnapshotError> { + let row = connection + .query_one_raw(statement( + &format!( + "SELECT * FROM {}($1,current_schema())", + function(installer, "mst2_route_select_snapshot") + ), + [sid.into()], + )) + .await + .map_err(internal)? + .ok_or_else(|| integrity("snapshot storage route selection is missing"))?; + let context: bool = row.try_get("", "context_present").map_err(internal)?; + let route: bool = row.try_get("", "route_present").map_err(internal)?; + let valid: bool = row.try_get("", "valid").map_err(internal)?; + if context != route || context && !valid { + return Err(integrity( + "snapshot storage route conflicts with its immutable generic session", + )); + } + Ok(()) +} + +pub(super) async fn lease( + connection: &C, + installer: &PostgresMetadataInstallRepository, + lease_id: &str, +) -> Result, SnapshotError> { + let row = connection + .query_one_raw(statement( + &format!( + "SELECT * FROM {}($1,current_schema())", + function(installer, "mst2_route_select_lease") + ), + [lease_id.into()], + )) + .await + .map_err(internal)? + .ok_or_else(|| integrity("lease storage route selection is missing"))?; + let sid: Option = row.try_get("", "actual_sid").map_err(internal)?; + let route: bool = row.try_get("", "route_present").map_err(internal)?; + let valid: bool = row.try_get("", "valid").map_err(internal)?; + if sid.is_some() != route || sid.is_some() && !valid { + return Err(integrity( + "lease storage route conflicts with its exact generic incarnation", + )); + } + Ok(sid) +} diff --git a/src/jupiter/storage/native_snapshot_session.rs b/src/jupiter/storage/native_snapshot_session.rs index d84d075a..bb374236 100644 --- a/src/jupiter/storage/native_snapshot_session.rs +++ b/src/jupiter/storage/native_snapshot_session.rs @@ -17,7 +17,6 @@ use super::{ MetadataInstallError, PostgresMetadataInstallRepository, PreparedMetadataReceipt, }, native_publication_storage::{NativePublicationHead, decode_native_observation}, - push_queue_storage::PushQueueStorage, }; use crate::{ callisto::mst2_snapshot_context, @@ -34,6 +33,7 @@ use crate::{ pub(crate) struct PostgresNativeSessionRepository { connection: DatabaseConnection, installer: OnceCell, + qualified_session_sql: OnceCell, restored: Mutex>>>, } @@ -42,6 +42,7 @@ impl PostgresNativeSessionRepository { Self { connection, installer: OnceCell::new(), + qualified_session_sql: OnceCell::new(), restored: Mutex::new(HashMap::new()), } } @@ -52,6 +53,12 @@ impl PostgresNativeSessionRepository { .await } + async fn session_sql(&self, installer: &PostgresMetadataInstallRepository) -> &str { + self.qualified_session_sql + .get_or_init(|| async { routes::session_sql(installer) }) + .await + } + pub(crate) async fn install( &self, built: &BuiltDescriptor, @@ -95,14 +102,14 @@ impl PostgresNativeSessionRepository { let installer = self.installer().await?; let txn = self.transaction().await?; let result = async { - PushQueueStorage::acquire_mono_write_lock(&txn).await.map_err(internal)?; - installer.verify_primary_connection(&txn).await?; + routes::enter(&txn, installer).await?; let current = MonoStorage::read_native_publication_head_from(&txn, &expected.instance_id) .await.map_err(|_| not_ready("native publication is not ready"))?; if current.root != expected.root || current.token != expected.token { return Err(not_ready("native publication advanced during preparation; retry resolve")); } retention_lock(&txn).await?; + routes::snapshot(&txn, installer, &built.snapshot_id).await?; let existing = mst2_snapshot_context::Entity::find_by_id(built.snapshot_id.clone()) .one(&txn).await.map_err(internal)?; let prepare_id = if let Some(existing) = existing { @@ -175,10 +182,17 @@ impl PostgresNativeSessionRepository { instance_id: &str, ) -> Result { let installer = self.installer().await?; + if routes::lease(&self.connection, installer, lease_id) + .await? + .as_deref() + != Some(snapshot_id) + { + return Err(expired()); + } let row = self .connection .query_one_raw(statement( - SESSION_SQL, + self.session_sql(installer).await, [snapshot_id.into(), lease_id.into()], )) .await @@ -191,6 +205,9 @@ impl PostgresNativeSessionRepository { if error.code == SnapshotErrorCode::LeaseExpired { let txn = self.transaction().await?; let result = async { + installer.verify_primary_connection(&txn).await?; + routes::lease(&txn, installer, lease_id).await?; + routes::generic_path(&txn, installer).await?; retention_lock(&txn).await?; expire_specific_locked(&txn, lease_id).await } @@ -218,15 +235,23 @@ impl PostgresNativeSessionRepository { seconds: u64, instance: &str, ) -> Result { + let installer = self.installer().await?; let txn = self.transaction().await?; let result = async { + installer.verify_primary_connection(&txn).await?; + let selected_sid = routes::lease(&txn, installer, lease_id).await? + .ok_or_else(|| SnapshotError::new(SnapshotErrorCode::LeaseUnknown,"unknown lease_id"))?; + routes::generic_path(&txn, installer).await?; retention_lock(&txn).await?; let lease = txn.query_one_raw(statement( "SELECT snapshot_id FROM mst2_snapshot_lease WHERE lease_id=$1 FOR UPDATE", [lease_id.into()], )).await.map_err(internal)?.ok_or_else(|| SnapshotError::new(SnapshotErrorCode::LeaseUnknown,"unknown lease_id"))?; let sid: String = lease.try_get("","snapshot_id").map_err(internal)?; - let row = txn.query_one_raw(statement(SESSION_SQL,[sid.clone().into(),lease_id.into()])) + if sid != selected_sid { + return Err(integrity("lease changed after immutable route selection")); + } + let row = txn.query_one_raw(statement(self.session_sql(installer).await,[sid.clone().into(),lease_id.into()])) .await.map_err(internal)?.ok_or_else(expired)?; self.installer().await?.verify_primary_scope_row(&row)?; if let Err(error)=decode_context(&row,&sid,lease_id,instance) { @@ -257,8 +282,12 @@ impl PostgresNativeSessionRepository { let installer = self.installer().await?; let txn = self.transaction().await?; let result = async { - retention_lock(&txn).await?; installer.verify_primary_connection(&txn).await?; + if routes::lease(&txn, installer, lease_id).await?.is_none() { + return Ok(false); + } + routes::generic_path(&txn, installer).await?; + retention_lock(&txn).await?; let changed=txn.query_one_raw(statement( "UPDATE mst2_snapshot_lease SET state='RELEASED' WHERE lease_id=$1 AND state='ACTIVE' RETURNING snapshot_id", [lease_id.into()], @@ -293,6 +322,9 @@ impl PostgresNativeSessionRepository { } } +#[path = "native_snapshot_routes.rs"] +mod routes; + 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,