Files
khadhroony-solana-project/crates/ksp-store-postgres-lib/src/raw_account.rs

1525 lines
90 KiB
Rust

// file: crates/ksp-store-postgres-lib/src/raw_account.rs
// version: 8
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 INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3))), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_observations) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.observation_key IS NOT NULL AS page_present, page.observation_key, page.account_pubkey, page.account_slot_text, page.account_state_hash, page.provider, page.protocol, page.acquisition_method, page.origin, page.received_at_unix_millis, page.capture_session_id, page.commitment, page.endpoint_id, page.filter_id, page.observed_at_unix_millis, page.source_payload_hash, page.source_payload_size_bytes, page.is_startup, page.transaction_signature, page.write_version_text FROM counts LEFT JOIN LATERAL (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 ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3)) ORDER BY received_at_unix_millis ASC, observation_key ASC LIMIT $4 OFFSET $5) AS page ON TRUE";
const INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3))), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_observations) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.observation_key IS NOT NULL AS page_present, page.observation_key, page.account_pubkey, page.account_slot_text, page.account_state_hash, page.provider, page.protocol, page.acquisition_method, page.origin, page.received_at_unix_millis, page.capture_session_id, page.commitment, page.endpoint_id, page.filter_id, page.observed_at_unix_millis, page.source_payload_hash, page.source_payload_size_bytes, page.is_startup, page.transaction_signature, page.write_version_text FROM counts LEFT JOIN LATERAL (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 ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3)) ORDER BY received_at_unix_millis DESC, observation_key DESC LIMIT $4 OFFSET $5) AS page ON TRUE";
const INSPECT_ACCOUNT_STATES_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_states WHERE ($1::BYTEA IS NULL OR pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL AND $2::TEXT IS NULL AND $3::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_states) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.pubkey IS NOT NULL AS page_present, page.pubkey, page.slot_text, page.state_hash, page.lamports_text, page.owner, page.executable, page.rent_epoch_text, page.data_length_bytes FROM counts LEFT JOIN LATERAL (SELECT account_row.pubkey, account_row.slot::TEXT AS slot_text, account_row.state_hash, account_row.lamports::TEXT AS lamports_text, account_row.owner, account_row.executable, account_row.rent_epoch::TEXT AS rent_epoch_text, OCTET_LENGTH(account_row.data)::BIGINT AS data_length_bytes FROM ksp_raw_account_states AS account_row WHERE ($1::BYTEA IS NULL OR account_row.pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR account_row.slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR account_row.slot <= $3::TEXT::NUMERIC) ORDER BY account_row.slot ASC, account_row.pubkey ASC, account_row.state_hash ASC LIMIT $4 OFFSET $5) AS page ON TRUE";
const INSPECT_ACCOUNT_STATES_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_states WHERE ($1::BYTEA IS NULL OR pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL AND $2::TEXT IS NULL AND $3::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_states) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.pubkey IS NOT NULL AS page_present, page.pubkey, page.slot_text, page.state_hash, page.lamports_text, page.owner, page.executable, page.rent_epoch_text, page.data_length_bytes FROM counts LEFT JOIN LATERAL (SELECT account_row.pubkey, account_row.slot::TEXT AS slot_text, account_row.state_hash, account_row.lamports::TEXT AS lamports_text, account_row.owner, account_row.executable, account_row.rent_epoch::TEXT AS rent_epoch_text, OCTET_LENGTH(account_row.data)::BIGINT AS data_length_bytes FROM ksp_raw_account_states AS account_row WHERE ($1::BYTEA IS NULL OR account_row.pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR account_row.slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR account_row.slot <= $3::TEXT::NUMERIC) ORDER BY account_row.slot DESC, account_row.pubkey DESC, account_row.state_hash DESC LIMIT $4 OFFSET $5) AS page ON TRUE";
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<u8>,
slot_text: std::string::String,
state_hash: std::vec::Vec<u8>,
}
struct RawAccountInspectionDbRow {
data_length_bytes: std::option::Option<i64>,
executable: std::option::Option<bool>,
filtered_count_text: std::string::String,
lamports_text: std::option::Option<std::string::String>,
owner: std::option::Option<std::vec::Vec<u8>>,
page_present: bool,
pubkey: std::option::Option<std::vec::Vec<u8>>,
rent_epoch_text: std::option::Option<std::string::String>,
slot_text: std::option::Option<std::string::String>,
state_hash: std::option::Option<std::vec::Vec<u8>>,
total_count_text: std::string::String,
}
struct RawAccountObservationInspectionDbRow {
account_pubkey: std::option::Option<std::vec::Vec<u8>>,
account_slot_text: std::option::Option<std::string::String>,
account_state_hash: std::option::Option<std::vec::Vec<u8>>,
acquisition_method: std::option::Option<std::string::String>,
capture_session_id: std::option::Option<std::string::String>,
commitment: std::option::Option<std::string::String>,
endpoint_id: std::option::Option<std::string::String>,
filter_id: std::option::Option<std::string::String>,
filtered_count_text: std::string::String,
is_startup: std::option::Option<bool>,
observation_key: std::option::Option<std::vec::Vec<u8>>,
observed_at_unix_millis: std::option::Option<i64>,
origin: std::option::Option<std::string::String>,
page_present: bool,
protocol: std::option::Option<std::string::String>,
provider: std::option::Option<std::string::String>,
received_at_unix_millis: std::option::Option<i64>,
source_payload_hash: std::option::Option<std::vec::Vec<u8>>,
source_payload_size_bytes: std::option::Option<i64>,
total_count_text: std::string::String,
transaction_signature: std::option::Option<std::vec::Vec<u8>>,
write_version_text: std::option::Option<std::string::String>,
}
struct RawAccountObservationDbRow {
account_pubkey: std::vec::Vec<u8>,
account_slot_text: std::string::String,
account_state_hash: std::vec::Vec<u8>,
acquisition_method: std::string::String,
capture_session_id: std::option::Option<std::string::String>,
commitment: std::option::Option<std::string::String>,
endpoint_id: std::option::Option<std::string::String>,
filter_id: std::option::Option<std::string::String>,
is_startup: std::option::Option<bool>,
observation_key: std::vec::Vec<u8>,
observed_at_unix_millis: std::option::Option<i64>,
origin: std::string::String,
protocol: std::string::String,
provider: std::string::String,
received_at_unix_millis: i64,
source_payload_hash: std::option::Option<std::vec::Vec<u8>>,
source_payload_size_bytes: std::option::Option<i64>,
transaction_signature: std::option::Option<std::vec::Vec<u8>>,
write_version_text: std::option::Option<std::string::String>,
}
struct RawAccountStateDbRow {
data: std::vec::Vec<u8>,
executable: bool,
lamports_text: std::string::String,
owner: std::vec::Vec<u8>,
pubkey: std::vec::Vec<u8>,
rent_epoch_text: std::string::String,
slot_text: std::string::String,
state_hash: std::vec::Vec<u8>,
}
/// 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<std::option::Option<ksp_store_api::RawAccountState>, 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<ksp_store_api::RawPage<ksp_store_api::RawAccountStateReference>, 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));
}
/// Inspects one data-free random-access RAW account-state window with exact logical counts.
pub(crate) async fn inspect_raw_account_states(
pool: &deadpool_postgres::Pool,
network: &ksp_store_api::RawNetworkId,
query: &ksp_store_api::RawAccountStateInspectionQuery,
) -> std::result::Result<ksp_store_api::RawInspectionPage<ksp_store_api::RawAccountStateSummary>, crate::PostgresBackendError> {
if query.network() != network {
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, "raw_account_inspection_network"));
}
let sql = match query.direction() {
ksp_store_api::RawSortDirection::Ascending => INSPECT_ACCOUNT_STATES_ASC_SQL,
ksp_store_api::RawSortDirection::Descending => INSPECT_ACCOUNT_STATES_DESC_SQL,
_ => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::QueryInvalid,
"raw_account_inspection_direction",
));
},
};
let (sql_limit, sql_offset) = match raw_account_inspection_sql_window(query.page()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => 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 = query.pubkey().map(|value| return value.to_bytes().to_vec());
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 rows_result = client.query(sql, &[&pubkey_bytes, &start_text, &end_text, &sql_limit, &sql_offset]).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_inspection_query"));
},
};
if rows.is_empty() {
return std::result::Result::Err(data_invalid("raw_account_inspection_cardinality"));
}
let row_count = rows.len();
let mut total_items = std::option::Option::None;
let mut filtered_items = std::option::Option::None;
let mut items = std::vec::Vec::new();
for row in rows {
let physical = match raw_account_inspection_db_row(&row) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let row_total = match decode_u64_decimal(physical.total_count_text.as_str(), "raw_account_inspection_total") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let row_filtered = match decode_u64_decimal(physical.filtered_count_text.as_str(), "raw_account_inspection_filtered") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match (total_items, filtered_items) {
(std::option::Option::None, std::option::Option::None) => {
total_items = std::option::Option::Some(row_total);
filtered_items = std::option::Option::Some(row_filtered);
},
(std::option::Option::Some(total), std::option::Option::Some(filtered)) if total == row_total && filtered == row_filtered => {},
_ => return std::result::Result::Err(data_invalid("raw_account_inspection_counts")),
}
if physical.page_present {
let summary = match decode_raw_account_inspection_summary(network, physical) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
items.push(summary);
} else if row_count != 1 || !raw_account_inspection_empty_page_is_clean(&physical) {
return std::result::Result::Err(data_invalid("raw_account_inspection_page"));
}
}
let total_items = match total_items {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_inspection_total")),
};
let filtered_items = match filtered_items {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_inspection_filtered")),
};
let item_count = match u64::try_from(items.len()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_page")),
};
if item_count > query.page().limit().get() {
return std::result::Result::Err(data_invalid("raw_account_inspection_page"));
}
return match ksp_store_api::RawInspectionPage::try_new(items, total_items, filtered_items) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_account_inspection_page")),
};
}
/// Inspects one safe random-access RAW account-observation window with exact counts.
pub(crate) async fn inspect_raw_account_observations(
pool: &deadpool_postgres::Pool,
network: &ksp_store_api::RawNetworkId,
query: &ksp_store_api::RawAccountObservationInspectionQuery,
) -> std::result::Result<ksp_store_api::RawInspectionPage<ksp_store_api::RawAccountObservationSummary>, crate::PostgresBackendError> {
if query.network() != network {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::WrongNetwork,
"raw_account_observation_inspection_network",
));
}
let sql = match query.direction() {
ksp_store_api::RawSortDirection::Ascending => INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL,
ksp_store_api::RawSortDirection::Descending => INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL,
_ => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::QueryInvalid,
"raw_account_observation_inspection_direction",
));
},
};
let (sql_limit, sql_offset) = match raw_account_inspection_sql_window(query.page()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let account_pubkey = query.account().map(|reference| return reference.pubkey().to_bytes().to_vec());
let account_slot_text = query.account().map(|reference| return reference.slot().to_string());
let account_state_hash = query.account().map(|reference| return reference.state_hash().as_bytes().to_vec());
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 rows_result = client.query(sql, &[&account_pubkey, &account_slot_text, &account_state_hash, &sql_limit, &sql_offset]).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_inspection_query",
));
},
};
if rows.is_empty() {
return std::result::Result::Err(data_invalid("raw_account_observation_inspection_cardinality"));
}
let row_count = rows.len();
let mut total_items = std::option::Option::None;
let mut filtered_items = std::option::Option::None;
let mut items = std::vec::Vec::new();
for row in rows {
let physical = match raw_account_observation_inspection_db_row(&row) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let row_total = match decode_u64_decimal(physical.total_count_text.as_str(), "raw_account_observation_inspection_total") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let row_filtered = match decode_u64_decimal(physical.filtered_count_text.as_str(), "raw_account_observation_inspection_filtered") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
match (total_items, filtered_items) {
(std::option::Option::None, std::option::Option::None) => {
total_items = std::option::Option::Some(row_total);
filtered_items = std::option::Option::Some(row_filtered);
},
(std::option::Option::Some(total), std::option::Option::Some(filtered)) if total == row_total && filtered == row_filtered => {},
_ => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_counts")),
}
if physical.page_present {
let observation = match decode_raw_account_observation_inspection_row(network, physical) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if let std::option::Option::Some(reference) = query.account()
&& observation.account() != reference
{
return std::result::Result::Err(data_invalid("raw_account_observation_inspection_reference"));
}
items.push(ksp_store_api::RawAccountObservationSummary::new(
observation.observation_key(),
observation.account().clone(),
observation.provenance().clone(),
observation.is_startup(),
observation.transaction_signature(),
observation.write_version(),
));
} else if row_count != 1 || !raw_account_observation_inspection_empty_page_is_clean(&physical) {
return std::result::Result::Err(data_invalid("raw_account_observation_inspection_page"));
}
}
let total_items = match total_items {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_total")),
};
let filtered_items = match filtered_items {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_filtered")),
};
let item_count = match u64::try_from(items.len()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_page")),
};
if item_count > query.page().limit().get() {
return std::result::Result::Err(data_invalid("raw_account_observation_inspection_page"));
}
return match ksp_store_api::RawInspectionPage::try_new(items, total_items, filtered_items) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_account_observation_inspection_page")),
};
}
/// 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<std::option::Option<ksp_store_api::RawAccountObservation>, 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<ksp_store_api::RawAcquisitionWriteOutcome, crate::PostgresBackendError> {
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<ksp_store_api::RawObservationWriteOutcome, crate::PostgresBackendError> {
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<RawAccountListDbRow, crate::PostgresBackendError> {
let pubkey = match row.try_get::<_, std::vec::Vec<u8>>("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<u8>>("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_inspection_empty_page_is_clean(row: &RawAccountObservationInspectionDbRow) -> bool {
return row.account_pubkey.is_none()
&& row.account_slot_text.is_none()
&& row.account_state_hash.is_none()
&& row.acquisition_method.is_none()
&& row.capture_session_id.is_none()
&& row.commitment.is_none()
&& row.endpoint_id.is_none()
&& row.filter_id.is_none()
&& row.is_startup.is_none()
&& row.observation_key.is_none()
&& row.observed_at_unix_millis.is_none()
&& row.origin.is_none()
&& row.protocol.is_none()
&& row.provider.is_none()
&& row.received_at_unix_millis.is_none()
&& row.source_payload_hash.is_none()
&& row.source_payload_size_bytes.is_none()
&& row.transaction_signature.is_none()
&& row.write_version_text.is_none();
}
fn raw_account_observation_inspection_db_row(
row: &tokio_postgres::Row,
) -> std::result::Result<RawAccountObservationInspectionDbRow, crate::PostgresBackendError> {
macro_rules! get_optional {
($name:literal, $ty:ty) => {
match row.try_get::<_, std::option::Option<$ty>>($name) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_decode")),
}
};
}
let filtered_count_text = match row.try_get::<_, std::string::String>("filtered_count_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_decode")),
};
let page_present = match row.try_get::<_, bool>("page_present") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_decode")),
};
let total_count_text = match row.try_get::<_, std::string::String>("total_count_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_inspection_decode")),
};
return std::result::Result::Ok(RawAccountObservationInspectionDbRow {
account_pubkey: get_optional!("account_pubkey", std::vec::Vec<u8>),
account_slot_text: get_optional!("account_slot_text", std::string::String),
account_state_hash: get_optional!("account_state_hash", std::vec::Vec<u8>),
acquisition_method: get_optional!("acquisition_method", std::string::String),
capture_session_id: get_optional!("capture_session_id", std::string::String),
commitment: get_optional!("commitment", std::string::String),
endpoint_id: get_optional!("endpoint_id", std::string::String),
filter_id: get_optional!("filter_id", std::string::String),
filtered_count_text,
is_startup: get_optional!("is_startup", bool),
observation_key: get_optional!("observation_key", std::vec::Vec<u8>),
observed_at_unix_millis: get_optional!("observed_at_unix_millis", i64),
origin: get_optional!("origin", std::string::String),
page_present,
protocol: get_optional!("protocol", std::string::String),
provider: get_optional!("provider", std::string::String),
received_at_unix_millis: get_optional!("received_at_unix_millis", i64),
source_payload_hash: get_optional!("source_payload_hash", std::vec::Vec<u8>),
source_payload_size_bytes: get_optional!("source_payload_size_bytes", i64),
total_count_text,
transaction_signature: get_optional!("transaction_signature", std::vec::Vec<u8>),
write_version_text: get_optional!("write_version_text", std::string::String),
});
}
fn decode_raw_account_observation_inspection_row(
network: &ksp_store_api::RawNetworkId,
row: RawAccountObservationInspectionDbRow,
) -> std::result::Result<ksp_store_api::RawAccountObservation, crate::PostgresBackendError> {
let account_pubkey = match inspection_required(row.account_pubkey, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let account_slot_text = match inspection_required(row.account_slot_text, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let account_state_hash = match inspection_required(row.account_state_hash, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let acquisition_method = match inspection_required(row.acquisition_method, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observation_key = match inspection_required(row.observation_key, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let origin = match inspection_required(row.origin, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let protocol = match inspection_required(row.protocol, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let provider = match inspection_required(row.provider, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let received_at_unix_millis = match inspection_required(row.received_at_unix_millis, "raw_account_observation_inspection_shape") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let physical = RawAccountObservationDbRow {
account_pubkey,
account_slot_text,
account_state_hash,
acquisition_method,
capture_session_id: row.capture_session_id,
commitment: row.commitment,
endpoint_id: row.endpoint_id,
filter_id: row.filter_id,
is_startup: row.is_startup,
observation_key,
observed_at_unix_millis: row.observed_at_unix_millis,
origin,
protocol,
provider,
received_at_unix_millis,
source_payload_hash: row.source_payload_hash,
source_payload_size_bytes: row.source_payload_size_bytes,
transaction_signature: row.transaction_signature,
write_version_text: row.write_version_text,
};
return decode_raw_account_observation_row(network, physical);
}
fn raw_account_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawAccountObservationDbRow, crate::PostgresBackendError> {
let account_pubkey: std::vec::Vec<u8> = 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<u8> = 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<std::string::String> = 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<std::string::String> = 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<std::string::String> = 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<std::string::String> = 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<bool> = 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<u8> = 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<i64> = 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<std::vec::Vec<u8>> = 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<i64> = 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<std::vec::Vec<u8>> = 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<std::string::String> = 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<RawAccountStateDbRow, crate::PostgresBackendError> {
let data: std::vec::Vec<u8> = 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<u8> = 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<u8> = 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<u8> = 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<ksp_store_api::RawAccountObservation, crate::PostgresBackendError> {
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<ksp_store_api::RawAccountState, crate::PostgresBackendError> {
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<bool, crate::PostgresBackendError> {
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<std::option::Option<ksp_store_api::RawAccountState>, 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<ksp_store_api::RawObservationWriteOutcome, crate::PostgresBackendError> {
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<ksp_store_api::RawAcquisitionOrigin, crate::PostgresBackendError> {
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<ksp_store_api::RawProvenanceCode, crate::PostgresBackendError> {
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<ksp_store_api::RawTimestamp, crate::PostgresBackendError> {
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 raw_account_inspection_empty_page_is_clean(row: &RawAccountInspectionDbRow) -> bool {
return !row.page_present
&& row.data_length_bytes.is_none()
&& row.executable.is_none()
&& row.lamports_text.is_none()
&& row.owner.is_none()
&& row.pubkey.is_none()
&& row.rent_epoch_text.is_none()
&& row.slot_text.is_none()
&& row.state_hash.is_none();
}
fn raw_account_inspection_sql_window(page: ksp_store_api::RawInspectionPageRequest) -> std::result::Result<(i64, i64), crate::PostgresBackendError> {
let sql_limit = match i64::try_from(page.limit().get()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(
crate::PostgresBackendErrorKind::PageLimitUnsupported,
"raw_account_inspection_limit",
));
},
};
let sql_offset = match i64::try_from(page.offset()) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => {
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::QueryInvalid, "raw_account_inspection_offset"));
},
};
return std::result::Result::Ok((sql_limit, sql_offset));
}
fn raw_account_inspection_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawAccountInspectionDbRow, crate::PostgresBackendError> {
let data_length_bytes = match row.try_get::<_, std::option::Option<i64>>("data_length_bytes") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let executable = match row.try_get::<_, std::option::Option<bool>>("executable") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let filtered_count_text = match row.try_get::<_, std::string::String>("filtered_count_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let lamports_text = match row.try_get::<_, std::option::Option<std::string::String>>("lamports_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let owner = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("owner") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let page_present = match row.try_get::<_, bool>("page_present") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let pubkey = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("pubkey") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let rent_epoch_text = match row.try_get::<_, std::option::Option<std::string::String>>("rent_epoch_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let slot_text = match row.try_get::<_, std::option::Option<std::string::String>>("slot_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let state_hash = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("state_hash") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
let total_count_text = match row.try_get::<_, std::string::String>("total_count_text") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_decode")),
};
return std::result::Result::Ok(RawAccountInspectionDbRow {
data_length_bytes,
executable,
filtered_count_text,
lamports_text,
owner,
page_present,
pubkey,
rent_epoch_text,
slot_text,
state_hash,
total_count_text,
});
}
fn decode_raw_account_inspection_summary(
network: &ksp_store_api::RawNetworkId,
row: RawAccountInspectionDbRow,
) -> std::result::Result<ksp_store_api::RawAccountStateSummary, crate::PostgresBackendError> {
if !row.page_present {
return std::result::Result::Err(data_invalid("raw_account_inspection_page"));
}
let pubkey = match inspection_required(row.pubkey, "raw_account_inspection_pubkey") {
std::result::Result::Ok(value) => match fixed_bytes::<32>(value, "raw_account_inspection_pubkey") {
std::result::Result::Ok(bytes) => ksp_store_api::Pubkey::new_from_array(bytes),
std::result::Result::Err(error) => return std::result::Result::Err(error),
},
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let slot_text = match inspection_required(row.slot_text, "raw_account_inspection_slot") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let slot = match decode_u64_decimal(slot_text.as_str(), "raw_account_inspection_slot") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let state_hash_raw = match inspection_required(row.state_hash, "raw_account_inspection_state_hash") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let state_hash = match fixed_bytes::<32>(state_hash_raw, "raw_account_inspection_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_text = match inspection_required(row.lamports_text, "raw_account_inspection_lamports") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let lamports = match decode_u64_decimal(lamports_text.as_str(), "raw_account_inspection_lamports") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let owner_raw = match inspection_required(row.owner, "raw_account_inspection_owner") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let owner = match fixed_bytes::<32>(owner_raw, "raw_account_inspection_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 executable = match inspection_required(row.executable, "raw_account_inspection_executable") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let rent_epoch_text = match inspection_required(row.rent_epoch_text, "raw_account_inspection_rent_epoch") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let rent_epoch = match decode_u64_decimal(rent_epoch_text.as_str(), "raw_account_inspection_rent_epoch") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let data_length_raw = match inspection_required(row.data_length_bytes, "raw_account_inspection_data_length") {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let data_length_bytes = match u64::try_from(data_length_raw) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_inspection_data_length")),
};
let reference = ksp_store_api::RawAccountStateReference::new(network.clone(), pubkey, slot, state_hash);
return match ksp_store_api::RawAccountStateSummary::try_new(reference, lamports, owner, executable, rent_epoch, data_length_bytes) {
std::result::Result::Ok(value) => std::result::Result::Ok(value),
std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_account_inspection_summary")),
};
}
fn inspection_required<T>(value: std::option::Option<T>, phase: &'static str) -> std::result::Result<T, crate::PostgresBackendError> {
return match value {
std::option::Option::Some(inner) => std::result::Result::Ok(inner),
std::option::Option::None => std::result::Result::Err(data_invalid(phase)),
};
}
fn decode_u64_decimal(value: &str, phase: &'static str) -> std::result::Result<u64, crate::PostgresBackendError> {
return match value.parse::<u64>() {
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<u64, crate::PostgresBackendError> {
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<const N: usize>(value: std::vec::Vec<u8>, 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;