// 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 { 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 { 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 { 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 { let parsed = uri.parse::(); 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 { 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::(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::(); 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 { 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::(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 { 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::(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 { 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::(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(()); }