use crate::Engine;
use std::{path::Path, sync::Arc};
use uqa_sql::catalog::{foreign_server::ForeignServerDefinition, roles::RoleIdentity};
pub(super) fn open(provider: usize, path: &Path) -> Engine {
match provider {
0 => Engine::new(),
1 => Engine::open(path).unwrap(),
2 => Engine::from_persistent_provider(Arc::new(
uqa_storage_sqlite::SQLiteKeyValueStorage::open(path).unwrap(),
))
.unwrap(),
3 => Engine::from_persistent_provider(Arc::new(
uqa_storage_redb::RedbStorage::open(path).unwrap(),
))
.unwrap(),
_ => unreachable!(),
}
}
fn sql(engine: &Engine, statement: &str) {
engine
.sql(statement, &[])
.unwrap_or_else(|error| panic!("{statement}: {error}"));
}
fn server(engine: &Engine, name: &str) -> ForeignServerDefinition {
engine.durable.foreign_servers.read()[name].clone()
}
#[test]
fn foreign_server_owner_and_identity_survive_rename_undo_refresh_and_reopen() {
for provider in 0..4 {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("servers.db");
let engine = open(provider, &path);
sql(&engine, "CREATE ROLE server_owner; SET ROLE server_owner; CREATE SERVER source TYPE '' VERSION '' FOREIGN DATA WRAPPER memory_fdw OPTIONS (metadata_json 'connection option'); RESET ROLE");
let expected = server(&engine, "source");
assert_eq!(expected.metadata.server_type.as_deref(), Some(""));
assert_eq!(expected.metadata.version.as_deref(), Some(""));
let owner = engine.durable.roles.read()["server_owner"].clone();
assert_eq!(
expected.metadata.owner,
RoleIdentity {
oid: owner.oid,
object_id: owner.object_id
}
);
assert_eq!(expected.options["metadata_json"], "connection option");
assert_eq!(
engine.foreign_server("source").unwrap().unwrap().options,
expected.options
);
let retained = engine.durable.foreign_servers.snapshot();
let pinned = engine.catalog_read_view();
sql(&engine, "ALTER ROLE server_owner RENAME TO renamed_owner; BEGIN; SAVEPOINT before_server; CREATE SERVER transient FOREIGN DATA WRAPPER memory_fdw");
assert_eq!(retained.len(), 1);
assert_eq!(pinned.snapshot().definitions.foreign_servers.len(), 1);
let transient = server(&engine, "transient");
assert_ne!(transient.metadata.oid, expected.metadata.oid);
assert_ne!(transient.metadata.object_id, expected.metadata.object_id);
sql(&engine, "ROLLBACK TO before_server; COMMIT");
assert!(engine.foreign_server("transient").unwrap().is_none());
assert_eq!(server(&engine, "source"), expected);
sql(
&engine,
"CREATE SERVER transient FOREIGN DATA WRAPPER memory_fdw",
);
assert_ne!(
server(&engine, "transient").metadata.object_id,
transient.metadata.object_id
);
assert_ne!(
server(&engine, "transient").metadata.oid,
transient.metadata.oid
);
assert_eq!(
engine
.sql("DROP ROLE renamed_owner", &[])
.unwrap_err()
.sqlstate(),
Some("2BP01")
);
if provider > 0 {
let peer = engine.new_session().unwrap();
sql(
&peer,
"CREATE SERVER peer_server FOREIGN DATA WRAPPER memory_fdw",
);
assert!(engine.foreign_server("peer_server").unwrap().is_some());
assert_eq!(server(&engine, "source"), expected);
drop(peer);
drop(pinned);
drop(engine);
let reopened = open(provider, &path);
assert_eq!(server(&reopened, "source"), expected);
assert_eq!(
reopened
.sql("DROP ROLE renamed_owner", &[])
.unwrap_err()
.sqlstate(),
Some("2BP01")
);
}
}
}
#[test]
fn foreign_server_initial_open_upgrades_legacy_metadata_once_for_each_persistent_provider() {
for provider in 1..4 {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("legacy-server.db");
let engine = open(provider, &path);
let factory = Arc::clone(engine.storage.provider.as_ref().unwrap());
let raw = factory.open_session().unwrap();
raw.backend.begin_transaction().unwrap();
raw.catalog
.delete_metadata("foreign-server-metadata-format")
.unwrap();
super::wrappers::remove_wrapper_format(raw.catalog.as_ref());
raw.catalog
.save_foreign_server("legacy", "memory_fdw", r#"{"key":"original"}"#)
.unwrap();
raw.backend.commit_transaction().unwrap();
assert!(
engine.new_session().is_err(),
"refresh must not migrate legacy metadata"
);
raw.catalog.set_metadata("sql_triggers_json", "{").unwrap();
drop(engine);
let failure = Engine::from_persistent_provider(Arc::clone(&factory))
.err()
.expect("malformed later catalog must reject open");
assert!(failure.to_string().contains("EOF"), "{failure}");
assert!(raw
.catalog
.get_metadata("foreign-server-metadata-format")
.unwrap()
.is_none());
assert!(raw.catalog.load_foreign_server_rows().unwrap()[0]
.metadata_json
.is_none());
assert!(raw
.catalog
.get_metadata("foreign-wrapper-catalog-format")
.unwrap()
.is_none());
assert_eq!(
raw.catalog
.metadata_with_prefix("foreign-wrapper/")
.unwrap()
.len(),
0
);
raw.catalog.delete_metadata("sql_triggers_json").unwrap();
let upgraded = Engine::from_persistent_provider(Arc::clone(&factory)).unwrap();
let expected = server(&upgraded, "legacy");
assert_eq!(expected.metadata.owner, RoleIdentity::BOOTSTRAP);
assert!(expected.metadata.oid >= 16_384);
assert_ne!(expected.metadata.object_id, [0; 16]);
assert_eq!(expected.options["key"], "original");
assert!(upgraded
.storage
.catalog
.as_ref()
.unwrap()
.load_foreign_server_rows()
.unwrap()[0]
.metadata_json
.is_some());
drop(upgraded);
drop(raw);
drop(factory);
let reopened = open(provider, &path);
assert_eq!(server(&reopened, "legacy"), expected);
}
}
#[test]
fn foreign_server_direct_options_remain_maps_and_empty_names_never_publish() {
let engine = Engine::new();
engine
.register_foreign_server(
"source".into(),
"memory_fdw".into(),
vec![
("a=b".into(), "first".into()),
("a=b".into(), "last".into()),
],
false,
)
.unwrap();
assert_eq!(
engine.foreign_server("source").unwrap().unwrap().options["a=b"],
"last"
);
assert!(engine
.register_foreign_server(String::new(), "memory_fdw".into(), vec![], false)
.is_err());
assert_eq!(engine.list_foreign_servers().unwrap(), ["source"]);
}
#[test]
fn foreign_server_creation_retains_owner_until_commit_or_rollback() {
use crate::tests::relation_lock_support::{after_shared_wait, sessions};
use uqa_execution::row_locks::shared_objects::SharedCatalogLock;
for provider in 0..3 {
for finish in ["COMMIT", "ROLLBACK"] {
let (_directory, first, second) = sessions(provider);
sql(&first, "CREATE ROLE server_owner");
sql(&first, "BEGIN; SET LOCAL ROLE server_owner; CREATE SERVER owned FOREIGN DATA WRAPPER memory_fdw");
let oid = u32::try_from(first.durable.roles.read()["server_owner"].oid).unwrap();
let (_second, result) = after_shared_wait(
&first,
second,
"DROP ROLE server_owner",
SharedCatalogLock::Object {
class_id: 1260,
oid,
},
finish,
);
if finish == "COMMIT" {
assert_eq!(result.unwrap_err().sqlstate(), Some("2BP01"));
} else {
result.unwrap();
}
}
}
}
#[test]
fn foreign_server_name_reservations_follow_competing_commit_or_rollback() {
use crate::tests::relation_lock_support::{after_shared_wait, sessions};
use uqa_execution::row_locks::shared_objects::SharedCatalogLock;
let oracle: serde_json::Value = serde_json::from_str(include_str!(
"../../../../../../tests/parity/pg18/foreign_server_concurrency_oracle.expected.json"
))
.unwrap();
for provider in 0..3 {
for schedule in oracle["schedules"].as_array().unwrap() {
let finish = schedule["finish_a"].as_str().unwrap();
let (_directory, first, second) = sessions(provider);
for statement in schedule["session_a"].as_array().unwrap() {
sql(&first, statement.as_str().unwrap());
}
let original = server(&first, "competing");
let (second, result) = after_shared_wait(
&first,
second,
schedule["session_b"].as_str().unwrap(),
SharedCatalogLock::Name {
class_id: 1417,
name: "competing",
},
finish,
);
if schedule["error"].is_null() {
let result = result.unwrap();
assert_eq!(
result.command_tag.as_deref(),
schedule["command_tag"].as_str()
);
assert_ne!(
server(&second, "competing").metadata.object_id,
original.metadata.object_id
);
} else {
let error = result.unwrap_err();
assert_eq!(
serde_json::json!({"sqlstate":error.sqlstate(),"message":error.to_string(),"detail":error.detail(),"hint":error.hint()}),
schedule["error"]
);
assert_eq!(server(&first, "competing"), original);
}
assert_eq!(second.take_sql_notices().len(), 0);
assert_eq!(
second.foreign_server("competing").unwrap().unwrap().options["source"],
schedule["final_source"]
);
}
}
}