Files
khadhroony-solana-project/crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs

1176 lines
61 KiB
Rust

// file: crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs
// version: 2
#![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<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 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(1) || 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(1) {
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| matches_inserted(value)).count();
let already = pair.iter().filter(|value| 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| 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(values) if values == [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(values) if values == [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<ksp_store_api::RawPageCursor> = 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<ksp_store_api::RawPageCursor>,
start: u64,
end: u64,
) -> std::result::Result<ksp_store_api::RawTransactionQuery, LiveFailure> {
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<ksp_store_api::RawPageCursor, LiveFailure> {
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<LivePersistResult, LiveFailure> {
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<LiveRetentionResult, LiveFailure> {
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<std::result::Result<LivePersistResult, LiveFailure>, tokio::task::JoinError>,
phase: &'static str,
) -> std::result::Result<LivePersistResult, LiveFailure> {
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<std::result::Result<LiveRetentionResult, LiveFailure>, tokio::task::JoinError>,
phase: &'static str,
) -> std::result::Result<LiveRetentionResult, LiveFailure> {
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<ksp_store_api::RawTransactionReference, LiveFailure> {
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<ksp_store_api::RawTransaction, LiveFailure> {
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<ksp_store_api::RawTransactionObservation, LiveFailure> {
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<ksp_store_api::RawProvenanceCode, LiveFailure> {
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<ksp_store_api::RawTransactionRetentionTransition, LiveFailure> {
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<ksp_store_api::RawNetworkId, LiveFailure> {
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<ksp_store_postgres_lib::PostgresBackendSettings, LiveFailure> {
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<ksp_store_postgres_lib::PostgresBackend, LiveFailure> {
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::Err(LiveFailure::new("backend_open")),
};
}
async fn open_backend_result(
uri: &str,
network_name: &str,
schema_autocreate: bool,
schema_autoupdate: bool,
) -> std::result::Result<std::result::Result<ksp_store_postgres_lib::PostgresBackend, ksp_store_postgres_lib::PostgresBackendError>, 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<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 = match row.try_get::<usize, std::string::String>(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::<u32>() {
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 = 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::<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 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<bool, LiveFailure> {
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::<usize, bool>(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<bool, LiveFailure> {
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::<usize, bool>(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<bool, LiveFailure> {
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::<usize, bool>(0) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("observation_exists_decode")),
};
}