Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
48 changes: 39 additions & 9 deletions src/jupiter/storage/native_metadata_install.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self, SnapshotError> {
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),
Expand Down Expand Up @@ -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",
));
Expand Down Expand Up @@ -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",
));
Expand Down Expand Up @@ -703,20 +709,44 @@ impl PostgresMetadataInstallRepository {
}
}

async fn capture_storage_schema<C: ConnectionTrait>(
connection: &C,
) -> Result<String, SnapshotError> {
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<C: ConnectionTrait>(
connection: &C,
schema: &str,
) -> Result<PrimaryStorageScope, SnapshotError> {
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
Expand Down
101 changes: 101 additions & 0 deletions src/jupiter/storage/native_metadata_install_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<String>("", "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;
Expand Down
Loading