// file: crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs // version: 4 #![warn(missing_docs)] #![deny(unreachable_pub)] #![forbid(unsafe_code)] //! Opt-in real PostgreSQL proof for the complete RawTransaction vertical slice. //! //! The test reads one dedicated PostgreSQL URI from stdin, refuses to start //! when any managed KSP Store table already exists, never prints the URI, and //! drops only the isolated schema it proved absent before the run. const LIVE_CANCEL_WAIT: std::time::Duration = std::time::Duration::from_millis(300); const LIVE_INDEX_EXISTS_SQL: &str = r#"SELECT EXISTS ( SELECT 1 FROM pg_indexes WHERE schemaname = current_schema() AND indexname = 'ix_ksp_raw_transactions_slot_signature' )"#; const LIVE_LOCK_OBSERVATION_SQL: &str = "SELECT observation_key FROM ksp_raw_transaction_observations WHERE observation_key = $1 FOR UPDATE"; 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_OBSERVATION_EXISTS_SQL: &str = "SELECT EXISTS (SELECT 1 FROM ksp_raw_transaction_observations WHERE observation_key = $1)"; const LIVE_TRANSACTION_EXISTS_SQL: &str = "SELECT EXISTS (SELECT 1 FROM ksp_raw_transactions WHERE signature = $1)"; #[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; } } #[derive(Debug)] enum LivePersistResult { BackendError(ksp_store_postgres_lib::PostgresBackendErrorKind), Outcome(ksp_store_api::RawAcquisitionWriteOutcome), } #[derive(Debug)] enum LiveRetentionResult { BackendError, Outcome(ksp_store_api::RawRetentionWriteOutcome), } #[test] #[ignore = "opt-in real PostgreSQL RawTransaction proof; reads one dedicated URI from stdin"] fn pre_009_real_postgres_raw_transaction_vertical_slice_is_atomic_concurrent_and_recoverable() { eprintln!("KSP Store RawTransaction 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 RawTransaction 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 RawTransaction 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 RawTransaction live 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 RawTransaction live proof: server major {major}"); let mut owns_schema = false; let scenario = run_raw_transaction_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; match remains_result { std::result::Result::Ok(false) => {}, std::result::Result::Ok(true) => return std::result::Result::Err(LiveFailure::new("cleanup_verification")), std::result::Result::Err(error) => return std::result::Result::Err(error), } return std::result::Result::Ok(()); } async fn run_raw_transaction_scenario(admin: &mut tokio_postgres::Client, uri: &str, owns_schema: &mut bool) -> std::result::Result<(), LiveFailure> { let initial_result = open_backend(uri, "devnet", true, true).await; let initial = match initial_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; *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 wrong_network_result = open_backend_result(uri, "testnet", true, true).await; match wrong_network_result { std::result::Result::Ok(std::result::Result::Err(error)) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationMismatch => {}, std::result::Result::Ok(std::result::Result::Ok(backend)) => { let _ = close_backend(backend).await; return std::result::Result::Err(LiveFailure::new("wrong_network_accepted")); }, _ => return std::result::Result::Err(LiveFailure::new("wrong_network_classification")), } let initial_close = close_backend(initial).await; if let std::result::Result::Err(error) = initial_close { return std::result::Result::Err(error); } let schema_result = prove_schema_update_policy(admin, uri).await; if let std::result::Result::Err(error) = schema_result { return std::result::Result::Err(error); } let backend_result = open_backend(uri, "devnet", true, true).await; let backend = match backend_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; for proof in [ prove_atomic_insert_and_reads(&backend).await, prove_additional_observation_and_atomic_rollback(&backend).await, prove_pagination(&backend).await, prove_retention_and_rehydrate(&backend).await, ] { if let std::result::Result::Err(error) = proof { return std::result::Result::Err(error); } } let backend_close = close_backend(backend).await; if let std::result::Result::Err(error) = backend_close { return std::result::Result::Err(error); } let identical_result = prove_concurrent_identical_insert(uri).await; if let std::result::Result::Err(error) = identical_result { return std::result::Result::Err(error); } let divergent_result = prove_concurrent_divergent_insert(admin, uri).await; if let std::result::Result::Err(error) = divergent_result { return std::result::Result::Err(error); } let race_result = prove_retention_races(uri).await; if let std::result::Result::Err(error) = race_result { return std::result::Result::Err(error); } let cancellation_result = prove_cancellation_rollback(admin, uri).await; if let std::result::Result::Err(error) = cancellation_result { return std::result::Result::Err(error); } let reopened_result = open_backend(uri, "devnet", true, true).await; let reopened = match reopened_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let reopened_health = reopened.health().await; if !reopened_health.is_ready() || reopened_health.migration_version() != std::option::Option::Some(2) { return std::result::Result::Err(LiveFailure::new("reopen_health")); } let reference_result = raw_reference(10); let reference = match reference_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let read_result = reopened.get_raw_transaction(&reference).await; match read_result { std::result::Result::Ok(std::option::Option::Some(value)) => { let exact = transaction_matches(&value, 10, 1_000, 10); if !exact { return std::result::Result::Err(LiveFailure::new("reopen_read_exact")); } }, _ => return std::result::Result::Err(LiveFailure::new("reopen_read")), } return close_backend(reopened).await; } async fn prove_schema_update_policy(admin: &tokio_postgres::Client, uri: &str) -> std::result::Result<(), LiveFailure> { let drop_result = admin.execute("DROP INDEX ix_ksp_raw_transactions_slot_signature", &[]).await; match drop_result { std::result::Result::Ok(_) => {}, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("schema_drop_index")), } let blocked_result = open_backend_result(uri, "devnet", true, false).await; match blocked_result { std::result::Result::Ok(std::result::Result::Err(error)) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::MigrationMismatch => {}, std::result::Result::Ok(std::result::Result::Ok(backend)) => { let _ = close_backend(backend).await; return std::result::Result::Err(LiveFailure::new("schema_autoupdate_disabled_accepted")); }, _ => return std::result::Result::Err(LiveFailure::new("schema_autoupdate_disabled_classification")), } let repaired_result = open_backend(uri, "devnet", true, true).await; let repaired = match repaired_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let index_result = index_exists(admin).await; match index_result { std::result::Result::Ok(true) => {}, std::result::Result::Ok(false) => return std::result::Result::Err(LiveFailure::new("schema_index_not_repaired")), std::result::Result::Err(error) => return std::result::Result::Err(error), } return close_backend(repaired).await; } async fn prove_atomic_insert_and_reads(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { let transaction_result = raw_transaction(10, 1_000, 10); let transaction = match transaction_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let observation_result = raw_observation(10, 10, 10); let observation = match observation_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let persisted = backend.persist_raw_transaction_acquisition(transaction, observation, ksp_store_api::RawTransactionAcquisitionMode::Normal).await; match persisted { std::result::Result::Ok(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted && outcome.observation() == ksp_store_api::RawObservationWriteOutcome::Inserted => {}, _ => return std::result::Result::Err(LiveFailure::new("atomic_insert")), } let reference = match raw_reference(10) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let read = backend.get_raw_transaction(&reference).await; match read { std::result::Result::Ok(std::option::Option::Some(value)) if transaction_matches(&value, 10, 1_000, 10) => {}, _ => return std::result::Result::Err(LiveFailure::new("atomic_get")), } let key = ksp_store_api::RawObservationKey::new([10; 32]); let observed = backend.get_raw_transaction_observation(&key).await; let expected = match raw_observation(10, 10, 10) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match observed { std::result::Result::Ok(std::option::Option::Some(value)) if value == expected => {}, _ => return std::result::Result::Err(LiveFailure::new("atomic_observation_get")), } return std::result::Result::Ok(()); } async fn prove_additional_observation_and_atomic_rollback(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { let additional = match raw_observation(41, 10, 41) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let first = backend.record_raw_transaction_observation(additional).await; if first != std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::Inserted) { return std::result::Result::Err(LiveFailure::new("observation_insert")); } let identical = match raw_observation(41, 10, 41) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let second = backend.record_raw_transaction_observation(identical).await; if second != std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::AlreadyPresent) { return std::result::Result::Err(LiveFailure::new("observation_idempotent")); } let divergent = match raw_observation(41, 10, 42) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let divergent_result = backend.record_raw_transaction_observation(divergent).await; match divergent_result { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {}, _ => return std::result::Result::Err(LiveFailure::new("observation_conflict")), } let rollback_transaction = match raw_transaction(42, 4_200, 42) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let rollback_observation = match raw_observation(41, 42, 43) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let rollback = backend .persist_raw_transaction_acquisition(rollback_transaction, rollback_observation, ksp_store_api::RawTransactionAcquisitionMode::Normal) .await; match rollback { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {}, _ => return std::result::Result::Err(LiveFailure::new("atomic_rollback_conflict")), } let rollback_reference = match raw_reference(42) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match backend.get_raw_transaction(&rollback_reference).await { std::result::Result::Ok(std::option::Option::None) => {}, _ => return std::result::Result::Err(LiveFailure::new("atomic_rollback_left_canonical")), } return std::result::Result::Ok(()); } async fn prove_concurrent_identical_insert(uri: &str) -> std::result::Result<(), LiveFailure> { let first = tokio::spawn(persist_once(uri.to_owned(), 20, 2_000, 20, 20)); let second = tokio::spawn(persist_once(uri.to_owned(), 20, 2_000, 20, 20)); let first_result = match joined_persist(first.await, "concurrent_identical_first") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let second_result = match joined_persist(second.await, "concurrent_identical_second") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let pair = [first_result, second_result]; let inserted = pair.iter().filter(|value| return matches_inserted(value)).count(); let already = pair.iter().filter(|value| return matches_already_present(value)).count(); if inserted != 1 || already != 1 { return std::result::Result::Err(LiveFailure::new("concurrent_identical_outcome")); } return std::result::Result::Ok(()); } async fn prove_concurrent_divergent_insert(admin: &tokio_postgres::Client, uri: &str) -> std::result::Result<(), LiveFailure> { let first = tokio::spawn(persist_once(uri.to_owned(), 30, 3_000, 30, 31)); let second = tokio::spawn(persist_once(uri.to_owned(), 30, 3_000, 31, 32)); let first_result = match joined_persist(first.await, "concurrent_divergent_first") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let second_result = match joined_persist(second.await, "concurrent_divergent_second") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let pair = [first_result, second_result]; let inserted = pair.iter().filter(|value| return matches_inserted(value)).count(); let conflicts = pair .iter() .filter(|value| matches!(value, LivePersistResult::BackendError(ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict))) .count(); if inserted != 1 || conflicts != 1 { return std::result::Result::Err(LiveFailure::new("concurrent_divergent_outcome")); } let first_key = ksp_store_api::RawObservationKey::new([31; 32]); let second_key = ksp_store_api::RawObservationKey::new([32; 32]); let first_exists = observation_exists(admin, &first_key).await; let second_exists = observation_exists(admin, &second_key).await; let count = match (first_exists, second_exists) { (std::result::Result::Ok(first_value), std::result::Result::Ok(second_value)) => usize::from(first_value) + usize::from(second_value), _ => return std::result::Result::Err(LiveFailure::new("concurrent_divergent_observation_probe")), }; if count != 1 { return std::result::Result::Err(LiveFailure::new("concurrent_divergent_rollback")); } return std::result::Result::Ok(()); } async fn prove_pagination(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { for (signature_seed, slot) in [(50_u8, 9_000_u64), (51, 9_000), (52, 9_001), (53, 9_002), (54, 9_002)] { let transaction = match raw_transaction(signature_seed, slot, signature_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let observation = match raw_observation(signature_seed, signature_seed, signature_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let result = backend.persist_raw_transaction_acquisition(transaction, observation, ksp_store_api::RawTransactionAcquisitionMode::Normal).await; match result { std::result::Result::Ok(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {}, _ => return std::result::Result::Err(LiveFailure::new("pagination_seed")), } } let asc = collect_pages(backend, ksp_store_api::RawSortDirection::Ascending).await; match asc { std::result::Result::Ok([50, 51, 52, 53, 54]) => {}, _ => return std::result::Result::Err(LiveFailure::new("pagination_ascending")), } let desc = collect_pages(backend, ksp_store_api::RawSortDirection::Descending).await; match desc { std::result::Result::Ok([54, 53, 52, 51, 50]) => {}, _ => return std::result::Result::Err(LiveFailure::new("pagination_descending")), } let first_query = match page_query(ksp_store_api::RawSortDirection::Ascending, std::option::Option::None, 9_000, 9_002) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let first_page = match backend.list_raw_transactions(&first_query).await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("pagination_replay_seed")), }; let cursor = match first_page.next_cursor() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(LiveFailure::new("pagination_replay_cursor")), }; let direction_cursor = match copy_cursor(cursor) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let direction_query = match page_query(ksp_store_api::RawSortDirection::Descending, std::option::Option::Some(direction_cursor), 9_000, 9_002) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match backend.list_raw_transactions(&direction_query).await { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::QueryInvalid => {}, _ => return std::result::Result::Err(LiveFailure::new("cursor_direction_replay")), } let range_cursor = match copy_cursor(cursor) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let range_query = match page_query(ksp_store_api::RawSortDirection::Ascending, std::option::Option::Some(range_cursor), 9_001, 9_002) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match backend.list_raw_transactions(&range_query).await { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::QueryInvalid => {}, _ => return std::result::Result::Err(LiveFailure::new("cursor_range_replay")), } return std::result::Result::Ok(()); } async fn prove_retention_and_rehydrate(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { let transaction = match raw_transaction(70, 7_000, 70) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let observation = match raw_observation(70, 70, 70) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let seed = backend.persist_raw_transaction_acquisition(transaction, observation, ksp_store_api::RawTransactionAcquisitionMode::Normal).await; if !matches!(seed, std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted) { return std::result::Result::Err(LiveFailure::new("retention_seed")); } let reference = match raw_reference(70) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let mismatch = match retention_transition(reference.clone(), ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if backend.transition_raw_transaction_retention(mismatch).await != std::result::Result::Ok(ksp_store_api::RawRetentionWriteOutcome::ExpectedStateMismatch) { return std::result::Result::Err(LiveFailure::new("retention_compare_mismatch")); } let archive = match retention_transition(reference.clone(), ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if backend.transition_raw_transaction_retention(archive).await != std::result::Result::Ok(ksp_store_api::RawRetentionWriteOutcome::Applied) { return std::result::Result::Err(LiveFailure::new("retention_archive")); } let archived_read = backend.get_raw_transaction(&reference).await; match archived_read { std::result::Result::Ok(std::option::Option::Some(value)) if transaction_matches(&value, 70, 7_000, 70) => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_archived_read")), } let already_archive = match retention_transition(reference.clone(), ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if backend.transition_raw_transaction_retention(already_archive).await != std::result::Result::Ok(ksp_store_api::RawRetentionWriteOutcome::AlreadyAtTarget) { return std::result::Result::Err(LiveFailure::new("retention_archive_idempotent")); } let purge = match retention_transition(reference.clone(), ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if backend.transition_raw_transaction_retention(purge).await != std::result::Result::Ok(ksp_store_api::RawRetentionWriteOutcome::Applied) { return std::result::Result::Err(LiveFailure::new("retention_purge")); } match backend.get_raw_transaction(&reference).await { std::result::Result::Ok(std::option::Option::None) => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_purged_get")), } let tombstone = backend.get_raw_transaction_tombstone(&reference).await; match tombstone { std::result::Result::Ok(std::option::Option::Some(value)) if value.reference() == &reference && value.slot() == 7_000 && value.format_id().as_str() == "ksp-live-raw-v1" && value.format_version() == 1 && value.content_hash() == ksp_store_api::RawContentHash::new([170; 32]) => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_tombstone")), } let normal_transaction = match raw_transaction(70, 7_000, 70) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let normal_observation = match raw_observation(71, 70, 71) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let normal = backend .persist_raw_transaction_acquisition(normal_transaction, normal_observation, ksp_store_api::RawTransactionAcquisitionMode::Normal) .await; match normal { std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::SkippedPurged && value.observation() == ksp_store_api::RawObservationWriteOutcome::NotRecorded => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_normal_tombstone")), } let skipped_key = ksp_store_api::RawObservationKey::new([71; 32]); match backend.get_raw_transaction_observation(&skipped_key).await { std::result::Result::Ok(std::option::Option::None) => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_normal_observation")), } let divergent_transaction = match raw_transaction(70, 7_000, 71) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let divergent_observation = match raw_observation(72, 70, 72) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let divergent_force = backend .persist_raw_transaction_acquisition(divergent_transaction, divergent_observation, ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate) .await; match divergent_force { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_force_divergent")), } let compatible_transaction = match raw_transaction(70, 7_000, 70) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let compatible_observation = match raw_observation(73, 70, 73) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let compatible_force = backend .persist_raw_transaction_acquisition(compatible_transaction, compatible_observation, ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate) .await; match compatible_force { std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Rehydrated && value.observation() == ksp_store_api::RawObservationWriteOutcome::Inserted => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_force_compatible")), } let compacted = match retention_transition(reference.clone(), ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Compacted) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match backend.transition_raw_transaction_retention(compacted).await { std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::RetentionCompactionUnsupported => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_compacted_rejection")), } match backend.get_raw_transaction_retention_state(&reference).await { std::result::Result::Ok(std::option::Option::Some(ksp_store_api::RawRetentionState::Full)) => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_compacted_state_changed")), } return std::result::Result::Ok(()); } async fn prove_retention_races(uri: &str) -> std::result::Result<(), LiveFailure> { let seed_result = persist_once(uri.to_owned(), 75, 7_500, 75, 75).await; match seed_result { std::result::Result::Ok(LivePersistResult::Outcome(value)) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_race_seed")), } let archive_first = tokio::spawn(transition_once(uri.to_owned(), 75, ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived)); let archive_second = tokio::spawn(transition_once(uri.to_owned(), 75, ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived)); let first_archive = match joined_retention(archive_first.await, "retention_race_archive_first") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let second_archive = match joined_retention(archive_second.await, "retention_race_archive_second") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if !one_applied_one_already([first_archive, second_archive]) { return std::result::Result::Err(LiveFailure::new("retention_race_archive_outcome")); } let purge_first = tokio::spawn(transition_once(uri.to_owned(), 75, ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged)); let purge_second = tokio::spawn(transition_once(uri.to_owned(), 75, ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged)); let first_purge = match joined_retention(purge_first.await, "retention_race_purge_first") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let second_purge = match joined_retention(purge_second.await, "retention_race_purge_second") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if !one_applied_one_already([first_purge, second_purge]) { return std::result::Result::Err(LiveFailure::new("retention_race_purge_outcome")); } let rehydrated_result = persist_once(uri.to_owned(), 75, 7_500, 75, 76).await; match rehydrated_result { std::result::Result::Ok(LivePersistResult::Outcome(value)) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Rehydrated => {}, _ => return std::result::Result::Err(LiveFailure::new("retention_race_rehydrate")), } return std::result::Result::Ok(()); } async fn prove_cancellation_rollback(admin: &mut tokio_postgres::Client, uri: &str) -> std::result::Result<(), LiveFailure> { let seed_result = persist_once(uri.to_owned(), 80, 8_000, 80, 80).await; match seed_result { std::result::Result::Ok(LivePersistResult::Outcome(value)) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {}, _ => return std::result::Result::Err(LiveFailure::new("cancellation_seed")), } let lock_transaction = match admin.transaction().await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("cancellation_lock_begin")), }; let key = ksp_store_api::RawObservationKey::new([80; 32]); let key_bytes: &[u8] = key.as_bytes(); let lock_result = lock_transaction.query_one(LIVE_LOCK_OBSERVATION_SQL, &[&key_bytes]).await; if lock_result.is_err() { return std::result::Result::Err(LiveFailure::new("cancellation_lock_observation")); } let cancel_backend_result = open_backend(uri, "devnet", true, true).await; let cancel_backend = match cancel_backend_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let transaction = match raw_transaction(81, 8_100, 81) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let observation = match raw_observation(80, 81, 81) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let mut task = tokio::spawn(async move { return cancel_backend.persist_raw_transaction_acquisition(transaction, observation, ksp_store_api::RawTransactionAcquisitionMode::Normal).await; }); let premature = tokio::time::timeout(LIVE_CANCEL_WAIT, &mut task).await; if premature.is_ok() { let _ = lock_transaction.rollback().await; return std::result::Result::Err(LiveFailure::new("cancellation_operation_not_blocked")); } task.abort(); let cancelled = task.await; match cancelled { std::result::Result::Err(error) if error.is_cancelled() => {}, _ => { let _ = lock_transaction.rollback().await; return std::result::Result::Err(LiveFailure::new("cancellation_join")); }, } let unlock_result = lock_transaction.rollback().await; if unlock_result.is_err() { return std::result::Result::Err(LiveFailure::new("cancellation_unlock")); } let reference = match raw_reference(81) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let exists_result = transaction_exists(admin, &reference).await; match exists_result { std::result::Result::Ok(false) => {}, std::result::Result::Ok(true) => return std::result::Result::Err(LiveFailure::new("cancellation_left_canonical")), std::result::Result::Err(error) => return std::result::Result::Err(error), } return std::result::Result::Ok(()); } async fn collect_pages( backend: &ksp_store_postgres_lib::PostgresBackend, direction: ksp_store_api::RawSortDirection, ) -> std::result::Result<[u8; 5], LiveFailure> { let mut output = [0_u8; 5]; let mut output_index = 0_usize; let mut cursor: std::option::Option = std::option::Option::None; loop { let query = match page_query(direction, cursor, 9_000, 9_002) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let page = match backend.list_raw_transactions(&query).await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("pagination_query")), }; for reference in page.items() { if output_index >= output.len() { return std::result::Result::Err(LiveFailure::new("pagination_cardinality")); } output[output_index] = reference.signature().as_bytes()[0]; output_index += 1; } cursor = match page.next_cursor() { std::option::Option::Some(value) => match copy_cursor(value) { std::result::Result::Ok(copied) => std::option::Option::Some(copied), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => break, }; } if output_index != output.len() { return std::result::Result::Err(LiveFailure::new("pagination_cardinality")); } return std::result::Result::Ok(output); } fn page_query( direction: ksp_store_api::RawSortDirection, cursor: std::option::Option, start: u64, end: u64, ) -> std::result::Result { let network = match network("devnet") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let range = match ksp_store_api::RawSlotRange::new(std::option::Option::Some(start), std::option::Option::Some(end)) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("pagination_range")), }; let limit = match ksp_store_api::RawPageLimit::new(2) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("pagination_limit")), }; let page = match cursor { std::option::Option::Some(value) => ksp_store_api::RawPageRequest::after(limit, value), std::option::Option::None => ksp_store_api::RawPageRequest::first(limit), }; return std::result::Result::Ok(ksp_store_api::RawTransactionQuery::new(network, range, direction, page)); } fn copy_cursor(cursor: &ksp_store_api::RawPageCursor) -> std::result::Result { let bytes = cursor.as_bytes().to_vec().into_boxed_slice(); return match ksp_store_api::RawPageCursor::try_new(bytes) { std::result::Result::Ok(value) => std::result::Result::Ok(value), std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("cursor_copy")), }; } async fn persist_once( uri: std::string::String, signature_seed: u8, slot: u64, payload_seed: u8, observation_seed: u8, ) -> std::result::Result { let backend = match open_backend(uri.as_str(), "devnet", true, true).await { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let transaction = match raw_transaction(signature_seed, slot, payload_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let observation = match raw_observation(observation_seed, signature_seed, observation_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let persisted = backend.persist_raw_transaction_acquisition(transaction, observation, ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate).await; let result = match persisted { std::result::Result::Ok(value) => LivePersistResult::Outcome(value), std::result::Result::Err(error) => LivePersistResult::BackendError(error.kind()), }; let close_result = close_backend(backend).await; if let std::result::Result::Err(error) = close_result { return std::result::Result::Err(error); } return std::result::Result::Ok(result); } async fn transition_once( uri: std::string::String, signature_seed: u8, expected: ksp_store_api::RawRetentionState, target: ksp_store_api::RawRetentionState, ) -> std::result::Result { let backend = match open_backend(uri.as_str(), "devnet", true, true).await { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let reference = match raw_reference(signature_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let transition = match retention_transition(reference, expected, target) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let transitioned = backend.transition_raw_transaction_retention(transition).await; let result = match transitioned { std::result::Result::Ok(value) => LiveRetentionResult::Outcome(value), std::result::Result::Err(_) => LiveRetentionResult::BackendError, }; let close_result = close_backend(backend).await; if let std::result::Result::Err(error) = close_result { return std::result::Result::Err(error); } return std::result::Result::Ok(result); } fn joined_persist( joined: std::result::Result, tokio::task::JoinError>, phase: &'static str, ) -> std::result::Result { return match joined { std::result::Result::Ok(std::result::Result::Ok(value)) => std::result::Result::Ok(value), _ => std::result::Result::Err(LiveFailure::new(phase)), }; } fn joined_retention( joined: std::result::Result, tokio::task::JoinError>, phase: &'static str, ) -> std::result::Result { return match joined { std::result::Result::Ok(std::result::Result::Ok(value)) => std::result::Result::Ok(value), _ => std::result::Result::Err(LiveFailure::new(phase)), }; } fn matches_inserted(value: &LivePersistResult) -> bool { return matches!( value, LivePersistResult::Outcome(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted && outcome.observation() == ksp_store_api::RawObservationWriteOutcome::Inserted ); } fn matches_already_present(value: &LivePersistResult) -> bool { return matches!( value, LivePersistResult::Outcome(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::AlreadyPresent && outcome.observation() == ksp_store_api::RawObservationWriteOutcome::AlreadyPresent ); } fn one_applied_one_already(values: [LiveRetentionResult; 2]) -> bool { let applied = values.iter().filter(|value| matches!(value, LiveRetentionResult::Outcome(ksp_store_api::RawRetentionWriteOutcome::Applied))).count(); let already = values .iter() .filter(|value| matches!(value, LiveRetentionResult::Outcome(ksp_store_api::RawRetentionWriteOutcome::AlreadyAtTarget))) .count(); let errors = values.iter().filter(|value| matches!(value, LiveRetentionResult::BackendError)).count(); return applied == 1 && already == 1 && errors == 0; } fn raw_reference(signature_seed: u8) -> std::result::Result { let network = match network("devnet") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return std::result::Result::Ok(ksp_store_api::RawTransactionReference::new(network, ksp_store_api::RawTransactionSignature::new([signature_seed; 64]))); } fn raw_transaction(signature_seed: u8, slot: u64, payload_seed: u8) -> std::result::Result { let reference = match raw_reference(signature_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_id = match ksp_store_api::RawFormatId::new("ksp-live-raw-v1") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_format")), }; let content_hash = ksp_store_api::RawContentHash::new([payload_seed.wrapping_add(100); 32]); let payload_bytes = [payload_seed, payload_seed.wrapping_add(1), payload_seed.wrapping_add(2), payload_seed.wrapping_add(3)].to_vec().into_boxed_slice(); let payload = match ksp_store_api::RawPayload::try_new(format_id, 1, payload_bytes, content_hash) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_payload")), }; let timestamp = match ksp_store_api::RawTimestamp::from_unix_millis(1_700_000_000_000_u64 + slot) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_block_time")), }; return std::result::Result::Ok(ksp_store_api::RawTransaction::new(reference, slot, std::option::Option::Some(timestamp), payload)); } fn raw_observation( observation_seed: u8, signature_seed: u8, provenance_seed: u8, ) -> std::result::Result { let transaction = match raw_reference(signature_seed) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let provider = match provenance_code("live-provider") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let protocol = match provenance_code("yellowstone") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let method = match provenance_code("transaction-subscribe") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let received = match ksp_store_api::RawTimestamp::from_unix_millis(1_800_000_000_000_u64 + u64::from(provenance_seed)) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_received_at")), }; let observed = match ksp_store_api::RawTimestamp::from_unix_millis(1_799_999_999_000_u64 + u64::from(provenance_seed)) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_observed_at")), }; let capture_session = match provenance_code("session-live") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let commitment = match provenance_code("confirmed") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let endpoint = match provenance_code("endpoint-live") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let filter = match provenance_code("filter-live") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let provenance = ksp_store_api::RawAcquisitionProvenance::new(provider, protocol, method, ksp_store_api::RawAcquisitionOrigin::Live, received) .with_capture_session_id(capture_session) .with_commitment(commitment) .with_endpoint_id(endpoint) .with_filter_id(filter) .with_source_payload_hash(ksp_store_api::RawContentHash::new([provenance_seed.wrapping_add(50); 32])); let provenance = match provenance.try_with_observed_at(observed) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_observed_order")), }; let provenance = match provenance.try_with_source_payload_size_bytes(1_024_u64 + u64::from(provenance_seed)) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("model_source_size")), }; return std::result::Result::Ok(ksp_store_api::RawTransactionObservation::new( ksp_store_api::RawObservationKey::new([observation_seed; 32]), transaction, provenance, )); } fn provenance_code(value: &str) -> std::result::Result { return match ksp_store_api::RawProvenanceCode::new(value) { std::result::Result::Ok(code) => std::result::Result::Ok(code), std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("model_provenance_code")), }; } fn retention_transition( reference: ksp_store_api::RawTransactionReference, expected: ksp_store_api::RawRetentionState, target: ksp_store_api::RawRetentionState, ) -> std::result::Result { return match ksp_store_api::RawTransactionRetentionTransition::try_new(reference, expected, target) { std::result::Result::Ok(value) => std::result::Result::Ok(value), std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("model_retention_transition")), }; } fn transaction_matches(transaction: &ksp_store_api::RawTransaction, signature_seed: u8, slot: u64, payload_seed: u8) -> bool { let expected_signature = ksp_store_api::RawTransactionSignature::new([signature_seed; 64]); let expected_hash = ksp_store_api::RawContentHash::new([payload_seed.wrapping_add(100); 32]); let expected_bytes = [payload_seed, payload_seed.wrapping_add(1), payload_seed.wrapping_add(2), payload_seed.wrapping_add(3)]; let expected_time = 1_700_000_000_000_u64 + slot; return transaction.reference().network().as_str() == "devnet" && transaction.reference().signature() == expected_signature && transaction.slot() == slot && transaction.block_time().map(|value| return value.unix_millis()) == std::option::Option::Some(expected_time) && transaction.payload().format_id().as_str() == "ksp-live-raw-v1" && transaction.payload().format_version() == 1 && transaction.payload().content_hash() == expected_hash && transaction.payload().bytes() == expected_bytes; } fn network(value: &str) -> std::result::Result { return match ksp_store_api::RawNetworkId::new(value) { std::result::Result::Ok(network) => std::result::Result::Ok(network), std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("network")), }; } fn settings( uri: &str, network_name: &str, schema_autocreate: bool, schema_autoupdate: bool, ) -> std::result::Result { let network = match network(network_name) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return std::result::Result::Ok(ksp_store_postgres_lib::PostgresBackendSettings::with_schema_policy( network, uri, 8, 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, schema_autocreate, schema_autoupdate, std::time::Duration::from_secs(30), std::time::Duration::from_secs(10), )); } async fn open_backend( uri: &str, network_name: &str, schema_autocreate: bool, schema_autoupdate: bool, ) -> std::result::Result { let opened = open_backend_result(uri, network_name, schema_autocreate, schema_autoupdate).await; return match opened { std::result::Result::Ok(std::result::Result::Ok(value)) => std::result::Result::Ok(value), std::result::Result::Ok(std::result::Result::Err(error)) => std::result::Result::Err(LiveFailure::new(error.phase())), std::result::Result::Err(error) => std::result::Result::Err(error), }; } async fn open_backend_result( uri: &str, network_name: &str, schema_autocreate: bool, schema_autoupdate: bool, ) -> std::result::Result, LiveFailure> { let backend_settings = match settings(uri, network_name, schema_autocreate, schema_autoupdate) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return std::result::Result::Ok(ksp_store_postgres_lib::PostgresBackend::open(backend_settings).await); } 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 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 = match row.try_get::(0) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("server_version_decode")), }; let version_num = match value.parse::() { 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 = match client.query_one(LIVE_MANAGED_SCHEMA_EXISTS_SQL, &[]).await { 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 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 index_exists(client: &tokio_postgres::Client) -> std::result::Result { let row = match client.query_one(LIVE_INDEX_EXISTS_SQL, &[]).await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("index_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("index_probe_decode")), }; } async fn transaction_exists(client: &tokio_postgres::Client, reference: &ksp_store_api::RawTransactionReference) -> std::result::Result { let signature = reference.signature(); let signature_bytes: &[u8] = signature.as_bytes(); let row = match client.query_one(LIVE_TRANSACTION_EXISTS_SQL, &[&signature_bytes]).await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("transaction_exists_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("transaction_exists_decode")), }; } async fn observation_exists(client: &tokio_postgres::Client, key: &ksp_store_api::RawObservationKey) -> std::result::Result { let key_bytes: &[u8] = key.as_bytes(); let row = match client.query_one(LIVE_OBSERVATION_EXISTS_SQL, &[&key_bytes]).await { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("observation_exists_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("observation_exists_decode")), }; }