From c33ed3cfc19d135029c030a9db48d7af1c3f1665 Mon Sep 17 00:00:00 2001 From: Xiaoyang Han Date: Thu, 8 Oct 2026 06:59:58 +0800 Subject: [PATCH] fix(v3): read metadata primary scope from its captured permanent schema A TEMP table with the same name could poison native metadata repository initialization. Capture the actual permanent schema through pg_catalog, qualify the storage identity read during initialization and mutation/recovery barriers, and continue observing and comparing the actual current database/schema/OIDs/server/replica identity. A matching fake TEMP UUID cannot authorize another primary schema. Add a single-connection regression with TEMP present before construction, primary schema switching and restoration. Runtime verification query counts and existing integrity assertions are preserved. --- .../storage/native_metadata_install.rs | 48 +++++++-- .../storage/native_metadata_install_tests.rs | 101 ++++++++++++++++++ 2 files changed, 140 insertions(+), 9 deletions(-) diff --git a/src/jupiter/storage/native_metadata_install.rs b/src/jupiter/storage/native_metadata_install.rs index 6ecf72c3..7d978fa0 100644 --- a/src/jupiter/storage/native_metadata_install.rs +++ b/src/jupiter/storage/native_metadata_install.rs @@ -156,7 +156,13 @@ struct PrimaryStorageScope { impl PostgresMetadataInstallRepository { /// The connection must target the deployment's primary PostgreSQL database. pub async fn new(connection: DatabaseConnection) -> Result { - let storage_scope = read_storage_scope(&connection).await?; + let schema = capture_storage_schema(&connection).await?; + let storage_scope = read_storage_scope(&connection, &schema).await?; + if storage_scope.schema != schema { + return Err(internal( + "metadata primary schema changed during storage scope capture", + )); + } Ok(Self { connection, barrier_timeout: Duration::from_secs(5), @@ -204,7 +210,7 @@ impl PostgresMetadataInstallRepository { &self, connection: &C, ) -> Result<(), SnapshotError> { - if read_storage_scope(connection).await? != self.storage_scope { + if read_storage_scope(connection, self.captured_schema()).await? != self.storage_scope { return Err(internal( "session mutation no longer targets its captured primary storage scope", )); @@ -659,7 +665,7 @@ impl PostgresMetadataInstallRepository { if let Some(family) = &self.qualified_family { family.enter(txn).await?; } - if read_storage_scope(txn).await? != self.storage_scope { + if read_storage_scope(txn, self.captured_schema()).await? != self.storage_scope { return Err(internal( "metadata recovery connection is outside the captured primary storage scope", )); @@ -703,20 +709,44 @@ impl PostgresMetadataInstallRepository { } } +async fn capture_storage_schema( + connection: &C, +) -> Result { + if connection.get_database_backend() != DbBackend::Postgres { + return Err(internal("metadata storage scope requires PostgreSQL")); + } + let row = connection + .query_one_raw(statement( + "SELECT n.nspname AS schema FROM pg_catalog.pg_namespace n + JOIN pg_catalog.pg_class c ON c.relnamespace=n.oid + AND c.relname='mst2_metadata_storage_scope' AND c.relkind='r' + WHERE n.nspname=pg_catalog.current_schema() AND n.nspname NOT LIKE 'pg_temp_%'", + [], + )) + .await + .map_err(internal)? + .ok_or_else(|| internal("metadata actual primary storage relation is missing"))?; + row.try_get("", "schema").map_err(internal) +} + async fn read_storage_scope( connection: &C, + schema: &str, ) -> Result { if connection.get_database_backend() != DbBackend::Postgres { return Err(internal("metadata storage scope requires PostgreSQL")); } let row = connection .query_one_raw(statement( - "SELECT s.storage_uuid, current_database() AS database, d.oid::bigint AS database_oid, - current_schema() AS schema, n.oid::bigint AS schema_oid, - inet_server_addr()::text AS server_address, inet_server_port() AS server_port, - pg_is_in_recovery() AS replica FROM mst2_metadata_storage_scope s - JOIN pg_catalog.pg_database d ON d.datname=current_database() - JOIN pg_catalog.pg_namespace n ON n.nspname=current_schema() WHERE s.singleton=1", + &format!( + "SELECT s.storage_uuid, pg_catalog.current_database() AS database, d.oid::bigint AS database_oid, + pg_catalog.current_schema() AS schema, n.oid::bigint AS schema_oid, + pg_catalog.inet_server_addr()::text AS server_address, pg_catalog.inet_server_port() AS server_port, + pg_catalog.pg_is_in_recovery() AS replica FROM \"{}\".mst2_metadata_storage_scope s + JOIN pg_catalog.pg_database d ON d.datname=pg_catalog.current_database() + JOIN pg_catalog.pg_namespace n ON n.nspname=pg_catalog.current_schema() WHERE s.singleton=1", + schema.replace('"', "\"\"") + ), [], )) .await diff --git a/src/jupiter/storage/native_metadata_install_tests.rs b/src/jupiter/storage/native_metadata_install_tests.rs index 4145afd8..4cb387dc 100644 --- a/src/jupiter/storage/native_metadata_install_tests.rs +++ b/src/jupiter/storage/native_metadata_install_tests.rs @@ -70,6 +70,107 @@ fn rejected(error: MetadataInstallError) -> SnapshotErrorCode { } } +#[tokio::test] +async fn native_metadata_scope_ignores_poisoned_temp_table_and_rejects_changed_primary_schema() { + let (first, _second, _schema, url) = fixture().await; + let (_other, _unused, other_schema, _other_url) = fixture().await; + let expected = PostgresMetadataInstallRepository::new(first) + .await + .unwrap() + .storage_scope; + let mut options = sea_orm::ConnectOptions::new(url); + options.max_connections(1).min_connections(1); + let connection = Database::connect(options).await.unwrap(); + connection + .execute_unprepared( + "CREATE TEMP TABLE mst2_metadata_storage_scope(singleton integer,storage_uuid text); + INSERT INTO pg_temp.mst2_metadata_storage_scope VALUES(1,'temp-poison')", + ) + .await + .unwrap(); + let repository = PostgresMetadataInstallRepository::new(connection.clone()) + .await + .unwrap(); + assert_eq!(repository.storage_scope, expected); + assert_eq!(repository.captured_schema(), _schema.schema()); + repository + .verify_primary_connection(&connection) + .await + .unwrap(); + let txn = connection.begin().await.unwrap(); + repository.barrier(&txn).await.unwrap(); + txn.rollback().await.unwrap(); + connection + .execute_raw(statement( + "UPDATE pg_temp.mst2_metadata_storage_scope SET storage_uuid=$1", + [repository.storage_scope.storage_uuid.clone().into()], + )) + .await + .unwrap(); + let txn = connection.begin().await.unwrap(); + txn.execute_unprepared(&format!( + "SET LOCAL search_path=\"{}\",pg_catalog,pg_temp", + other_schema.schema().replace('"', "\"\"") + )) + .await + .unwrap(); + let error = repository + .verify_primary_connection(&txn) + .await + .unwrap_err(); + assert_eq!(error.code, SnapshotErrorCode::Internal); + assert!( + error + .message + .contains("session mutation no longer targets its captured primary storage scope") + ); + let error = repository.barrier(&txn).await.unwrap_err(); + assert_eq!(error.code, SnapshotErrorCode::Internal); + assert!( + error + .message + .contains("metadata recovery connection is outside the captured primary storage scope") + ); + txn.rollback().await.unwrap(); + connection + .execute_unprepared("SET search_path=pg_temp,pg_catalog") + .await + .unwrap(); + let error = PostgresMetadataInstallRepository::new(connection.clone()) + .await + .err() + .unwrap(); + assert_eq!(error.code, SnapshotErrorCode::Internal); + assert_eq!( + error.message, + "metadata actual primary storage relation is missing" + ); + connection + .execute_unprepared(&format!( + "SET search_path=\"{}\",pg_catalog,pg_temp", + repository.captured_schema().replace('"', "\"\"") + )) + .await + .unwrap(); + repository + .verify_primary_connection(&connection) + .await + .unwrap(); + assert_eq!( + connection + .query_one_raw(statement( + "SELECT storage_uuid FROM pg_temp.mst2_metadata_storage_scope WHERE singleton=1", + [], + )) + .await + .unwrap() + .unwrap() + .try_get::("", "storage_uuid") + .unwrap(), + expected.storage_uuid + ); +} + #[tokio::test] async fn native_metadata_two_primary_connections_replay_one_receipt_without_recounting() { let (first, second, _schema, _url) = fixture().await;