|
|
|
|
@@ -1,5 +1,5 @@
|
|
|
|
|
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
|
|
|
|
|
// version: 4
|
|
|
|
|
// version: 5
|
|
|
|
|
|
|
|
|
|
pub(crate) mod cursor;
|
|
|
|
|
|
|
|
|
|
@@ -12,6 +12,8 @@ const GET_TRANSACTION_SQL: &str = "SELECT transaction_row.signature, transaction
|
|
|
|
|
const INSERT_ARCHIVE_PAYLOAD_SQL: &str = "INSERT INTO ksp_raw_transaction_archive_payloads (signature, payload) VALUES ($1, $2)";
|
|
|
|
|
const INSERT_OBSERVATION_SQL: &str = "INSERT INTO ksp_raw_transaction_observations (observation_key, transaction_signature, 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) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14) ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key";
|
|
|
|
|
const INSERT_TRANSACTION_SQL: &str = "INSERT INTO ksp_raw_transactions (signature, slot, block_time_unix_millis, format_id, format_version, content_hash, payload, retention_state) VALUES ($1, $2::TEXT::NUMERIC, $3, $4, $5, $6, $7, 'full') ON CONFLICT (signature) DO NOTHING RETURNING signature";
|
|
|
|
|
const INSPECT_TRANSACTIONS_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_transactions WHERE ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::TEXT IS NULL AND $2::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_transactions) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.signature IS NOT NULL AS page_present, page.signature, page.slot_text, page.block_time_unix_millis, page.format_id, page.format_version, page.content_hash, page.retention_state, page.payload_size_bytes, page.hot_payload_present, page.archive_payload_present FROM counts LEFT JOIN LATERAL (SELECT transaction_row.signature, transaction_row.slot::TEXT AS slot_text, transaction_row.block_time_unix_millis, transaction_row.format_id, transaction_row.format_version, transaction_row.content_hash, transaction_row.retention_state, CASE WHEN transaction_row.retention_state = 'full' THEN OCTET_LENGTH(transaction_row.payload)::BIGINT WHEN transaction_row.retention_state = 'archived' THEN OCTET_LENGTH(archive_row.payload)::BIGINT ELSE NULL END AS payload_size_bytes, transaction_row.payload IS NOT NULL AS hot_payload_present, archive_row.payload IS NOT NULL AS archive_payload_present FROM ksp_raw_transactions AS transaction_row LEFT JOIN ksp_raw_transaction_archive_payloads AS archive_row ON archive_row.signature = transaction_row.signature WHERE ($1::TEXT IS NULL OR transaction_row.slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR transaction_row.slot <= $2::TEXT::NUMERIC) ORDER BY transaction_row.slot ASC, transaction_row.signature ASC LIMIT $3 OFFSET $4) AS page ON TRUE";
|
|
|
|
|
const INSPECT_TRANSACTIONS_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_transactions WHERE ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::TEXT IS NULL AND $2::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_transactions) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.signature IS NOT NULL AS page_present, page.signature, page.slot_text, page.block_time_unix_millis, page.format_id, page.format_version, page.content_hash, page.retention_state, page.payload_size_bytes, page.hot_payload_present, page.archive_payload_present FROM counts LEFT JOIN LATERAL (SELECT transaction_row.signature, transaction_row.slot::TEXT AS slot_text, transaction_row.block_time_unix_millis, transaction_row.format_id, transaction_row.format_version, transaction_row.content_hash, transaction_row.retention_state, CASE WHEN transaction_row.retention_state = 'full' THEN OCTET_LENGTH(transaction_row.payload)::BIGINT WHEN transaction_row.retention_state = 'archived' THEN OCTET_LENGTH(archive_row.payload)::BIGINT ELSE NULL END AS payload_size_bytes, transaction_row.payload IS NOT NULL AS hot_payload_present, archive_row.payload IS NOT NULL AS archive_payload_present FROM ksp_raw_transactions AS transaction_row LEFT JOIN ksp_raw_transaction_archive_payloads AS archive_row ON archive_row.signature = transaction_row.signature WHERE ($1::TEXT IS NULL OR transaction_row.slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR transaction_row.slot <= $2::TEXT::NUMERIC) ORDER BY transaction_row.slot DESC, transaction_row.signature DESC LIMIT $3 OFFSET $4) AS page ON TRUE";
|
|
|
|
|
const LIST_TRANSACTIONS_ASC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($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, signature) > ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot ASC, signature ASC LIMIT $5";
|
|
|
|
|
const LIST_TRANSACTIONS_DESC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($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, signature) < ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot DESC, signature DESC LIMIT $5";
|
|
|
|
|
const LOCK_OBSERVATION_SQL: &str = "SELECT observation_key, transaction_signature, 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 FROM ksp_raw_transaction_observations WHERE observation_key = $1 FOR UPDATE";
|
|
|
|
|
@@ -30,6 +32,22 @@ struct RawListDbRow {
|
|
|
|
|
slot_text: std::string::String,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct RawInspectionDbRow {
|
|
|
|
|
archive_payload_present: std::option::Option<bool>,
|
|
|
|
|
block_time_unix_millis: std::option::Option<i64>,
|
|
|
|
|
content_hash: std::option::Option<std::vec::Vec<u8>>,
|
|
|
|
|
filtered_count_text: std::string::String,
|
|
|
|
|
format_id: std::option::Option<std::string::String>,
|
|
|
|
|
format_version: std::option::Option<i64>,
|
|
|
|
|
hot_payload_present: std::option::Option<bool>,
|
|
|
|
|
page_present: bool,
|
|
|
|
|
payload_size_bytes: std::option::Option<i64>,
|
|
|
|
|
retention_state: std::option::Option<std::string::String>,
|
|
|
|
|
signature: std::option::Option<std::vec::Vec<u8>>,
|
|
|
|
|
slot_text: std::option::Option<std::string::String>,
|
|
|
|
|
total_count_text: std::string::String,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
struct RawObservationDbRow {
|
|
|
|
|
acquisition_method: std::string::String,
|
|
|
|
|
capture_session_id: std::option::Option<std::string::String>,
|
|
|
|
|
@@ -192,6 +210,103 @@ pub(crate) async fn list_raw_transactions(
|
|
|
|
|
return std::result::Result::Ok(ksp_store_api::RawPage::new(items, next_cursor));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Inspects one payload-free random-access RAW transaction window with exact counts.
|
|
|
|
|
pub(crate) async fn inspect_raw_transactions(
|
|
|
|
|
pool: &deadpool_postgres::Pool,
|
|
|
|
|
network: &ksp_store_api::RawNetworkId,
|
|
|
|
|
query: &ksp_store_api::RawTransactionInspectionQuery,
|
|
|
|
|
) -> std::result::Result<ksp_store_api::RawInspectionPage<ksp_store_api::RawTransactionSummary>, crate::PostgresBackendError> {
|
|
|
|
|
if query.network() != network {
|
|
|
|
|
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, "raw_transaction_inspection_network"));
|
|
|
|
|
}
|
|
|
|
|
let sql = match query.direction() {
|
|
|
|
|
ksp_store_api::RawSortDirection::Ascending => INSPECT_TRANSACTIONS_ASC_SQL,
|
|
|
|
|
ksp_store_api::RawSortDirection::Descending => INSPECT_TRANSACTIONS_DESC_SQL,
|
|
|
|
|
_ => {
|
|
|
|
|
return std::result::Result::Err(crate::PostgresBackendError::new(
|
|
|
|
|
crate::PostgresBackendErrorKind::QueryInvalid,
|
|
|
|
|
"raw_transaction_inspection_direction",
|
|
|
|
|
));
|
|
|
|
|
},
|
|
|
|
|
};
|
|
|
|
|
let (sql_limit, sql_offset) = match raw_transaction_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 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, &[&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_transaction_inspection_query"));
|
|
|
|
|
},
|
|
|
|
|
};
|
|
|
|
|
if rows.is_empty() {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_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_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_transaction_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_transaction_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_transaction_inspection_counts")),
|
|
|
|
|
}
|
|
|
|
|
if physical.page_present {
|
|
|
|
|
let summary = match decode_raw_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_inspection_empty_page_is_clean(&physical) {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_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_transaction_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_transaction_inspection_page")),
|
|
|
|
|
};
|
|
|
|
|
if item_count > query.page().limit().get() {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_inspection_page")),
|
|
|
|
|
};
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/// Reads one RAW transaction observation from the physical PostgreSQL backend.
|
|
|
|
|
pub(crate) async fn get_raw_transaction_observation(
|
|
|
|
|
pool: &deadpool_postgres::Pool,
|
|
|
|
|
@@ -1009,6 +1124,221 @@ fn write_failed(phase: &'static str) -> crate::PostgresBackendError {
|
|
|
|
|
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WriteFailed, phase);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn raw_inspection_empty_page_is_clean(row: &RawInspectionDbRow) -> bool {
|
|
|
|
|
return !row.page_present
|
|
|
|
|
&& row.archive_payload_present.is_none()
|
|
|
|
|
&& row.block_time_unix_millis.is_none()
|
|
|
|
|
&& row.content_hash.is_none()
|
|
|
|
|
&& row.format_id.is_none()
|
|
|
|
|
&& row.format_version.is_none()
|
|
|
|
|
&& row.hot_payload_present.is_none()
|
|
|
|
|
&& row.payload_size_bytes.is_none()
|
|
|
|
|
&& row.retention_state.is_none()
|
|
|
|
|
&& row.signature.is_none()
|
|
|
|
|
&& row.slot_text.is_none();
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn raw_transaction_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_transaction_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_transaction_inspection_offset",
|
|
|
|
|
));
|
|
|
|
|
},
|
|
|
|
|
};
|
|
|
|
|
return std::result::Result::Ok((sql_limit, sql_offset));
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn raw_inspection_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawInspectionDbRow, crate::PostgresBackendError> {
|
|
|
|
|
let archive_payload_present = match row.try_get::<_, std::option::Option<bool>>("archive_payload_present") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let block_time_unix_millis = match row.try_get::<_, std::option::Option<i64>>("block_time_unix_millis") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let content_hash = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("content_hash") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let format_id = match row.try_get::<_, std::option::Option<std::string::String>>("format_id") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let format_version = match row.try_get::<_, std::option::Option<i64>>("format_version") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let hot_payload_present = match row.try_get::<_, std::option::Option<bool>>("hot_payload_present") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let payload_size_bytes = match row.try_get::<_, std::option::Option<i64>>("payload_size_bytes") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let retention_state = match row.try_get::<_, std::option::Option<std::string::String>>("retention_state") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
let signature = match row.try_get::<_, std::option::Option<std::vec::Vec<u8>>>("signature") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_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_transaction_inspection_decode")),
|
|
|
|
|
};
|
|
|
|
|
return std::result::Result::Ok(RawInspectionDbRow {
|
|
|
|
|
archive_payload_present,
|
|
|
|
|
block_time_unix_millis,
|
|
|
|
|
content_hash,
|
|
|
|
|
filtered_count_text,
|
|
|
|
|
format_id,
|
|
|
|
|
format_version,
|
|
|
|
|
hot_payload_present,
|
|
|
|
|
page_present,
|
|
|
|
|
payload_size_bytes,
|
|
|
|
|
retention_state,
|
|
|
|
|
signature,
|
|
|
|
|
slot_text,
|
|
|
|
|
total_count_text,
|
|
|
|
|
});
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fn decode_raw_inspection_summary(
|
|
|
|
|
network: &ksp_store_api::RawNetworkId,
|
|
|
|
|
row: RawInspectionDbRow,
|
|
|
|
|
) -> std::result::Result<ksp_store_api::RawTransactionSummary, crate::PostgresBackendError> {
|
|
|
|
|
if !row.page_present {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_inspection_page"));
|
|
|
|
|
}
|
|
|
|
|
let signature = match inspection_required(row.signature, "raw_transaction_inspection_signature") {
|
|
|
|
|
std::result::Result::Ok(value) => match fixed_bytes::<64>(value) {
|
|
|
|
|
std::result::Result::Ok(bytes) => ksp_store_api::RawTransactionSignature::new(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_transaction_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_transaction_inspection_slot") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let block_time = match decode_optional_timestamp(row.block_time_unix_millis, "raw_transaction_inspection_block_time") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let format_id_raw = match inspection_required(row.format_id, "raw_transaction_inspection_format_id") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let format_id = match decode_format_id(format_id_raw) {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let format_version_raw = match inspection_required(row.format_version, "raw_transaction_inspection_format_version") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let format_version = match decode_u32_i64(format_version_raw, "raw_transaction_inspection_format_version") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let content_hash_raw = match inspection_required(row.content_hash, "raw_transaction_inspection_content_hash") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let content_hash = match fixed_bytes::<32>(content_hash_raw) {
|
|
|
|
|
std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value),
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let retention_raw = match inspection_required(row.retention_state, "raw_transaction_inspection_retention") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let retention_state = match decode_retention_state(retention_raw.as_str()) {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let hot_payload_present = match inspection_required(row.hot_payload_present, "raw_transaction_inspection_hot_payload") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let archive_payload_present = match inspection_required(row.archive_payload_present, "raw_transaction_inspection_archive_payload") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
|
|
|
|
};
|
|
|
|
|
let payload_size_bytes = match row.payload_size_bytes {
|
|
|
|
|
std::option::Option::Some(value) => match u64::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_transaction_inspection_payload_size")),
|
|
|
|
|
},
|
|
|
|
|
std::option::Option::None => std::option::Option::None,
|
|
|
|
|
};
|
|
|
|
|
match retention_state {
|
|
|
|
|
ksp_store_api::RawRetentionState::Full if !hot_payload_present || archive_payload_present => {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_inspection_full_shape"));
|
|
|
|
|
},
|
|
|
|
|
ksp_store_api::RawRetentionState::Archived if hot_payload_present || !archive_payload_present => {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_inspection_archived_shape"));
|
|
|
|
|
},
|
|
|
|
|
ksp_store_api::RawRetentionState::Purged if hot_payload_present || archive_payload_present || block_time.is_some() => {
|
|
|
|
|
return std::result::Result::Err(data_invalid("raw_transaction_inspection_purged_shape"));
|
|
|
|
|
},
|
|
|
|
|
ksp_store_api::RawRetentionState::Full | ksp_store_api::RawRetentionState::Archived | ksp_store_api::RawRetentionState::Purged => {},
|
|
|
|
|
_ => return std::result::Result::Err(data_invalid("raw_transaction_inspection_retention")),
|
|
|
|
|
}
|
|
|
|
|
let reference = ksp_store_api::RawTransactionReference::new(network.clone(), signature);
|
|
|
|
|
return match ksp_store_api::RawTransactionSummary::try_new(
|
|
|
|
|
reference,
|
|
|
|
|
slot,
|
|
|
|
|
block_time,
|
|
|
|
|
format_id,
|
|
|
|
|
format_version,
|
|
|
|
|
content_hash,
|
|
|
|
|
payload_size_bytes,
|
|
|
|
|
retention_state,
|
|
|
|
|
) {
|
|
|
|
|
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
|
|
|
|
std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_transaction_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 raw_list_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawListDbRow, crate::PostgresBackendError> {
|
|
|
|
|
let signature = match row.try_get::<_, std::vec::Vec<u8>>("signature") {
|
|
|
|
|
std::result::Result::Ok(value) => value,
|
|
|
|
|
|