Files
2026-08-30 21:11:23 +02:00

446 lines
21 KiB
Rust

// file: crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs
// version: 4
#![warn(missing_docs)]
#![deny(unreachable_pub)]
#![forbid(unsafe_code)]
//! Opt-in real PostgreSQL proof for the Store foundation runtime.
//!
//! The test reads one dedicated PostgreSQL URI from stdin, refuses to start
//! when any KSP Store table managed by V000/V001 already exists, never prints
//! the URI, and cleans up only the isolated schema it proved absent first. It
//! validates migration/bootstrap behavior, not RawTransaction capabilities.
const LIVE_BOOTSTRAP_SQL: &str = include_str!("../migrations/v000_bootstrap/tables/001_ksp_store_schema_migrations.sql");
const LIVE_BROKEN_CHECKSUM_A: &str = "0000000000000000000000000000000000000000000000000000000000000000";
const LIVE_BROKEN_CHECKSUM_B: &str = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_observations;
DROP TABLE IF EXISTS ksp_raw_transaction_archive_payloads;
DROP TABLE IF EXISTS ksp_raw_transactions;
DROP TABLE IF EXISTS ksp_store_identity;
DROP TABLE IF EXISTS ksp_store_schema_migrations;"#;
const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM information_schema.tables
WHERE table_schema = current_schema()
AND table_name IN (
'ksp_store_schema_migrations',
'ksp_store_identity',
'ksp_raw_transactions',
'ksp_raw_transaction_observations',
'ksp_raw_transaction_archive_payloads'
)
AND table_type = 'BASE TABLE'
)"#;
const LIVE_MAX_URI_BYTES: usize = 4_096;
const LIVE_METADATA_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM information_schema.tables
WHERE table_schema = current_schema()
AND table_name = 'ksp_store_schema_migrations'
AND table_type = 'BASE TABLE'
)"#;
const LIVE_SENTINEL_CHECKSUM_SQL: &str = "SELECT checksum FROM ksp_store_schema_migrations WHERE version = 0";
const LIVE_SENTINEL_INSERT_SQL: &str =
"INSERT INTO ksp_store_schema_migrations (version, name, checksum, applied_at) VALUES (0, 'bootstrap', 'pre008_rollback_injected', CURRENT_TIMESTAMP)";
const LIVE_SENTINEL_UPDATE_SQL: &str = "UPDATE ksp_store_schema_migrations SET checksum = $1 WHERE version = 0";
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
struct LiveFailure {
phase: &'static str,
}
impl LiveFailure {
const fn new(phase: &'static str) -> Self {
return Self { phase };
}
const fn phase(&self) -> &'static str {
return self.phase;
}
}
#[test]
#[ignore = "opt-in real PostgreSQL foundation proof; reads one dedicated URI from stdin"]
fn pre_008_real_postgres_foundation_is_safe_idempotent_concurrent_and_recoverable() {
eprintln!("KSP Store PostgreSQL live proof: reading one dedicated URI from stdin without echoing it from the test.");
let uri_result = read_uri_from_stdin();
let uri = match uri_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => panic!("PostgreSQL live input rejected at phase {}", error.phase()),
};
let runtime_result = tokio::runtime::Builder::new_current_thread().enable_all().build();
let runtime = match runtime_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => panic!("PostgreSQL live runtime could not be constructed"),
};
let outcome = runtime.block_on(run_live_test(uri.as_str()));
if let std::result::Result::Err(error) = outcome {
panic!("PostgreSQL live foundation proof failed at safe phase {}", error.phase());
}
return;
}
fn read_uri_from_stdin() -> std::result::Result<std::string::String, LiveFailure> {
let mut input = std::string::String::new();
let read_result = std::io::stdin().read_line(&mut input);
match read_result {
std::result::Result::Ok(0) | std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("stdin_read")),
std::result::Result::Ok(_) => {},
}
let uri = input.trim().to_owned();
if uri.is_empty() || uri.len() > LIVE_MAX_URI_BYTES {
return std::result::Result::Err(LiveFailure::new("stdin_uri"));
}
return std::result::Result::Ok(uri);
}
async fn run_live_test(uri: &str) -> std::result::Result<(), LiveFailure> {
let admin_result = connect_admin(uri).await;
let mut admin = match admin_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let preexisting_result = managed_schema_exists(&admin).await;
let preexisting = match preexisting_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if preexisting {
return std::result::Result::Err(LiveFailure::new("managed_schema_preexisting_refusal"));
}
let major_result = postgres_major(&admin).await;
let major = match major_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if major < 15 {
return std::result::Result::Err(LiveFailure::new("postgres_major_unsupported"));
}
eprintln!("KSP Store PostgreSQL live proof: server major {major}");
let mut owns_schema = false;
let scenario = run_foundation_scenario(&mut admin, uri, &mut owns_schema).await;
let cleanup = if owns_schema { drop_managed_schema(&admin).await } else { std::result::Result::Ok(()) };
if let std::result::Result::Err(error) = cleanup {
return std::result::Result::Err(error);
}
if let std::result::Result::Err(error) = scenario {
return std::result::Result::Err(error);
}
let remains_result = managed_schema_exists(&admin).await;
let remains = match remains_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if remains {
return std::result::Result::Err(LiveFailure::new("cleanup_verification"));
}
return std::result::Result::Ok(());
}
async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str, owns_schema: &mut bool) -> std::result::Result<(), LiveFailure> {
let initial_result = open_backend(uri).await;
let initial = match initial_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let created_result = metadata_exists(admin).await;
let created = match created_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if !created {
return std::result::Result::Err(LiveFailure::new("initial_bootstrap_metadata"));
}
*owns_schema = true;
let initial_health = initial.health().await;
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(2) || initial_health.pending_migration_count() != 0 {
return std::result::Result::Err(LiveFailure::new("initial_health"));
}
let initial_close = close_backend(initial).await;
if let std::result::Result::Err(error) = initial_close {
return std::result::Result::Err(error);
}
let idempotent_result = open_backend(uri).await;
let idempotent = match idempotent_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let idempotent_health = idempotent.health().await;
if !idempotent_health.is_ready()
|| idempotent_health.migration_version() != std::option::Option::Some(2)
|| idempotent_health.pending_migration_count() != 0
{
return std::result::Result::Err(LiveFailure::new("idempotent_health"));
}
let idempotent_close = close_backend(idempotent).await;
if let std::result::Result::Err(error) = idempotent_close {
return std::result::Result::Err(error);
}
let reset_result = drop_managed_schema(admin).await;
if let std::result::Result::Err(error) = reset_result {
return std::result::Result::Err(error);
}
let concurrent_result = concurrent_bootstrap(uri).await;
if let std::result::Result::Err(error) = concurrent_result {
return std::result::Result::Err(error);
}
let after_concurrent_result = metadata_exists(admin).await;
match after_concurrent_result {
std::result::Result::Ok(true) => {},
std::result::Result::Ok(false) => return std::result::Result::Err(LiveFailure::new("concurrent_bootstrap_metadata")),
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
let checksum_result = sentinel_checksum(admin).await;
let checksum = match checksum_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let broken_checksum = if checksum == LIVE_BROKEN_CHECKSUM_A { LIVE_BROKEN_CHECKSUM_B } else { LIVE_BROKEN_CHECKSUM_A };
let corrupt_result = set_sentinel_checksum(admin, broken_checksum).await;
if let std::result::Result::Err(error) = corrupt_result {
return std::result::Result::Err(error);
}
let mismatch_settings = match settings(uri) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let mismatch = ksp_store_postgres_lib::PostgresBackend::open(mismatch_settings).await;
match mismatch {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationMismatch => {},
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("checksum_mismatch_classification")),
std::result::Result::Ok(backend) => {
let _ = close_backend(backend).await;
return std::result::Result::Err(LiveFailure::new("checksum_mismatch_accepted"));
},
}
let restore_result = set_sentinel_checksum(admin, checksum.as_str()).await;
if let std::result::Result::Err(error) = restore_result {
return std::result::Result::Err(error);
}
let recovered_result = open_backend(uri).await;
let recovered = match recovered_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let recovered_health = recovered.health().await;
if !recovered_health.is_ready() {
return std::result::Result::Err(LiveFailure::new("checksum_recovery_health"));
}
let recovered_close = close_backend(recovered).await;
if let std::result::Result::Err(error) = recovered_close {
return std::result::Result::Err(error);
}
let rollback_reset = drop_managed_schema(admin).await;
if let std::result::Result::Err(error) = rollback_reset {
return std::result::Result::Err(error);
}
let rollback_result = prove_transaction_rollback(admin).await;
if let std::result::Result::Err(error) = rollback_result {
return std::result::Result::Err(error);
}
let absent_after_rollback = metadata_exists(admin).await;
match absent_after_rollback {
std::result::Result::Ok(false) => {},
std::result::Result::Ok(true) => return std::result::Result::Err(LiveFailure::new("rollback_left_metadata")),
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
let final_result = open_backend(uri).await;
let final_backend = match final_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let final_health = final_backend.health().await;
if !final_health.is_ready() || final_health.migration_version() != std::option::Option::Some(2) || final_health.pending_migration_count() != 0 {
return std::result::Result::Err(LiveFailure::new("final_health"));
}
return close_backend(final_backend).await;
}
fn settings(uri: &str) -> std::result::Result<ksp_store_postgres_lib::PostgresBackendSettings, LiveFailure> {
let network_result = ksp_store_api::RawNetworkId::new("devnet");
let network = match network_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("network")),
};
return std::result::Result::Ok(ksp_store_postgres_lib::PostgresBackendSettings::new(
network,
uri,
4,
std::time::Duration::from_secs(10),
std::time::Duration::from_secs(5),
std::time::Duration::from_secs(10),
std::time::Duration::from_secs(5),
ksp_store_postgres_lib::PostgresBackendTlsMode::Disabled,
true,
std::time::Duration::from_secs(30),
std::time::Duration::from_secs(10),
));
}
async fn open_backend(uri: &str) -> std::result::Result<ksp_store_postgres_lib::PostgresBackend, LiveFailure> {
let settings_result = settings(uri);
let value = match settings_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return match ksp_store_postgres_lib::PostgresBackend::open(value).await {
std::result::Result::Ok(backend) => std::result::Result::Ok(backend),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("backend_open")),
};
}
async fn close_backend(backend: ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> {
return match backend.close(std::time::Duration::from_secs(5)).await {
std::result::Result::Ok(()) => std::result::Result::Ok(()),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("backend_close")),
};
}
async fn concurrent_bootstrap(uri: &str) -> std::result::Result<(), LiveFailure> {
let first_settings = match settings(uri) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let second_settings = match settings(uri) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let first = tokio::spawn(async move {
return ksp_store_postgres_lib::PostgresBackend::open(first_settings).await;
});
let second = tokio::spawn(async move {
return ksp_store_postgres_lib::PostgresBackend::open(second_settings).await;
});
let first_joined = first.await;
let second_joined = second.await;
let first_backend = match first_joined {
std::result::Result::Ok(std::result::Result::Ok(value)) => value,
_ => return std::result::Result::Err(LiveFailure::new("concurrent_first")),
};
let second_backend = match second_joined {
std::result::Result::Ok(std::result::Result::Ok(value)) => value,
_ => {
let _ = close_backend(first_backend).await;
return std::result::Result::Err(LiveFailure::new("concurrent_second"));
},
};
let first_close = close_backend(first_backend).await;
let second_close = close_backend(second_backend).await;
if first_close.is_err() || second_close.is_err() {
return std::result::Result::Err(LiveFailure::new("concurrent_close"));
}
return std::result::Result::Ok(());
}
async fn connect_admin(uri: &str) -> std::result::Result<tokio_postgres::Client, LiveFailure> {
let parsed = uri.parse::<tokio_postgres::Config>();
let mut config = match parsed {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("admin_config")),
};
config.ssl_mode(tokio_postgres::config::SslMode::Disable);
config.ssl_negotiation(tokio_postgres::config::SslNegotiation::Postgres);
let connected = config.connect(tokio_postgres::NoTls).await;
let (client, connection) = match connected {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("admin_connect")),
};
let _connection_task = tokio::spawn(async move {
let _result = connection.await;
return;
});
return std::result::Result::Ok(client);
}
async fn postgres_major(client: &tokio_postgres::Client) -> std::result::Result<u32, LiveFailure> {
let row_result = client.query_one("SHOW server_version_num", &[]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("server_version")),
};
let value_result = row.try_get::<usize, std::string::String>(0);
let value = match value_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("server_version_decode")),
};
let parsed = value.parse::<u32>();
let version_num = match parsed {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("server_version_parse")),
};
return std::result::Result::Ok(version_num / 10_000);
}
async fn managed_schema_exists(client: &tokio_postgres::Client) -> std::result::Result<bool, LiveFailure> {
let row_result = client.query_one(LIVE_MANAGED_SCHEMA_EXISTS_SQL, &[]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("managed_schema_probe")),
};
return match row.try_get::<usize, bool>(0) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("managed_schema_probe_decode")),
};
}
async fn metadata_exists(client: &tokio_postgres::Client) -> std::result::Result<bool, LiveFailure> {
let row_result = client.query_one(LIVE_METADATA_EXISTS_SQL, &[]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("metadata_probe")),
};
return match row.try_get::<usize, bool>(0) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("metadata_probe_decode")),
};
}
async fn drop_managed_schema(client: &tokio_postgres::Client) -> std::result::Result<(), LiveFailure> {
return match client.batch_execute(LIVE_MANAGED_SCHEMA_DROP_SQL).await {
std::result::Result::Ok(()) => std::result::Result::Ok(()),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("managed_schema_cleanup")),
};
}
async fn sentinel_checksum(client: &tokio_postgres::Client) -> std::result::Result<std::string::String, LiveFailure> {
let row_result = client.query_one(LIVE_SENTINEL_CHECKSUM_SQL, &[]).await;
let row = match row_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("sentinel_read")),
};
return match row.try_get::<usize, std::string::String>(0) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("sentinel_decode")),
};
}
async fn set_sentinel_checksum(client: &tokio_postgres::Client, checksum: &str) -> std::result::Result<(), LiveFailure> {
let update = client.execute(LIVE_SENTINEL_UPDATE_SQL, &[&checksum]).await;
return match update {
std::result::Result::Ok(1) => std::result::Result::Ok(()),
_ => std::result::Result::Err(LiveFailure::new("sentinel_update")),
};
}
async fn prove_transaction_rollback(client: &mut tokio_postgres::Client) -> std::result::Result<(), LiveFailure> {
let transaction_result = client.transaction().await;
let transaction = match transaction_result {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("rollback_begin")),
};
let create = transaction.batch_execute(LIVE_BOOTSTRAP_SQL).await;
if create.is_err() {
return std::result::Result::Err(LiveFailure::new("rollback_create"));
}
let insert = transaction.batch_execute(LIVE_SENTINEL_INSERT_SQL).await;
if insert.is_err() {
return std::result::Result::Err(LiveFailure::new("rollback_insert"));
}
let injected = transaction.batch_execute("SELECT 1 / 0").await;
if injected.is_ok() {
return std::result::Result::Err(LiveFailure::new("rollback_injection_missing"));
}
drop(transaction);
return std::result::Result::Ok(());
}