// file: crates/ksp-store-postgres-lib/src/raw_account.rs // version: 5 pub(crate) mod cursor; const GET_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key, account_pubkey, account_slot::text AS account_slot_text, account_state_hash, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes, is_startup, transaction_signature, write_version::text AS write_version_text FROM ksp_raw_account_observations WHERE observation_key = $1"; const GET_ACCOUNT_STATE_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash, lamports::text AS lamports_text, owner, executable, rent_epoch::text AS rent_epoch_text, data FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3"; const INSERT_ACCOUNT_OBSERVATION_SQL: &str = "INSERT INTO ksp_raw_account_observations (observation_key, account_pubkey, account_slot, account_state_hash, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes, is_startup, transaction_signature, write_version) VALUES ($1, $2, $3::TEXT::NUMERIC, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19::TEXT::NUMERIC) ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key"; const INSERT_ACCOUNT_STATE_SQL: &str = "INSERT INTO ksp_raw_account_states (pubkey, slot, state_hash, lamports, owner, executable, rent_epoch, data) VALUES ($1, $2::TEXT::NUMERIC, $3, $4::TEXT::NUMERIC, $5, $6, $7::TEXT::NUMERIC, $8) ON CONFLICT (pubkey, slot, state_hash) DO NOTHING RETURNING pubkey"; const LIST_ACCOUNT_STATES_ASC_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash FROM ksp_raw_account_states WHERE ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, pubkey, state_hash) > ($3::TEXT::NUMERIC, $4::BYTEA, $5::BYTEA)) ORDER BY slot ASC, pubkey ASC, state_hash ASC LIMIT $6"; const LIST_ACCOUNT_STATES_BY_PUBKEY_ASC_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash FROM ksp_raw_account_states WHERE pubkey = $1 AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC) AND ($4::TEXT IS NULL OR (slot, pubkey, state_hash) > ($4::TEXT::NUMERIC, $5::BYTEA, $6::BYTEA)) ORDER BY slot ASC, pubkey ASC, state_hash ASC LIMIT $7"; const LIST_ACCOUNT_STATES_BY_PUBKEY_DESC_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash FROM ksp_raw_account_states WHERE pubkey = $1 AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC) AND ($4::TEXT IS NULL OR (slot, pubkey, state_hash) < ($4::TEXT::NUMERIC, $5::BYTEA, $6::BYTEA)) ORDER BY slot DESC, pubkey DESC, state_hash DESC LIMIT $7"; const LIST_ACCOUNT_STATES_DESC_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash FROM ksp_raw_account_states WHERE ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, pubkey, state_hash) < ($3::TEXT::NUMERIC, $4::BYTEA, $5::BYTEA)) ORDER BY slot DESC, pubkey DESC, state_hash DESC LIMIT $6"; const LOCK_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key, account_pubkey, account_slot::text AS account_slot_text, account_state_hash, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes, is_startup, transaction_signature, write_version::text AS write_version_text FROM ksp_raw_account_observations WHERE observation_key = $1 FOR UPDATE"; const LOCK_ACCOUNT_REFERENCE_SQL: &str = "SELECT 1 FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3 FOR KEY SHARE"; const LOCK_ACCOUNT_STATE_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash, lamports::text AS lamports_text, owner, executable, rent_epoch::text AS rent_epoch_text, data FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3 FOR UPDATE"; struct RawAccountListDbRow { pubkey: std::vec::Vec, slot_text: std::string::String, state_hash: std::vec::Vec, } struct RawAccountObservationDbRow { account_pubkey: std::vec::Vec, account_slot_text: std::string::String, account_state_hash: std::vec::Vec, acquisition_method: std::string::String, capture_session_id: std::option::Option, commitment: std::option::Option, endpoint_id: std::option::Option, filter_id: std::option::Option, is_startup: std::option::Option, observation_key: std::vec::Vec, observed_at_unix_millis: std::option::Option, origin: std::string::String, protocol: std::string::String, provider: std::string::String, received_at_unix_millis: i64, source_payload_hash: std::option::Option>, source_payload_size_bytes: std::option::Option, transaction_signature: std::option::Option>, write_version_text: std::option::Option, } struct RawAccountStateDbRow { data: std::vec::Vec, executable: bool, lamports_text: std::string::String, owner: std::vec::Vec, pubkey: std::vec::Vec, rent_epoch_text: std::string::String, slot_text: std::string::String, state_hash: std::vec::Vec, } /// Reads one complete canonical RAW account state from the physical PostgreSQL backend. pub(crate) async fn get_raw_account_state( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawAccountStateReference, ) -> std::result::Result, crate::PostgresBackendError> { let network_result = ensure_network(network, reference, "raw_account_state_network"); if let std::result::Result::Err(error) = network_result { return std::result::Result::Err(error); } let client_result = pool.get().await; let client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; 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 rows_result = client.query(GET_ACCOUNT_STATE_SQL, &[&pubkey_bytes, &slot_text, &state_hash_bytes]).await; let rows = match rows_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => { return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReadFailed, "raw_account_state_query")); }, }; if rows.is_empty() { return std::result::Result::Ok(std::option::Option::None); } if rows.len() != 1 { return std::result::Result::Err(data_invalid("raw_account_state_cardinality")); } let row = match rows.first() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_state_cardinality")), }; let physical = match raw_account_state_db_row(row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let decoded = match decode_raw_account_state_row(network, physical) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if decoded.reference() != reference { return std::result::Result::Err(data_invalid("raw_account_state_reference")); } return std::result::Result::Ok(std::option::Option::Some(decoded)); } /// Lists deterministic canonical RAW account-state references using PostgreSQL keyset pagination. pub(crate) async fn list_raw_account_states( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, query: &ksp_store_api::RawAccountStateQuery, ) -> std::result::Result, crate::PostgresBackendError> { if query.network() != network { return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, "raw_account_list_network")); } let (requested_usize, sql_limit) = match crate::raw_account_physical_page_limit(query.page().limit().get()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let decoded_cursor = match query.page().cursor() { std::option::Option::Some(value) => match crate::decode_raw_account_cursor(query, value) { std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => std::option::Option::None, }; let client_result = pool.get().await; let client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; let slots = query.slots(); let start_text = slots.start_inclusive().map(|value| return value.to_string()); let end_text = slots.end_inclusive().map(|value| return value.to_string()); let cursor_slot_text = decoded_cursor.as_ref().map(|value| return value.last_slot.to_string()); let cursor_pubkey = decoded_cursor.as_ref().map(|value| return value.last_pubkey.to_vec()); let cursor_state_hash = decoded_cursor.as_ref().map(|value| return value.last_state_hash.to_vec()); let rows_result = match query.pubkey() { std::option::Option::Some(pubkey) => { let pubkey_bytes: &[u8] = pubkey.as_ref(); let sql = match query.direction() { ksp_store_api::RawSortDirection::Ascending => LIST_ACCOUNT_STATES_BY_PUBKEY_ASC_SQL, ksp_store_api::RawSortDirection::Descending => LIST_ACCOUNT_STATES_BY_PUBKEY_DESC_SQL, _ => { return std::result::Result::Err(crate::PostgresBackendError::new( crate::PostgresBackendErrorKind::QueryInvalid, "raw_account_list_direction", )); }, }; client.query(sql, &[&pubkey_bytes, &start_text, &end_text, &cursor_slot_text, &cursor_pubkey, &cursor_state_hash, &sql_limit]).await }, std::option::Option::None => { let sql = match query.direction() { ksp_store_api::RawSortDirection::Ascending => LIST_ACCOUNT_STATES_ASC_SQL, ksp_store_api::RawSortDirection::Descending => LIST_ACCOUNT_STATES_DESC_SQL, _ => { return std::result::Result::Err(crate::PostgresBackendError::new( crate::PostgresBackendErrorKind::QueryInvalid, "raw_account_list_direction", )); }, }; client.query(sql, &[&start_text, &end_text, &cursor_slot_text, &cursor_pubkey, &cursor_state_hash, &sql_limit]).await }, }; let rows = match rows_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => { return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReadFailed, "raw_account_list_query")); }, }; let mut decoded = std::vec::Vec::with_capacity(rows.len()); for row in rows { let physical = match raw_account_list_db_row(&row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let item = match decode_raw_account_list_row(network, physical) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; decoded.push(item); } let has_more = decoded.len() > requested_usize; if has_more { decoded.truncate(requested_usize); } let next_cursor = if has_more { let last = match decoded.last() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_list_page")), }; match crate::encode_raw_account_cursor(query, last.0, last.1.pubkey(), &last.1.state_hash()) { std::result::Result::Ok(value) => std::option::Option::Some(value), std::result::Result::Err(error) => return std::result::Result::Err(error), } } else { std::option::Option::None }; let items = decoded.into_iter().map(|value| return value.1).collect(); return std::result::Result::Ok(ksp_store_api::RawPage::new(items, next_cursor)); } /// Reads one persisted RAW account observation by producer-owned idempotence key. pub(crate) async fn get_raw_account_observation( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, observation_key: &ksp_store_api::RawObservationKey, ) -> std::result::Result, crate::PostgresBackendError> { let client_result = pool.get().await; let client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; let observation_key_bytes: &[u8] = observation_key.as_bytes(); let rows_result = client.query(GET_ACCOUNT_OBSERVATION_SQL, &[&observation_key_bytes]).await; let rows = match rows_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => { return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReadFailed, "raw_account_observation_query")); }, }; if rows.is_empty() { return std::result::Result::Ok(std::option::Option::None); } if rows.len() != 1 { return std::result::Result::Err(data_invalid("raw_account_observation_cardinality")); } let row = match rows.first() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_observation_cardinality")), }; let physical = match raw_account_observation_db_row(row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let decoded = match decode_raw_account_observation_row(network, physical) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if decoded.observation_key() != *observation_key { return std::result::Result::Err(data_invalid("raw_account_observation_key")); } return std::result::Result::Ok(std::option::Option::Some(decoded)); } /// Persists one canonical RAW account state and its acquisition observation atomically. pub(crate) async fn persist_raw_account_acquisition( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, state: ksp_store_api::RawAccountState, observation: ksp_store_api::RawAccountObservation, ) -> std::result::Result { let input_result = ensure_acquisition_inputs(network, &state, &observation); if let std::result::Result::Err(error) = input_result { return std::result::Result::Err(error); } let client_result = pool.get().await; let mut client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; let sql_transaction_result = client.transaction().await; let sql_transaction = match sql_transaction_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_acquisition_begin")), }; let insert_result = insert_account_state(&sql_transaction, &state).await; let inserted = match insert_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let entity_outcome = if inserted { ksp_store_api::RawEntityWriteOutcome::Inserted } else { let locked_result = load_locked_account_state(&sql_transaction, network, state.reference()).await; let locked = match locked_result { std::result::Result::Ok(std::option::Option::Some(value)) => value, std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(data_invalid("raw_account_acquisition_conflict_missing")), std::result::Result::Err(error) => return std::result::Result::Err(error), }; if !raw_account_states_equal(&locked, &state) { return std::result::Result::Err(conflict("raw_account_acquisition_content_conflict")); } ksp_store_api::RawEntityWriteOutcome::AlreadyPresent }; let observation_result = persist_account_observation_row(&sql_transaction, network, &observation).await; let observation_outcome = match observation_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let commit_result = sql_transaction.commit().await; if commit_result.is_err() { return std::result::Result::Err(write_failed("raw_account_acquisition_commit")); } return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new(entity_outcome, observation_outcome)); } /// Persists one additional acquisition observation for an already durable RAW account state. pub(crate) async fn record_raw_account_observation( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, observation: ksp_store_api::RawAccountObservation, ) -> std::result::Result { let input_result = ensure_observation_write_input(network, &observation); if let std::result::Result::Err(error) = input_result { return std::result::Result::Err(error); } let client_result = pool.get().await; let mut client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; let sql_transaction_result = client.transaction().await; let sql_transaction = match sql_transaction_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_observation_begin")), }; let account = observation.account(); let account_pubkey_bytes: &[u8] = account.pubkey().as_ref(); let account_slot_text = account.slot().to_string(); let account_state_hash = account.state_hash(); let account_state_hash_bytes: &[u8] = account_state_hash.as_bytes(); let reference_result = sql_transaction .query_opt(LOCK_ACCOUNT_REFERENCE_SQL, &[&account_pubkey_bytes, &account_slot_text.as_str(), &account_state_hash_bytes]) .await; match reference_result { std::result::Result::Ok(std::option::Option::Some(_)) => {}, std::result::Result::Ok(std::option::Option::None) => { return std::result::Result::Err(reference_not_found("raw_account_observation_reference")); }, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_observation_lock_reference")), } let observation_result = persist_account_observation_row(&sql_transaction, network, &observation).await; let outcome = match observation_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let commit_result = sql_transaction.commit().await; if commit_result.is_err() { return std::result::Result::Err(write_failed("raw_account_observation_commit")); } return std::result::Result::Ok(outcome); } fn raw_account_list_db_row(row: &tokio_postgres::Row) -> std::result::Result { let pubkey = match row.try_get::<_, std::vec::Vec>("pubkey") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_list_decode")), }; let slot_text = match row.try_get::<_, std::string::String>("slot_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_list_decode")), }; let state_hash = match row.try_get::<_, std::vec::Vec>("state_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_list_decode")), }; return std::result::Result::Ok(RawAccountListDbRow { pubkey, slot_text, state_hash }); } fn decode_raw_account_list_row( network: &ksp_store_api::RawNetworkId, row: RawAccountListDbRow, ) -> std::result::Result<(u64, ksp_store_api::RawAccountStateReference), crate::PostgresBackendError> { let pubkey = match fixed_bytes::<32>(row.pubkey, "raw_account_list_pubkey") { std::result::Result::Ok(value) => ksp_store_api::Pubkey::new_from_array(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let slot = match decode_u64_decimal(row.slot_text.as_str(), "raw_account_list_slot") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let state_hash = match fixed_bytes::<32>(row.state_hash, "raw_account_list_state_hash") { std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let reference = ksp_store_api::RawAccountStateReference::new(network.clone(), pubkey, slot, state_hash); return std::result::Result::Ok((slot, reference)); } fn raw_account_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result { let account_pubkey: std::vec::Vec = match row.try_get("account_pubkey") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let account_slot_text: std::string::String = match row.try_get("account_slot_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let account_state_hash: std::vec::Vec = match row.try_get("account_state_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let acquisition_method: std::string::String = match row.try_get("acquisition_method") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let capture_session_id: std::option::Option = match row.try_get("capture_session_id") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let commitment: std::option::Option = match row.try_get("commitment") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let endpoint_id: std::option::Option = match row.try_get("endpoint_id") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let filter_id: std::option::Option = match row.try_get("filter_id") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let is_startup: std::option::Option = match row.try_get("is_startup") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let observation_key: std::vec::Vec = match row.try_get("observation_key") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let observed_at_unix_millis: std::option::Option = match row.try_get("observed_at_unix_millis") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let origin: std::string::String = match row.try_get("origin") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let protocol: std::string::String = match row.try_get("protocol") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let provider: std::string::String = match row.try_get("provider") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let received_at_unix_millis: i64 = match row.try_get("received_at_unix_millis") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let source_payload_hash: std::option::Option> = match row.try_get("source_payload_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let source_payload_size_bytes: std::option::Option = match row.try_get("source_payload_size_bytes") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let transaction_signature: std::option::Option> = match row.try_get("transaction_signature") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; let write_version_text: std::option::Option = match row.try_get("write_version_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_decode")), }; return std::result::Result::Ok(RawAccountObservationDbRow { account_pubkey, account_slot_text, account_state_hash, acquisition_method, capture_session_id, commitment, endpoint_id, filter_id, is_startup, observation_key, observed_at_unix_millis, origin, protocol, provider, received_at_unix_millis, source_payload_hash, source_payload_size_bytes, transaction_signature, write_version_text, }); } fn raw_account_state_db_row(row: &tokio_postgres::Row) -> std::result::Result { let data: std::vec::Vec = match row.try_get("data") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let executable: bool = match row.try_get("executable") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let lamports_text: std::string::String = match row.try_get("lamports_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let owner: std::vec::Vec = match row.try_get("owner") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let pubkey: std::vec::Vec = match row.try_get("pubkey") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let rent_epoch_text: std::string::String = match row.try_get("rent_epoch_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let slot_text: std::string::String = match row.try_get("slot_text") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; let state_hash: std::vec::Vec = match row.try_get("state_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_state_decode")), }; return std::result::Result::Ok(RawAccountStateDbRow { data, executable, lamports_text, owner, pubkey, rent_epoch_text, slot_text, state_hash }); } fn decode_raw_account_observation_row( network: &ksp_store_api::RawNetworkId, row: RawAccountObservationDbRow, ) -> std::result::Result { let observation_key = match fixed_bytes::<32>(row.observation_key, "raw_account_observation_key") { std::result::Result::Ok(value) => ksp_store_api::RawObservationKey::new(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let pubkey = match fixed_bytes::<32>(row.account_pubkey, "raw_account_observation_pubkey") { std::result::Result::Ok(value) => ksp_store_api::Pubkey::new_from_array(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let slot = match decode_u64_decimal(row.account_slot_text.as_str(), "raw_account_observation_slot") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let state_hash = match fixed_bytes::<32>(row.account_state_hash, "raw_account_observation_state_hash") { std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let provider = match decode_provenance_code(row.provider) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let protocol = match decode_provenance_code(row.protocol) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let acquisition_method = match decode_provenance_code(row.acquisition_method) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let origin = match decode_origin(row.origin.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let received_at = match decode_timestamp_i64(row.received_at_unix_millis, "raw_account_observation_received_at") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let mut provenance = ksp_store_api::RawAcquisitionProvenance::new(provider, protocol, acquisition_method, origin, received_at); provenance = match row.capture_session_id { std::option::Option::Some(value) => match decode_provenance_code(value) { std::result::Result::Ok(code) => provenance.with_capture_session_id(code), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => provenance, }; provenance = match row.commitment { std::option::Option::Some(value) => match decode_provenance_code(value) { std::result::Result::Ok(code) => provenance.with_commitment(code), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => provenance, }; provenance = match row.endpoint_id { std::option::Option::Some(value) => match decode_provenance_code(value) { std::result::Result::Ok(code) => provenance.with_endpoint_id(code), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => provenance, }; provenance = match row.filter_id { std::option::Option::Some(value) => match decode_provenance_code(value) { std::result::Result::Ok(code) => provenance.with_filter_id(code), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => provenance, }; provenance = match row.observed_at_unix_millis { std::option::Option::Some(value) => { let timestamp = match decode_timestamp_i64(value, "raw_account_observation_observed_at") { std::result::Result::Ok(decoded) => decoded, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match provenance.try_with_observed_at(timestamp) { std::result::Result::Ok(updated) => updated, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_time_order")), } }, std::option::Option::None => provenance, }; provenance = match row.source_payload_hash { std::option::Option::Some(value) => { let hash = match fixed_bytes::<32>(value, "raw_account_observation_source_hash") { std::result::Result::Ok(decoded) => ksp_store_api::RawContentHash::new(decoded), std::result::Result::Err(error) => return std::result::Result::Err(error), }; provenance.with_source_payload_hash(hash) }, std::option::Option::None => provenance, }; provenance = match row.source_payload_size_bytes { std::option::Option::Some(value) => { let size = match decode_u64_i64(value, "raw_account_observation_source_size") { std::result::Result::Ok(decoded) => decoded, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match provenance.try_with_source_payload_size_bytes(size) { std::result::Result::Ok(updated) => updated, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_source_size")), } }, std::option::Option::None => provenance, }; let reference = ksp_store_api::RawAccountStateReference::new(network.clone(), pubkey, slot, state_hash); let mut observation = ksp_store_api::RawAccountObservation::new(observation_key, reference, provenance); observation = match row.is_startup { std::option::Option::Some(value) => observation.with_is_startup(value), std::option::Option::None => observation, }; observation = match row.transaction_signature { std::option::Option::Some(value) => { let signature = match fixed_bytes::<64>(value, "raw_account_observation_transaction_signature") { std::result::Result::Ok(decoded) => ksp_store_api::RawTransactionSignature::new(decoded), std::result::Result::Err(error) => return std::result::Result::Err(error), }; observation.with_transaction_signature(signature) }, std::option::Option::None => observation, }; observation = match row.write_version_text { std::option::Option::Some(value) => match decode_u64_decimal(value.as_str(), "raw_account_observation_write_version") { std::result::Result::Ok(decoded) => observation.with_write_version(decoded), std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::option::Option::None => observation, }; return std::result::Result::Ok(observation); } fn decode_raw_account_state_row( network: &ksp_store_api::RawNetworkId, row: RawAccountStateDbRow, ) -> std::result::Result { let pubkey = match fixed_bytes::<32>(row.pubkey, "raw_account_state_pubkey") { std::result::Result::Ok(value) => ksp_store_api::Pubkey::new_from_array(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let slot = match decode_u64_decimal(row.slot_text.as_str(), "raw_account_state_slot") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let state_hash = match fixed_bytes::<32>(row.state_hash, "raw_account_state_hash") { std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let lamports = match decode_u64_decimal(row.lamports_text.as_str(), "raw_account_state_lamports") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let owner = match fixed_bytes::<32>(row.owner, "raw_account_state_owner") { std::result::Result::Ok(value) => ksp_store_api::Pubkey::new_from_array(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let rent_epoch = match decode_u64_decimal(row.rent_epoch_text.as_str(), "raw_account_state_rent_epoch") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let reference = ksp_store_api::RawAccountStateReference::new(network.clone(), pubkey, slot, state_hash); return match ksp_store_api::RawAccountState::try_new(reference, lamports, owner, row.executable, rent_epoch, row.data.into_boxed_slice()) { std::result::Result::Ok(value) => std::result::Result::Ok(value), std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_account_state_model")), }; } async fn insert_account_state( sql_transaction: &deadpool_postgres::Transaction<'_>, state: &ksp_store_api::RawAccountState, ) -> std::result::Result { let reference = state.reference(); 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 lamports_text = state.lamports().to_string(); let owner_bytes: &[u8] = state.owner().as_ref(); let executable = state.executable(); let rent_epoch_text = state.rent_epoch().to_string(); let data = state.data(); let row_result = sql_transaction .query_opt( INSERT_ACCOUNT_STATE_SQL, &[&pubkey_bytes, &slot_text.as_str(), &state_hash_bytes, &lamports_text.as_str(), &owner_bytes, &executable, &rent_epoch_text.as_str(), &data], ) .await; return match row_result { std::result::Result::Ok(std::option::Option::Some(_)) => std::result::Result::Ok(true), std::result::Result::Ok(std::option::Option::None) => std::result::Result::Ok(false), std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_account_acquisition_insert_state")), }; } async fn load_locked_account_state( sql_transaction: &deadpool_postgres::Transaction<'_>, network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawAccountStateReference, ) -> std::result::Result, crate::PostgresBackendError> { 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_result = sql_transaction.query_opt(LOCK_ACCOUNT_STATE_SQL, &[&pubkey_bytes, &slot_text.as_str(), &state_hash_bytes]).await; let row = match row_result { std::result::Result::Ok(std::option::Option::Some(value)) => value, std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Ok(std::option::Option::None), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_acquisition_lock_state")), }; let physical = match raw_account_state_db_row(&row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return match decode_raw_account_state_row(network, physical) { std::result::Result::Ok(value) => std::result::Result::Ok(std::option::Option::Some(value)), std::result::Result::Err(error) => std::result::Result::Err(error), }; } async fn persist_account_observation_row( sql_transaction: &deadpool_postgres::Transaction<'_>, network: &ksp_store_api::RawNetworkId, observation: &ksp_store_api::RawAccountObservation, ) -> std::result::Result { let provenance = observation.provenance(); let observation_key = observation.observation_key(); let observation_key_bytes: &[u8] = observation_key.as_bytes(); let account = observation.account(); let account_pubkey_bytes: &[u8] = account.pubkey().as_ref(); let account_slot_text = account.slot().to_string(); let account_state_hash = account.state_hash(); let account_state_hash_bytes: &[u8] = account_state_hash.as_bytes(); let origin = match encode_origin(provenance.origin()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let received_at = match i64::try_from(provenance.received_at().unix_millis()) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_received_at_encode")), }; let observed_at = match provenance.observed_at() { std::option::Option::Some(value) => match i64::try_from(value.unix_millis()) { std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_observed_at_encode")), }, std::option::Option::None => std::option::Option::None, }; let source_payload_size = match provenance.source_payload_size_bytes() { std::option::Option::Some(value) => match i64::try_from(value) { std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_source_size_encode")), }, std::option::Option::None => std::option::Option::None, }; let capture_session_id = provenance.capture_session_id().map(|value| return value.as_str()); let commitment = provenance.commitment().map(|value| return value.as_str()); let endpoint_id = provenance.endpoint_id().map(|value| return value.as_str()); let filter_id = provenance.filter_id().map(|value| return value.as_str()); let source_payload_hash = provenance.source_payload_hash(); let source_payload_hash_bytes: std::option::Option<&[u8]> = source_payload_hash.as_ref().map(|value| return &value.as_bytes()[..]); let is_startup = observation.is_startup(); let transaction_signature = observation.transaction_signature(); let transaction_signature_bytes: std::option::Option<&[u8]> = transaction_signature.as_ref().map(|value| return &value.as_bytes()[..]); let write_version_text = observation.write_version().map(|value| return value.to_string()); let row_result = sql_transaction .query_opt( INSERT_ACCOUNT_OBSERVATION_SQL, &[ &observation_key_bytes, &account_pubkey_bytes, &account_slot_text.as_str(), &account_state_hash_bytes, &provenance.provider().as_str(), &provenance.protocol().as_str(), &provenance.acquisition_method().as_str(), &origin, &received_at, &capture_session_id, &commitment, &endpoint_id, &filter_id, &observed_at, &source_payload_hash_bytes, &source_payload_size, &is_startup, &transaction_signature_bytes, &write_version_text, ], ) .await; let inserted = match row_result { std::result::Result::Ok(std::option::Option::Some(_)) => true, std::result::Result::Ok(std::option::Option::None) => false, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_observation_insert")), }; if inserted { return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::Inserted); } let existing_result = sql_transaction.query_opt(LOCK_ACCOUNT_OBSERVATION_SQL, &[&observation_key_bytes]).await; let existing_row = match existing_result { std::result::Result::Ok(std::option::Option::Some(value)) => value, std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(data_invalid("raw_account_observation_conflict_missing")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_account_observation_conflict_query")), }; let physical = match raw_account_observation_db_row(&existing_row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let stored = match decode_raw_account_observation_row(network, physical) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if stored.eq(observation) { return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::AlreadyPresent); } return std::result::Result::Err(conflict("raw_account_observation_content_conflict")); } fn ensure_acquisition_inputs( network: &ksp_store_api::RawNetworkId, state: &ksp_store_api::RawAccountState, observation: &ksp_store_api::RawAccountObservation, ) -> std::result::Result<(), crate::PostgresBackendError> { let state_network_result = ensure_network(network, state.reference(), "raw_account_acquisition_state_network"); if let std::result::Result::Err(error) = state_network_result { return std::result::Result::Err(error); } let observation_network_result = ensure_network(network, observation.account(), "raw_account_acquisition_observation_network"); if let std::result::Result::Err(error) = observation_network_result { return std::result::Result::Err(error); } if observation.account() != state.reference() { return std::result::Result::Err(conflict("raw_account_acquisition_reference_mismatch")); } return std::result::Result::Ok(()); } fn ensure_observation_write_input( network: &ksp_store_api::RawNetworkId, observation: &ksp_store_api::RawAccountObservation, ) -> std::result::Result<(), crate::PostgresBackendError> { return ensure_network(network, observation.account(), "raw_account_observation_write_network"); } fn raw_account_states_equal(left: &ksp_store_api::RawAccountState, right: &ksp_store_api::RawAccountState) -> bool { return left.reference() == right.reference() && left.lamports() == right.lamports() && left.owner() == right.owner() && left.executable() == right.executable() && left.rent_epoch() == right.rent_epoch() && left.data() == right.data(); } fn encode_origin(origin: ksp_store_api::RawAcquisitionOrigin) -> std::result::Result<&'static str, crate::PostgresBackendError> { return match origin { ksp_store_api::RawAcquisitionOrigin::Backfill => std::result::Result::Ok("backfill"), ksp_store_api::RawAcquisitionOrigin::Import => std::result::Result::Ok("import"), ksp_store_api::RawAcquisitionOrigin::Live => std::result::Result::Ok("live"), ksp_store_api::RawAcquisitionOrigin::Repair => std::result::Result::Ok("repair"), ksp_store_api::RawAcquisitionOrigin::Replay => std::result::Result::Ok("replay"), _ => std::result::Result::Err(data_invalid("raw_account_observation_origin_encode")), }; } fn conflict(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::Conflict, phase); } fn decode_origin(value: &str) -> std::result::Result { return match value { "backfill" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Backfill), "import" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Import), "live" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Live), "repair" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Repair), "replay" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Replay), _ => std::result::Result::Err(data_invalid("raw_account_observation_origin")), }; } fn decode_provenance_code(value: std::string::String) -> std::result::Result { return match ksp_store_api::RawProvenanceCode::new(value) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_account_provenance_code")), }; } fn decode_timestamp_i64(value: i64, phase: &'static str) -> std::result::Result { let unsigned = match u64::try_from(value) { std::result::Result::Ok(decoded) => decoded, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid(phase)), }; return match ksp_store_api::RawTimestamp::from_unix_millis(unsigned) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn decode_u64_decimal(value: &str, phase: &'static str) -> std::result::Result { return match value.parse::() { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn decode_u64_i64(value: i64, phase: &'static str) -> std::result::Result { return match u64::try_from(value) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn fixed_bytes(value: std::vec::Vec, phase: &'static str) -> std::result::Result<[u8; N], crate::PostgresBackendError> { return match <[u8; N]>::try_from(value.as_slice()) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn ensure_network( network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawAccountStateReference, phase: &'static str, ) -> std::result::Result<(), crate::PostgresBackendError> { if reference.network() != network { return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, phase)); } return std::result::Result::Ok(()); } fn data_invalid(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::DataInvalid, phase); } fn reference_not_found(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReferenceNotFound, phase); } fn write_failed(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WriteFailed, phase); } #[cfg(test)] #[path = "../unit_tests/raw_account.rs"] mod tests;