v0.3.8-pre.005
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/lib.rs
|
||||
// version: 22
|
||||
// version: 23
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -33,8 +33,9 @@
|
||||
//! deterministic account keyset pagination with the fixed `KSPA` cursor. `0.3.4-pre.008`
|
||||
//! implements the four `RawAccount*` capabilities directly on `PostgresBackend`, completing
|
||||
//! the backend RAW capability inventory at ten without exposing physical PostgreSQL types.
|
||||
//! `0.3.8-pre.004` adds the payload-free random-access RawTransaction inspection
|
||||
//! capability with exact counts while preserving keyset traversal unchanged.
|
||||
//! `0.3.8-pre.004` adds payload-free random-access RawTransaction inspection;
|
||||
//! `0.3.8-pre.005` adds the corresponding data-free RawAccountState inspection
|
||||
//! with exact counts while preserving both keyset traversal families unchanged.
|
||||
//!
|
||||
//! This crate depends on `ksp-store-api` and never on `ksp-store-lib`. The
|
||||
//! common facade consumes only this crate's narrow backend bridge and never
|
||||
@@ -86,6 +87,8 @@ pub(crate) use self::raw_account::get_raw_account_observation;
|
||||
pub(crate) use self::raw_account::get_raw_account_state;
|
||||
/// Private RAW account-state list reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_account::list_raw_account_states;
|
||||
/// Private data-free RAW account-state inspection reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_account::inspect_raw_account_states;
|
||||
/// Private atomic RAW account acquisition writer consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_account::persist_raw_account_acquisition;
|
||||
/// Private additional RAW account observation writer consumed by the physical backend runtime.
|
||||
@@ -104,10 +107,10 @@ pub(crate) use self::raw_transaction::get_raw_transaction_observation;
|
||||
pub(crate) use self::raw_transaction::get_raw_transaction_retention_state;
|
||||
/// Private RAW transaction tombstone reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_transaction::get_raw_transaction_tombstone;
|
||||
/// Private payload-free RAW transaction inspection reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_transaction::inspect_raw_transactions;
|
||||
/// Private RAW transaction list reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_transaction::list_raw_transactions;
|
||||
/// Private payload-free RAW transaction inspection reader consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_transaction::inspect_raw_transactions;
|
||||
/// Private atomic RAW transaction acquisition writer consumed by the physical backend runtime.
|
||||
pub(crate) use self::raw_transaction::persist_raw_transaction_acquisition;
|
||||
/// Private additional RAW transaction observation writer consumed by the physical backend runtime.
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/raw_account.rs
|
||||
// version: 5
|
||||
// version: 6
|
||||
|
||||
pub(crate) mod cursor;
|
||||
|
||||
@@ -7,6 +7,8 @@ const GET_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key, account_pubke
|
||||
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_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";
|
||||
@@ -22,6 +24,20 @@ struct RawAccountListDbRow {
|
||||
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 RawAccountObservationDbRow {
|
||||
account_pubkey: std::vec::Vec<u8>,
|
||||
account_slot_text: std::string::String,
|
||||
@@ -203,6 +219,107 @@ pub(crate) async fn list_raw_account_states(
|
||||
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")),
|
||||
};
|
||||
}
|
||||
|
||||
/// Reads one persisted RAW account observation by producer-owned idempotence key.
|
||||
pub(crate) async fn get_raw_account_observation(
|
||||
pool: &deadpool_postgres::Pool,
|
||||
@@ -913,6 +1030,182 @@ fn decode_timestamp_i64(value: i64, phase: &'static str) -> std::result::Result<
|
||||
};
|
||||
}
|
||||
|
||||
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),
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/runtime.rs
|
||||
// version: 16
|
||||
// version: 17
|
||||
|
||||
const APPLICATION_NAME: &str = "ksp-store";
|
||||
const MAX_CONNECTION_URI_BYTES: usize = 4_096;
|
||||
@@ -348,6 +348,14 @@ impl PostgresBackend {
|
||||
return crate::list_raw_account_states(&self.pool, &self.network, query).await;
|
||||
}
|
||||
|
||||
/// Inspects one data-free random-access RAW account-state window with exact logical counts.
|
||||
pub async fn inspect_raw_account_states(
|
||||
&self,
|
||||
query: &ksp_store_api::RawAccountStateInspectionQuery,
|
||||
) -> std::result::Result<ksp_store_api::RawInspectionPage<ksp_store_api::RawAccountStateSummary>, crate::PostgresBackendError> {
|
||||
return crate::inspect_raw_account_states(&self.pool, &self.network, query).await;
|
||||
}
|
||||
|
||||
/// Reads one persisted RAW account observation by producer-owned idempotence key.
|
||||
pub async fn get_raw_account_observation(
|
||||
&self,
|
||||
@@ -523,6 +531,18 @@ impl ksp_store_api::RawAccountStateRead for PostgresBackend {
|
||||
}
|
||||
}
|
||||
|
||||
impl ksp_store_api::RawAccountStateInspectionRead for PostgresBackend {
|
||||
fn inspect_raw_account_states<'a>(
|
||||
&'a self,
|
||||
query: &'a ksp_store_api::RawAccountStateInspectionQuery,
|
||||
) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result<ksp_store_api::RawInspectionPage<ksp_store_api::RawAccountStateSummary>>> {
|
||||
return std::boxed::Box::pin(async move {
|
||||
let result = PostgresBackend::inspect_raw_account_states(self, query).await;
|
||||
return result.map_err(map_capability_error);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
impl ksp_store_api::RawAccountStateWrite for PostgresBackend {
|
||||
fn persist_raw_account_acquisition<'a>(
|
||||
&'a self,
|
||||
|
||||
Reference in New Issue
Block a user