Files
khadhroony-solana-project/crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs
2026-08-30 23:39:43 +02:00

1217 lines
62 KiB
Rust

// file: crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs
// version: 1
#![warn(missing_docs)]
#![deny(unreachable_pub)]
#![forbid(unsafe_code)]
//! Opt-in real PostgreSQL proof for the complete RawAccountState 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_ACCOUNT_INDEX_EXISTS_SQL: &str = r#"SELECT EXISTS (
SELECT 1 FROM pg_indexes
WHERE schemaname = current_schema()
AND indexname = 'ix_ksp_raw_account_states_slot_pubkey_state_hash'
)"#;
const LIVE_ACCOUNT_STATE_EXISTS_SQL: &str =
"SELECT EXISTS (SELECT 1 FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3)";
const LIVE_CANCEL_WAIT: std::time::Duration = std::time::Duration::from_millis(300);
const LIVE_LOCK_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key FROM ksp_raw_account_observations WHERE observation_key = $1 FOR UPDATE";
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_account_observations;
DROP TABLE IF EXISTS ksp_raw_account_states;
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',
'ksp_raw_account_states',
'ksp_raw_account_observations'
)
AND table_type = 'BASE TABLE'
)"#;
const LIVE_MAX_URI_BYTES: usize = 4_096;
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 LiveAccountPersistResult {
BackendError(ksp_store_postgres_lib::PostgresBackendErrorKind),
Outcome(ksp_store_api::RawAcquisitionWriteOutcome),
}
#[test]
#[ignore = "opt-in real PostgreSQL RawAccountState proof; reads one dedicated URI from stdin"]
fn pre_009_real_postgres_raw_account_vertical_slice_is_atomic_concurrent_and_recoverable() {
eprintln!("KSP Store RawAccountState 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 RawAccountState 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 RawAccountState 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 RawAccountState 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 RawAccountState live proof: server major {major}");
let mut owns_schema = false;
let scenario = run_raw_account_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_account_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_state_extremes(&backend).await,
prove_atomic_idempotence_and_conflicts(&backend).await,
prove_additional_observations(&backend).await,
prove_pagination_and_cursors(&backend).await,
prove_cross_family_coexistence(&backend, admin).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(uri).await;
if let std::result::Result::Err(error) = divergent_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 = match account_reference(10, 100, 10) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match reopened.get_raw_account_state(&reference).await {
std::result::Result::Ok(std::option::Option::Some(value)) if account_state_matches(&value, 10, 100, 10, 1_000, 11, false, 12, 4, 10) => {},
_ => return std::result::Result::Err(LiveFailure::new("reopen_account_read")),
}
let transaction_reference = match transaction_reference(90) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match reopened.get_raw_transaction(&transaction_reference).await {
std::result::Result::Ok(std::option::Option::Some(_)) => {},
_ => return std::result::Result::Err(LiveFailure::new("reopen_transaction_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_account_states_slot_pubkey_state_hash", &[]).await;
match drop_result {
std::result::Result::Ok(_) => {},
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("schema_drop_account_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 = account_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_account_index_not_repaired")),
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
return close_backend(repaired).await;
}
async fn prove_state_extremes(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> {
let max_state = match account_state(1, u64::MAX, 1, u64::MAX, 2, true, u64::MAX, 0, 1) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let max_observation = match account_observation(1, 1, u64::MAX, 1, 1, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let max_persist = backend.persist_raw_account_acquisition(max_state, max_observation).await;
match max_persist {
std::result::Result::Ok(value)
if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted
&& value.observation() == ksp_store_api::RawObservationWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("u64_max_insert")),
}
let max_reference = match account_reference(1, u64::MAX, 1) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.get_raw_account_state(&max_reference).await {
std::result::Result::Ok(std::option::Option::Some(value)) if account_state_matches(&value, 1, u64::MAX, 1, u64::MAX, 2, true, u64::MAX, 0, 1) => {},
_ => return std::result::Result::Err(LiveFailure::new("u64_max_read")),
}
let max_observation_key = ksp_store_api::RawObservationKey::new([1; 32]);
match backend.get_raw_account_observation(&max_observation_key).await {
std::result::Result::Ok(std::option::Option::Some(value))
if value.is_startup() == std::option::Option::Some(true)
&& value.write_version() == std::option::Option::Some(u64::MAX)
&& value.transaction_signature() == std::option::Option::Some(ksp_store_api::RawTransactionSignature::new([21; 64])) => {},
_ => return std::result::Result::Err(LiveFailure::new("u64_max_observation_read")),
}
let full_state = match account_state(2, 2, 2, 2, 3, false, 2, ksp_store_api::MAX_RAW_ACCOUNT_DATA_BYTES, 2) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let full_observation = match account_observation(2, 2, 2, 2, 2, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let full_persist = backend.persist_raw_account_acquisition(full_state, full_observation).await;
match full_persist {
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("max_data_insert")),
}
let full_reference = match account_reference(2, 2, 2) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.get_raw_account_state(&full_reference).await {
std::result::Result::Ok(std::option::Option::Some(value))
if account_state_matches(&value, 2, 2, 2, 2, 3, false, 2, ksp_store_api::MAX_RAW_ACCOUNT_DATA_BYTES, 2) => {},
_ => return std::result::Result::Err(LiveFailure::new("max_data_read")),
}
return std::result::Result::Ok(());
}
async fn prove_atomic_idempotence_and_conflicts(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> {
let state = match account_state(10, 100, 10, 1_000, 11, false, 12, 4, 10) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation = match account_observation(10, 10, 100, 10, 10, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let first = backend.persist_raw_account_acquisition(state, observation).await;
match first {
std::result::Result::Ok(value)
if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted
&& value.observation() == ksp_store_api::RawObservationWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("atomic_insert")),
}
let identical_state = match account_state(10, 100, 10, 1_000, 11, false, 12, 4, 10) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let identical_observation = match account_observation(10, 10, 100, 10, 10, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let second = backend.persist_raw_account_acquisition(identical_state, identical_observation).await;
match second {
std::result::Result::Ok(value)
if value.entity() == ksp_store_api::RawEntityWriteOutcome::AlreadyPresent
&& value.observation() == ksp_store_api::RawObservationWriteOutcome::AlreadyPresent => {},
_ => return std::result::Result::Err(LiveFailure::new("atomic_idempotent")),
}
let divergent_state = match account_state(10, 100, 10, 1_001, 11, false, 12, 4, 99) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let divergent_observation = match account_observation(11, 10, 100, 10, 11, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.persist_raw_account_acquisition(divergent_state, divergent_observation).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {},
_ => return std::result::Result::Err(LiveFailure::new("state_conflict")),
}
let same_slot_first = match account_state(20, 120, 20, 2_000, 21, false, 22, 3, 20) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let same_slot_first_observation = match account_observation(20, 20, 120, 20, 20, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let same_slot_second = match account_state(20, 120, 21, 2_001, 21, false, 22, 3, 21) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let same_slot_second_observation = match account_observation(21, 20, 120, 21, 21, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
for result in [
backend.persist_raw_account_acquisition(same_slot_first, same_slot_first_observation).await,
backend.persist_raw_account_acquisition(same_slot_second, same_slot_second_observation).await,
] {
match result {
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("same_slot_hash_insert")),
}
}
let collision_state = match account_state(30, 130, 30, 3_000, 31, false, 32, 3, 30) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let collision_observation = match account_observation(10, 30, 130, 30, 30, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.persist_raw_account_acquisition(collision_state, collision_observation).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {},
_ => return std::result::Result::Err(LiveFailure::new("observation_collision_conflict")),
}
let collision_reference = match account_reference(30, 130, 30) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.get_raw_account_state(&collision_reference).await {
std::result::Result::Ok(std::option::Option::None) => {},
_ => return std::result::Result::Err(LiveFailure::new("observation_collision_rollback")),
}
return std::result::Result::Ok(());
}
async fn prove_additional_observations(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> {
let additional = match account_observation(41, 10, 100, 10, 41, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let expected = additional.clone();
let first = backend.record_raw_account_observation(additional).await;
if first != std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::Inserted) {
return std::result::Result::Err(LiveFailure::new("additional_observation_insert"));
}
let identical = match account_observation(41, 10, 100, 10, 41, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let second = backend.record_raw_account_observation(identical).await;
if second != std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::AlreadyPresent) {
return std::result::Result::Err(LiveFailure::new("additional_observation_idempotent"));
}
let key = ksp_store_api::RawObservationKey::new([41; 32]);
match backend.get_raw_account_observation(&key).await {
std::result::Result::Ok(std::option::Option::Some(value)) if value == expected => {},
_ => return std::result::Result::Err(LiveFailure::new("additional_observation_round_trip")),
}
let divergent = match account_observation(41, 10, 100, 10, 42, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.record_raw_account_observation(divergent).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict => {},
_ => return std::result::Result::Err(LiveFailure::new("additional_observation_conflict")),
}
let missing = match account_observation(42, 99, 999, 99, 42, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.record_raw_account_observation(missing).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::ReferenceNotFound => {},
_ => return std::result::Result::Err(LiveFailure::new("missing_reference")),
}
return std::result::Result::Ok(());
}
async fn prove_pagination_and_cursors(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> {
let rows = [(50_u8, 200_u64, 50_u8), (50_u8, 200_u64, 51_u8), (50_u8, 201_u64, 52_u8), (51_u8, 201_u64, 53_u8), (52_u8, 202_u64, 54_u8)];
for (index, (pubkey_seed, slot, hash_seed)) in rows.into_iter().enumerate() {
let state = match account_state(pubkey_seed, slot, hash_seed, 5_000_u64 + slot, 60, false, 61, 2, hash_seed) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation_seed = 50_u8.wrapping_add(index as u8);
let observation = match account_observation(observation_seed, pubkey_seed, slot, hash_seed, observation_seed, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.persist_raw_account_acquisition(state, observation).await {
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("pagination_seed")),
}
}
let ascending = match collect_account_pages(backend, std::option::Option::None, 200, 202, ksp_store_api::RawSortDirection::Ascending).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let expected_ascending = [(200, 50, 50), (200, 50, 51), (201, 50, 52), (201, 51, 53), (202, 52, 54)];
if ascending.as_slice() != expected_ascending.as_slice() {
return std::result::Result::Err(LiveFailure::new("pagination_ascending"));
}
let descending = match collect_account_pages(backend, std::option::Option::None, 200, 202, ksp_store_api::RawSortDirection::Descending).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let expected_descending = [(202, 52, 54), (201, 51, 53), (201, 50, 52), (200, 50, 51), (200, 50, 50)];
if descending.as_slice() != expected_descending.as_slice() {
return std::result::Result::Err(LiveFailure::new("pagination_descending"));
}
let filtered = match collect_account_pages(
backend,
std::option::Option::Some(ksp_store_api::Pubkey::new_from_array([50; 32])),
200,
201,
ksp_store_api::RawSortDirection::Ascending,
)
.await
{
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let expected_filtered = [(200, 50, 50), (200, 50, 51), (201, 50, 52)];
if filtered.as_slice() != expected_filtered.as_slice() {
return std::result::Result::Err(LiveFailure::new("pagination_pubkey_filter"));
}
let first_query = match account_query(std::option::Option::None, 200, 202, ksp_store_api::RawSortDirection::Ascending, std::option::Option::None) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let first_page = match backend.list_raw_account_states(&first_query).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("cursor_first_page")),
};
let cursor = match first_page.next_cursor() {
std::option::Option::Some(value) => value.clone(),
std::option::Option::None => return std::result::Result::Err(LiveFailure::new("cursor_missing")),
};
let other_direction =
match account_query(std::option::Option::None, 200, 202, ksp_store_api::RawSortDirection::Descending, std::option::Option::Some(cursor.clone())) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.list_raw_account_states(&other_direction).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::QueryInvalid => {},
_ => return std::result::Result::Err(LiveFailure::new("cursor_cross_direction")),
}
let mut hostile_bytes = cursor.as_bytes().to_vec();
let last_index = match hostile_bytes.len().checked_sub(1) {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(LiveFailure::new("cursor_hostile_index")),
};
hostile_bytes[last_index] ^= 0xff;
let hostile_cursor = match ksp_store_api::RawPageCursor::try_new(hostile_bytes.into_boxed_slice()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("cursor_hostile_model")),
};
let hostile_query =
match account_query(std::option::Option::None, 200, 202, ksp_store_api::RawSortDirection::Ascending, std::option::Option::Some(hostile_cursor)) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.list_raw_account_states(&hostile_query).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::QueryInvalid => {},
_ => return std::result::Result::Err(LiveFailure::new("cursor_hostile_digest")),
}
let mut transaction_family_bytes = vec![0_u8; 109];
transaction_family_bytes[0..4].copy_from_slice(b"KSPT");
transaction_family_bytes[4] = 1;
let transaction_family_cursor = match ksp_store_api::RawPageCursor::try_new(transaction_family_bytes.into_boxed_slice()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("cursor_cross_family_model")),
};
let transaction_family_query = match account_query(
std::option::Option::None,
200,
202,
ksp_store_api::RawSortDirection::Ascending,
std::option::Option::Some(transaction_family_cursor),
) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.list_raw_account_states(&transaction_family_query).await {
std::result::Result::Err(error) if error.kind() == ksp_store_postgres_lib::PostgresBackendErrorKind::QueryInvalid => {},
_ => return std::result::Result::Err(LiveFailure::new("cursor_cross_family")),
}
return std::result::Result::Ok(());
}
async fn prove_cross_family_coexistence(
backend: &ksp_store_postgres_lib::PostgresBackend,
admin: &tokio_postgres::Client,
) -> std::result::Result<(), LiveFailure> {
let account_state = match account_state(90, 900, 90, 9_000, 91, false, 92, 3, 90) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let account_observation = match account_observation(90, 90, 900, 90, 90, true) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.persist_raw_account_acquisition(account_state, account_observation).await {
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("coexistence_account_insert")),
}
let transaction = match raw_transaction(90, 900, 90) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let transaction_observation = match raw_transaction_observation(90, 90, 90) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend
.persist_raw_transaction_acquisition(transaction, transaction_observation, ksp_store_api::RawTransactionAcquisitionMode::Normal)
.await
{
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("coexistence_transaction_insert")),
}
let transaction_reference = match transaction_reference(90) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match backend.get_raw_transaction(&transaction_reference).await {
std::result::Result::Ok(std::option::Option::Some(_)) => {},
_ => return std::result::Result::Err(LiveFailure::new("coexistence_transaction_read")),
}
let account_key = ksp_store_api::RawObservationKey::new([90; 32]);
match backend.get_raw_account_observation(&account_key).await {
std::result::Result::Ok(std::option::Option::Some(_)) => {},
_ => return std::result::Result::Err(LiveFailure::new("coexistence_account_observation")),
}
let transaction_exists_result = transaction_exists(admin, &transaction_reference).await;
match transaction_exists_result {
std::result::Result::Ok(true) => {},
_ => return std::result::Result::Err(LiveFailure::new("coexistence_transaction_physical")),
}
return std::result::Result::Ok(());
}
async fn prove_concurrent_identical_insert(uri: &str) -> std::result::Result<(), LiveFailure> {
let first = tokio::spawn(persist_account_once(uri.to_owned(), 70, 700, 70, 7_000, 70, 70));
let second = tokio::spawn(persist_account_once(uri.to_owned(), 70, 700, 70, 7_000, 70, 70));
let first_result = match joined_account_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_account_persist(second.await, "concurrent_identical_second") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if !one_inserted_one_already([first_result, second_result]) {
return std::result::Result::Err(LiveFailure::new("concurrent_identical_outcome"));
}
return std::result::Result::Ok(());
}
async fn prove_concurrent_divergent_insert(uri: &str) -> std::result::Result<(), LiveFailure> {
let first = tokio::spawn(persist_account_once(uri.to_owned(), 71, 710, 71, 7_100, 71, 71));
let second = tokio::spawn(persist_account_once(uri.to_owned(), 71, 710, 71, 7_101, 72, 72));
let first_result = match joined_account_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_account_persist(second.await, "concurrent_divergent_second") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let inserted = [&first_result, &second_result]
.into_iter()
.filter(|value| matches!(value, LiveAccountPersistResult::Outcome(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted))
.count();
let conflicts = [&first_result, &second_result]
.into_iter()
.filter(|value| matches!(value, LiveAccountPersistResult::BackendError(ksp_store_postgres_lib::PostgresBackendErrorKind::Conflict)))
.count();
if inserted != 1 || conflicts != 1 {
return std::result::Result::Err(LiveFailure::new("concurrent_divergent_outcome"));
}
return std::result::Result::Ok(());
}
async fn prove_cancellation_rollback(admin: &mut tokio_postgres::Client, uri: &str) -> std::result::Result<(), LiveFailure> {
let seed_state = match account_state(80, 800, 80, 8_000, 81, false, 82, 3, 80) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let seed_observation = match account_observation(80, 80, 800, 80, 80, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let seed_backend = match open_backend(uri, "devnet", true, true).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let seeded = seed_backend.persist_raw_account_acquisition(seed_state, seed_observation).await;
match seeded {
std::result::Result::Ok(value) if value.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted => {},
_ => return std::result::Result::Err(LiveFailure::new("cancellation_seed")),
}
let close_seed = close_backend(seed_backend).await;
if let std::result::Result::Err(error) = close_seed {
return std::result::Result::Err(error);
}
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_ACCOUNT_OBSERVATION_SQL, &[&key_bytes]).await;
if lock_result.is_err() {
return std::result::Result::Err(LiveFailure::new("cancellation_lock_observation"));
}
let cancel_backend = match open_backend(uri, "devnet", true, true).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let state = match account_state(81, 810, 81, 8_100, 82, false, 83, 3, 81) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation = match account_observation(80, 81, 810, 81, 81, false) {
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_account_acquisition(state, observation).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 account_reference(81, 810, 81) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let exists_result = account_state_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_state")),
std::result::Result::Err(error) => return std::result::Result::Err(error),
}
return std::result::Result::Ok(());
}
async fn collect_account_pages(
backend: &ksp_store_postgres_lib::PostgresBackend,
pubkey: std::option::Option<ksp_store_api::Pubkey>,
start: u64,
end: u64,
direction: ksp_store_api::RawSortDirection,
) -> std::result::Result<std::vec::Vec<(u64, u8, u8)>, LiveFailure> {
let mut output = std::vec::Vec::new();
let mut cursor = std::option::Option::None;
loop {
let query = match account_query(pubkey, start, end, direction, cursor) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let page = match backend.list_raw_account_states(&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() {
let pubkey_bytes: &[u8] = reference.pubkey().as_ref();
let pubkey_seed = match pubkey_bytes.first() {
std::option::Option::Some(value) => *value,
std::option::Option::None => return std::result::Result::Err(LiveFailure::new("pagination_pubkey_decode")),
};
let hash_seed = reference.state_hash().as_bytes()[0];
output.push((reference.slot(), pubkey_seed, hash_seed));
}
cursor = page.next_cursor().cloned();
if cursor.is_none() {
break;
}
}
return std::result::Result::Ok(output);
}
async fn persist_account_once(
uri: std::string::String,
pubkey_seed: u8,
slot: u64,
hash_seed: u8,
lamports: u64,
data_seed: u8,
observation_seed: u8,
) -> std::result::Result<LiveAccountPersistResult, 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 state = match account_state(pubkey_seed, slot, hash_seed, lamports, 100, false, 101, 4, data_seed) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation = match account_observation(observation_seed, pubkey_seed, slot, hash_seed, observation_seed, false) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let persisted = backend.persist_raw_account_acquisition(state, observation).await;
let result = match persisted {
std::result::Result::Ok(value) => LiveAccountPersistResult::Outcome(value),
std::result::Result::Err(error) => LiveAccountPersistResult::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);
}
fn joined_account_persist(
joined: std::result::Result<std::result::Result<LiveAccountPersistResult, LiveFailure>, tokio::task::JoinError>,
phase: &'static str,
) -> std::result::Result<LiveAccountPersistResult, 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 one_inserted_one_already(values: [LiveAccountPersistResult; 2]) -> bool {
let inserted = values
.iter()
.filter(|value| matches!(value, LiveAccountPersistResult::Outcome(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::Inserted))
.count();
let already = values
.iter()
.filter(|value| matches!(value, LiveAccountPersistResult::Outcome(outcome) if outcome.entity() == ksp_store_api::RawEntityWriteOutcome::AlreadyPresent))
.count();
let errors = values.iter().filter(|value| matches!(value, LiveAccountPersistResult::BackendError(_))).count();
return inserted == 1 && already == 1 && errors == 0;
}
fn account_query(
pubkey: std::option::Option<ksp_store_api::Pubkey>,
start: u64,
end: u64,
direction: ksp_store_api::RawSortDirection,
cursor: std::option::Option<ksp_store_api::RawPageCursor>,
) -> std::result::Result<ksp_store_api::RawAccountStateQuery, LiveFailure> {
let network = match network("devnet") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let slots = 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("model_slot_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("model_page_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::RawAccountStateQuery::new(network, pubkey, slots, direction, page));
}
fn account_reference(pubkey_seed: u8, slot: u64, hash_seed: u8) -> std::result::Result<ksp_store_api::RawAccountStateReference, 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::RawAccountStateReference::new(
network,
ksp_store_api::Pubkey::new_from_array([pubkey_seed; 32]),
slot,
ksp_store_api::RawContentHash::new([hash_seed; 32]),
));
}
fn account_state(
pubkey_seed: u8,
slot: u64,
hash_seed: u8,
lamports: u64,
owner_seed: u8,
executable: bool,
rent_epoch: u64,
data_len: usize,
data_seed: u8,
) -> std::result::Result<ksp_store_api::RawAccountState, LiveFailure> {
let reference = match account_reference(pubkey_seed, slot, hash_seed) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let data = vec![data_seed; data_len].into_boxed_slice();
return match ksp_store_api::RawAccountState::try_new(
reference,
lamports,
ksp_store_api::Pubkey::new_from_array([owner_seed; 32]),
executable,
rent_epoch,
data,
) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("model_account_state")),
};
}
fn account_observation(
observation_seed: u8,
pubkey_seed: u8,
slot: u64,
hash_seed: u8,
provenance_seed: u8,
yellowstone_metadata: bool,
) -> std::result::Result<ksp_store_api::RawAccountObservation, LiveFailure> {
let reference = match account_reference(pubkey_seed, slot, hash_seed) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let provenance = match acquisition_provenance(provenance_seed, "account-subscribe") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation = ksp_store_api::RawAccountObservation::new(ksp_store_api::RawObservationKey::new([observation_seed; 32]), reference, provenance);
if !yellowstone_metadata {
return std::result::Result::Ok(observation);
}
return std::result::Result::Ok(
observation
.with_is_startup(true)
.with_transaction_signature(ksp_store_api::RawTransactionSignature::new([provenance_seed.wrapping_add(20); 64]))
.with_write_version(if slot == u64::MAX { u64::MAX } else { 10_000_u64 + u64::from(provenance_seed) }),
);
}
fn acquisition_provenance(provenance_seed: u8, method_name: &str) -> std::result::Result<ksp_store_api::RawAcquisitionProvenance, LiveFailure> {
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-grpc") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let method = match provenance_code(method_name) {
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")),
};
return match provenance.try_with_source_payload_size_bytes(1_024_u64 + u64::from(provenance_seed)) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("model_source_size")),
};
}
fn account_state_matches(
state: &ksp_store_api::RawAccountState,
pubkey_seed: u8,
slot: u64,
hash_seed: u8,
lamports: u64,
owner_seed: u8,
executable: bool,
rent_epoch: u64,
data_len: usize,
data_seed: u8,
) -> bool {
return state.reference().network().as_str() == "devnet"
&& state.reference().pubkey() == &ksp_store_api::Pubkey::new_from_array([pubkey_seed; 32])
&& state.reference().slot() == slot
&& state.reference().state_hash() == ksp_store_api::RawContentHash::new([hash_seed; 32])
&& state.lamports() == lamports
&& state.owner() == &ksp_store_api::Pubkey::new_from_array([owner_seed; 32])
&& state.executable() == executable
&& state.rent_epoch() == rent_epoch
&& state.data_len() == data_len
&& state.data().iter().all(|value| return *value == data_seed);
}
fn raw_transaction(signature_seed: u8, slot: u64, payload_seed: u8) -> std::result::Result<ksp_store_api::RawTransaction, LiveFailure> {
let reference = match transaction_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_transaction_observation(
observation_seed: u8,
signature_seed: u8,
provenance_seed: u8,
) -> std::result::Result<ksp_store_api::RawTransactionObservation, LiveFailure> {
let transaction = match transaction_reference(signature_seed) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let provenance = match acquisition_provenance(provenance_seed, "transaction-subscribe") {
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::RawTransactionObservation::new(
ksp_store_api::RawObservationKey::new([observation_seed; 32]),
transaction,
provenance,
));
}
fn transaction_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 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 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::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<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 account_index_exists(client: &tokio_postgres::Client) -> std::result::Result<bool, LiveFailure> {
let row = match client.query_one(LIVE_ACCOUNT_INDEX_EXISTS_SQL, &[]).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("account_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("account_index_probe_decode")),
};
}
async fn account_state_exists(client: &tokio_postgres::Client, reference: &ksp_store_api::RawAccountStateReference) -> std::result::Result<bool, LiveFailure> {
let pubkey_bytes: &[u8] = reference.pubkey().as_ref();
let slot_text = reference.slot().to_string();
let state_hash = reference.state_hash();
let state_hash_bytes: &[u8] = state_hash.as_bytes();
let row = match client.query_one(LIVE_ACCOUNT_STATE_EXISTS_SQL, &[&pubkey_bytes, &slot_text, &state_hash_bytes]).await {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("account_state_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("account_state_exists_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")),
};
}