v0.3.3-pre.006
This commit is contained in:
@@ -1,5 +1,7 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
pub(crate) mod cursor;
|
||||
|
||||
const GET_ARCHIVE_PAYLOAD_SQL: &str = "SELECT payload FROM ksp_raw_transaction_archive_payloads WHERE signature = $1";
|
||||
const GET_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";
|
||||
@@ -8,12 +10,19 @@ const GET_TOMBSTONE_SQL: &str = "SELECT signature, slot::text AS slot_text, bloc
|
||||
const GET_TRANSACTION_SQL: &str = "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.payload, transaction_row.retention_state, archive_row.payload AS archive_payload 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 transaction_row.signature = $1";
|
||||
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 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";
|
||||
const LOCK_TRANSACTION_SQL: &str = "SELECT signature, slot::text AS slot_text, block_time_unix_millis, format_id, format_version, content_hash, payload, retention_state, NULL::BYTEA AS archive_payload FROM ksp_raw_transactions WHERE signature = $1 FOR UPDATE";
|
||||
const LOCK_TRANSACTION_STATE_SQL: &str = "SELECT retention_state FROM ksp_raw_transactions WHERE signature = $1 FOR UPDATE";
|
||||
const REHYDRATE_TRANSACTION_SQL: &str =
|
||||
"UPDATE ksp_raw_transactions SET block_time_unix_millis = $2, payload = $3, retention_state = 'full' WHERE signature = $1";
|
||||
|
||||
struct RawListDbRow {
|
||||
signature: std::vec::Vec<u8>,
|
||||
slot_text: std::string::String,
|
||||
}
|
||||
|
||||
struct RawObservationDbRow {
|
||||
acquisition_method: std::string::String,
|
||||
capture_session_id: std::option::Option<std::string::String>,
|
||||
@@ -94,6 +103,82 @@ pub(crate) async fn get_raw_transaction(
|
||||
return decode_raw_transaction_row(network, physical);
|
||||
}
|
||||
|
||||
/// Lists deterministic canonical RAW transaction references using PostgreSQL keyset pagination.
|
||||
pub(crate) async fn list_raw_transactions(
|
||||
pool: &deadpool_postgres::Pool,
|
||||
network: &ksp_store_api::RawNetworkId,
|
||||
query: &ksp_store_api::RawTransactionQuery,
|
||||
) -> std::result::Result<ksp_store_api::RawPage<ksp_store_api::RawTransactionReference>, crate::PostgresBackendError> {
|
||||
if query.network() != network {
|
||||
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, "raw_transaction_list_network"));
|
||||
}
|
||||
let (requested_usize, sql_limit) = match crate::raw_transaction_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_transaction_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_signature = decoded_cursor.as_ref().map(|value| return value.last_signature.to_vec());
|
||||
let sql = match query.direction() {
|
||||
ksp_store_api::RawSortDirection::Ascending => LIST_TRANSACTIONS_ASC_SQL,
|
||||
ksp_store_api::RawSortDirection::Descending => LIST_TRANSACTIONS_DESC_SQL,
|
||||
_ => {
|
||||
return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::QueryInvalid, "raw_transaction_list_direction"));
|
||||
},
|
||||
};
|
||||
let rows_result = client.query(sql, &[&start_text, &end_text, &cursor_slot_text, &cursor_signature, &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_transaction_list_query"));
|
||||
},
|
||||
};
|
||||
let mut decoded = std::vec::Vec::with_capacity(rows.len());
|
||||
for row in rows {
|
||||
let physical = match raw_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_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_transaction_list_page")),
|
||||
};
|
||||
match crate::encode_raw_transaction_cursor(query, last.0, &last.1.signature()) {
|
||||
std::result::Result::Ok(value) => std::option::Option::Some(value),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
}
|
||||
} else {
|
||||
std::option::Option::None
|
||||
};
|
||||
let items = decoded.into_iter().map(|value| return value.1).collect();
|
||||
return std::result::Result::Ok(ksp_store_api::RawPage::new(items, next_cursor));
|
||||
}
|
||||
|
||||
/// Reads one RAW transaction observation from the physical PostgreSQL backend.
|
||||
pub(crate) async fn get_raw_transaction_observation(
|
||||
pool: &deadpool_postgres::Pool,
|
||||
@@ -655,6 +740,18 @@ fn write_failed(phase: &'static str) -> crate::PostgresBackendError {
|
||||
return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WriteFailed, 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,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_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_transaction_list_decode")),
|
||||
};
|
||||
return std::result::Result::Ok(RawListDbRow { signature, slot_text });
|
||||
}
|
||||
|
||||
fn raw_transaction_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawTransactionDbRow, crate::PostgresBackendError> {
|
||||
let signature = match row.try_get::<_, std::vec::Vec<u8>>("signature") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
@@ -772,6 +869,22 @@ fn raw_tombstone_db_row(row: &tokio_postgres::Row) -> std::result::Result<RawTom
|
||||
});
|
||||
}
|
||||
|
||||
fn decode_raw_list_row(
|
||||
network: &ksp_store_api::RawNetworkId,
|
||||
row: RawListDbRow,
|
||||
) -> std::result::Result<(u64, ksp_store_api::RawTransactionReference), crate::PostgresBackendError> {
|
||||
let signature = match fixed_bytes::<64>(row.signature) {
|
||||
std::result::Result::Ok(value) => ksp_store_api::RawTransactionSignature::new(value),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let slot = match decode_u64_decimal(row.slot_text.as_str(), "raw_transaction_list_slot") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let reference = ksp_store_api::RawTransactionReference::new(network.clone(), signature);
|
||||
return std::result::Result::Ok((slot, reference));
|
||||
}
|
||||
|
||||
fn decode_raw_transaction_row(
|
||||
network: &ksp_store_api::RawNetworkId,
|
||||
row: RawTransactionDbRow,
|
||||
|
||||
Reference in New Issue
Block a user