// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs // version: 4 pub(crate) mod cursor; const DELETE_ARCHIVE_PAYLOAD_SQL: &str = "DELETE FROM ksp_raw_transaction_archive_payloads WHERE signature = $1"; 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"; const GET_RETENTION_SQL: &str = "SELECT retention_state FROM ksp_raw_transactions WHERE signature = $1"; const GET_TOMBSTONE_SQL: &str = "SELECT signature, slot::text AS slot_text, block_time_unix_millis, format_id, format_version, content_hash, retention_state FROM ksp_raw_transactions WHERE signature = $1"; 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_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 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_RETENTION_TRANSACTION_SQL: &str = "SELECT block_time_unix_millis, payload, retention_state FROM ksp_raw_transactions WHERE signature = $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"; const UPDATE_ARCHIVED_TRANSACTION_SQL: &str = "UPDATE ksp_raw_transactions SET payload = NULL, retention_state = 'archived' WHERE signature = $1 AND retention_state = 'full'"; const UPDATE_PURGED_TRANSACTION_SQL: &str = "UPDATE ksp_raw_transactions SET block_time_unix_millis = NULL, payload = NULL, retention_state = 'purged' WHERE signature = $1 AND retention_state = 'archived'"; struct RawListDbRow { signature: std::vec::Vec, slot_text: std::string::String, } struct RawObservationDbRow { acquisition_method: std::string::String, capture_session_id: std::option::Option, commitment: std::option::Option, endpoint_id: std::option::Option, filter_id: std::option::Option, observation_key: std::vec::Vec, observed_at_unix_millis: std::option::Option, origin: std::string::String, protocol: std::string::String, provider: std::string::String, received_at_unix_millis: i64, source_payload_hash: std::option::Option>, source_payload_size_bytes: std::option::Option, transaction_signature: std::vec::Vec, } struct RawRetentionDbRow { block_time_unix_millis: std::option::Option, payload: std::option::Option>, retention_state: std::string::String, } struct RawTransactionDbRow { archive_payload: std::option::Option>, block_time_unix_millis: std::option::Option, content_hash: std::vec::Vec, format_id: std::string::String, format_version: i64, payload: std::option::Option>, retention_state: std::string::String, signature: std::vec::Vec, slot_text: std::string::String, } struct RawTombstoneDbRow { block_time_unix_millis: std::option::Option, content_hash: std::vec::Vec, format_id: std::string::String, format_version: i64, retention_state: std::string::String, signature: std::vec::Vec, slot_text: std::string::String, } /// Reads one canonical RAW transaction from the physical PostgreSQL backend. pub(crate) async fn get_raw_transaction( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawTransactionReference, ) -> std::result::Result, crate::PostgresBackendError> { let network_result = ensure_network(network, reference, "raw_transaction_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 signature = reference.signature(); let signature_bytes: &[u8] = signature.as_bytes(); let rows_result = client.query(GET_TRANSACTION_SQL, &[&signature_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_transaction_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_transaction_cardinality")); } let row = match rows.first() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_transaction_cardinality")), }; let physical = match raw_transaction_db_row(row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; 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, 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, network: &ksp_store_api::RawNetworkId, observation_key: &ksp_store_api::RawObservationKey, ) -> std::result::Result, crate::PostgresBackendError> { let client_result = pool.get().await; let client = match client_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), }; let observation_key_bytes: &[u8] = observation_key.as_bytes(); let rows_result = client.query(GET_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_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_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_observation_cardinality")), }; let physical = match raw_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_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_observation_identity")); } return std::result::Result::Ok(std::option::Option::Some(decoded)); } /// Reads one RAW transaction retention state from the physical PostgreSQL backend. pub(crate) async fn get_raw_transaction_retention_state( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawTransactionReference, ) -> std::result::Result, crate::PostgresBackendError> { let network_result = ensure_network(network, reference, "raw_retention_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 signature = reference.signature(); let signature_bytes: &[u8] = signature.as_bytes(); let rows_result = client.query(GET_RETENTION_SQL, &[&signature_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_retention_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_retention_cardinality")); } let row = match rows.first() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_cardinality")), }; let state_result = row.try_get::<_, std::string::String>("retention_state"); let state = match state_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_decode")), }; let decoded = match decode_retention_state(state.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return std::result::Result::Ok(std::option::Option::Some(decoded)); } /// Reads one minimal RAW transaction tombstone from the physical PostgreSQL backend. pub(crate) async fn get_raw_transaction_tombstone( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawTransactionReference, ) -> std::result::Result, crate::PostgresBackendError> { let network_result = ensure_network(network, reference, "raw_tombstone_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 signature = reference.signature(); let signature_bytes: &[u8] = signature.as_bytes(); let rows_result = client.query(GET_TOMBSTONE_SQL, &[&signature_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_tombstone_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_tombstone_cardinality")); } let row = match rows.first() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_tombstone_cardinality")), }; let physical = match raw_tombstone_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_tombstone_row(network, physical) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let tombstone = match decoded { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Ok(std::option::Option::None), }; if tombstone.reference() != reference { return std::result::Result::Err(data_invalid("raw_tombstone_identity")); } return std::result::Result::Ok(std::option::Option::Some(tombstone)); } /// Persists one canonical RAW transaction and one observation atomically. pub(crate) async fn persist_raw_transaction_acquisition( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, raw_transaction: ksp_store_api::RawTransaction, observation: ksp_store_api::RawTransactionObservation, mode: ksp_store_api::RawTransactionAcquisitionMode, ) -> std::result::Result { let input_result = ensure_acquisition_inputs(network, &raw_transaction, &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_acquisition_begin")), }; let insert_result = insert_canonical_transaction(&sql_transaction, &raw_transaction).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_transaction_row(&sql_transaction, raw_transaction.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_acquisition_conflict_missing")), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let comparison_result = compare_existing_transaction(network, locked, &raw_transaction); let comparison = match comparison_result { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match comparison { ExistingTransactionMatch::Active => ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, ExistingTransactionMatch::Purged => { if mode == ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate { let rehydrate_result = rehydrate_transaction(&sql_transaction, &raw_transaction).await; if let std::result::Result::Err(error) = rehydrate_result { return std::result::Result::Err(error); } ksp_store_api::RawEntityWriteOutcome::Rehydrated } else { let commit_result = sql_transaction.commit().await; if commit_result.is_err() { return std::result::Result::Err(write_failed("raw_acquisition_commit")); } return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new( ksp_store_api::RawEntityWriteOutcome::SkippedPurged, ksp_store_api::RawObservationWriteOutcome::NotRecorded, )); } }, } }; let observation_result = persist_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_acquisition_commit")); } return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new(entity_outcome, observation_outcome)); } /// Persists one additional observation for an already known RAW transaction. pub(crate) async fn record_raw_transaction_observation( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, observation: ksp_store_api::RawTransactionObservation, ) -> std::result::Result { let network_result = ensure_network(network, observation.transaction(), "raw_observation_write_network"); if let std::result::Result::Err(error) = network_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_observation_begin")), }; let signature = observation.transaction().signature(); let signature_bytes: &[u8] = signature.as_bytes(); let state_row_result = sql_transaction.query_opt(LOCK_TRANSACTION_STATE_SQL, &[&signature_bytes]).await; let state_row = match state_row_result { std::result::Result::Ok(std::option::Option::Some(value)) => value, std::result::Result::Ok(std::option::Option::None) => return std::result::Result::Err(reference_not_found("raw_observation_transaction")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_lock_transaction")), }; let state_result = state_row.try_get::<_, std::string::String>("retention_state"); let state = match state_result { std::result::Result::Ok(value) => match decode_retention_state(value.as_str()) { std::result::Result::Ok(decoded) => decoded, std::result::Result::Err(error) => return std::result::Result::Err(error), }, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_transaction_state")), }; if state == ksp_store_api::RawRetentionState::Purged { let commit_result = sql_transaction.commit().await; if commit_result.is_err() { return std::result::Result::Err(write_failed("raw_observation_commit")); } return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::NotRecorded); } let observation_result = persist_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_observation_commit")); } return std::result::Result::Ok(outcome); } /// Applies one atomic compare-and-transition RAW transaction retention mutation. pub(crate) async fn transition_raw_transaction_retention( pool: &deadpool_postgres::Pool, network: &ksp_store_api::RawNetworkId, transition: ksp_store_api::RawTransactionRetentionTransition, ) -> std::result::Result { let input_result = ensure_retention_transition_inputs(network, &transition); 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_retention_begin")), }; let signature = transition.reference().signature(); let signature_bytes: &[u8] = signature.as_bytes(); let row_result = sql_transaction.query_opt(LOCK_RETENTION_TRANSACTION_SQL, &[&signature_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::Err(reference_not_found("raw_retention_transition_reference")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_lock_transaction")), }; let physical = match raw_retention_db_row(&row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let current = match decode_retention_state(physical.retention_state.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let shape_result = validate_retention_shape(&sql_transaction, signature_bytes, &physical, current).await; if let std::result::Result::Err(error) = shape_result { return std::result::Result::Err(error); } let decision = match retention_transition_decision(current, transition.expected(), transition.target()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; match decision { RetentionTransitionDecision::AlreadyAtTarget => { return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::AlreadyAtTarget).await; }, RetentionTransitionDecision::ExpectedStateMismatch => { return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::ExpectedStateMismatch).await; }, RetentionTransitionDecision::Archive => { let payload = match physical.payload.as_ref() { std::option::Option::Some(value) => value.as_slice(), std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_full_payload")), }; let archive_result = archive_full_transaction(&sql_transaction, signature_bytes, payload).await; if let std::result::Result::Err(error) = archive_result { return std::result::Result::Err(error); } }, RetentionTransitionDecision::Purge => { let purge_result = purge_archived_transaction(&sql_transaction, signature_bytes).await; if let std::result::Result::Err(error) = purge_result { return std::result::Result::Err(error); } }, } return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::Applied).await; } #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum ExistingTransactionMatch { Active, Purged, } #[derive(Clone, Copy, Debug, Eq, PartialEq)] enum RetentionTransitionDecision { AlreadyAtTarget, Archive, ExpectedStateMismatch, Purge, } async fn archive_full_transaction( sql_transaction: &deadpool_postgres::Transaction<'_>, signature_bytes: &[u8], payload: &[u8], ) -> std::result::Result<(), crate::PostgresBackendError> { let insert_result = sql_transaction.execute(INSERT_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes, &payload]).await; match insert_result { std::result::Result::Ok(1) => {}, std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_retention_archive_insert_cardinality")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_archive_insert")), } let update_result = sql_transaction.execute(UPDATE_ARCHIVED_TRANSACTION_SQL, &[&signature_bytes]).await; return match update_result { std::result::Result::Ok(1) => std::result::Result::Ok(()), std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_retention_archive_update_cardinality")), std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_retention_archive_update")), }; } async fn commit_retention_outcome( sql_transaction: deadpool_postgres::Transaction<'_>, outcome: ksp_store_api::RawRetentionWriteOutcome, ) -> std::result::Result { let commit_result = sql_transaction.commit().await; if commit_result.is_err() { return std::result::Result::Err(write_failed("raw_retention_commit")); } return std::result::Result::Ok(outcome); } async fn load_retention_archive_payload( sql_transaction: &deadpool_postgres::Transaction<'_>, signature_bytes: &[u8], ) -> std::result::Result>, crate::PostgresBackendError> { let row_result = sql_transaction.query_opt(GET_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await; let row = match row_result { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_archive_query")), }; let archive_row = match row { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Ok(std::option::Option::None), }; let payload_result = archive_row.try_get::<_, std::vec::Vec>("payload"); return match payload_result { std::result::Result::Ok(value) => std::result::Result::Ok(std::option::Option::Some(value)), std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_retention_archive_decode")), }; } async fn purge_archived_transaction( sql_transaction: &deadpool_postgres::Transaction<'_>, signature_bytes: &[u8], ) -> std::result::Result<(), crate::PostgresBackendError> { let delete_result = sql_transaction.execute(DELETE_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await; match delete_result { std::result::Result::Ok(1) => {}, std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_retention_purge_archive_cardinality")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_retention_purge_archive")), } let update_result = sql_transaction.execute(UPDATE_PURGED_TRANSACTION_SQL, &[&signature_bytes]).await; return match update_result { std::result::Result::Ok(1) => std::result::Result::Ok(()), std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_retention_purge_update_cardinality")), std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_retention_purge_update")), }; } fn raw_retention_db_row(row: &tokio_postgres::Row) -> std::result::Result { let block_time_unix_millis = match row.try_get::<_, std::option::Option>("block_time_unix_millis") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")), }; let payload = match row.try_get::<_, std::option::Option>>("payload") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")), }; let retention_state = match row.try_get::<_, std::string::String>("retention_state") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_retention_transition_decode")), }; return std::result::Result::Ok(RawRetentionDbRow { block_time_unix_millis, payload, retention_state }); } fn retention_transition_decision( current: ksp_store_api::RawRetentionState, expected: ksp_store_api::RawRetentionState, target: ksp_store_api::RawRetentionState, ) -> std::result::Result { if current == target { return std::result::Result::Ok(RetentionTransitionDecision::AlreadyAtTarget); } if current != expected { return std::result::Result::Ok(RetentionTransitionDecision::ExpectedStateMismatch); } return match (expected, target) { (ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Archived) => std::result::Result::Ok(RetentionTransitionDecision::Archive), (ksp_store_api::RawRetentionState::Archived, ksp_store_api::RawRetentionState::Purged) => std::result::Result::Ok(RetentionTransitionDecision::Purge), (ksp_store_api::RawRetentionState::Full, ksp_store_api::RawRetentionState::Compacted) | (ksp_store_api::RawRetentionState::Compacted, ksp_store_api::RawRetentionState::Archived) => { std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction")) }, _ => std::result::Result::Err(data_invalid("raw_retention_transition_unreachable")), }; } async fn validate_retention_shape( sql_transaction: &deadpool_postgres::Transaction<'_>, signature_bytes: &[u8], row: &RawRetentionDbRow, state: ksp_store_api::RawRetentionState, ) -> std::result::Result<(), crate::PostgresBackendError> { let archive_payload = match load_retention_archive_payload(sql_transaction, signature_bytes).await { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; return match state { ksp_store_api::RawRetentionState::Full => { let payload = match row.payload.as_ref() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_full_shape")), }; let payload_result = validate_retained_payload_bytes(payload.as_slice(), "raw_retention_full_payload"); if let std::result::Result::Err(error) = payload_result { return std::result::Result::Err(error); } if archive_payload.is_some() { return std::result::Result::Err(data_invalid("raw_retention_full_archive_residue")); } std::result::Result::Ok(()) }, ksp_store_api::RawRetentionState::Archived => { if row.payload.is_some() { return std::result::Result::Err(data_invalid("raw_retention_archived_hot_payload")); } let payload = match archive_payload.as_ref() { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_retention_archived_payload_missing")), }; validate_retained_payload_bytes(payload.as_slice(), "raw_retention_archived_payload") }, ksp_store_api::RawRetentionState::Purged => { if row.payload.is_some() || archive_payload.is_some() || row.block_time_unix_millis.is_some() { return std::result::Result::Err(data_invalid("raw_retention_purged_shape")); } std::result::Result::Ok(()) }, ksp_store_api::RawRetentionState::Compacted => std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction")), _ => std::result::Result::Err(data_invalid("raw_retention_state")), }; } fn validate_retained_payload_bytes(payload: &[u8], phase: &'static str) -> std::result::Result<(), crate::PostgresBackendError> { if payload.is_empty() || payload.len() > ksp_store_api::MAX_RAW_PAYLOAD_BYTES { return std::result::Result::Err(data_invalid(phase)); } return std::result::Result::Ok(()); } async fn insert_canonical_transaction( sql_transaction: &deadpool_postgres::Transaction<'_>, raw_transaction: &ksp_store_api::RawTransaction, ) -> std::result::Result { let signature = raw_transaction.reference().signature(); let signature_bytes: &[u8] = signature.as_bytes(); let slot_text = raw_transaction.slot().to_string(); let block_time = match raw_transaction.block_time() { 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_acquisition_block_time")), }, std::option::Option::None => std::option::Option::None, }; let format_version = i64::from(raw_transaction.payload().format_version()); let content_hash = raw_transaction.payload().content_hash(); let content_hash_bytes: &[u8] = content_hash.as_bytes(); let payload_bytes = raw_transaction.payload().bytes(); let row_result = sql_transaction .query_opt( INSERT_TRANSACTION_SQL, &[ &signature_bytes, &slot_text.as_str(), &block_time, &raw_transaction.payload().format_id().as_str(), &format_version, &content_hash_bytes, &payload_bytes, ], ) .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_acquisition_insert_transaction")), }; } async fn load_locked_transaction_row( sql_transaction: &deadpool_postgres::Transaction<'_>, reference: &ksp_store_api::RawTransactionReference, ) -> std::result::Result, crate::PostgresBackendError> { let signature = reference.signature(); let signature_bytes: &[u8] = signature.as_bytes(); let row_result = sql_transaction.query_opt(LOCK_TRANSACTION_SQL, &[&signature_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_acquisition_lock_transaction")), }; let mut physical = match raw_transaction_db_row(&row) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let archive_result = sql_transaction.query_opt(GET_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await; physical.archive_payload = match archive_result { std::result::Result::Ok(std::option::Option::Some(archive_row)) => match archive_row.try_get::<_, std::vec::Vec>("payload") { std::result::Result::Ok(value) => std::option::Option::Some(value), std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_acquisition_archive_decode")), }, std::result::Result::Ok(std::option::Option::None) => std::option::Option::None, std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_acquisition_archive_query")), }; return std::result::Result::Ok(std::option::Option::Some(physical)); } fn compare_existing_transaction( network: &ksp_store_api::RawNetworkId, row: RawTransactionDbRow, incoming: &ksp_store_api::RawTransaction, ) -> std::result::Result { let state = match decode_retention_state(row.retention_state.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if state == ksp_store_api::RawRetentionState::Purged { if row.payload.is_some() || row.archive_payload.is_some() || row.block_time_unix_millis.is_some() { return std::result::Result::Err(data_invalid("raw_acquisition_purged_shape")); } 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_acquisition_purged_slot") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_id = match decode_format_id(row.format_id) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_version = match decode_u32_i64(row.format_version, "raw_acquisition_purged_format_version") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let content_hash = match fixed_bytes::<32>(row.content_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::RawTransactionReference::new(network.clone(), signature); let matches = reference.eq(incoming.reference()) && slot == incoming.slot() && format_id.as_str() == incoming.payload().format_id().as_str() && format_version == incoming.payload().format_version() && content_hash == incoming.payload().content_hash(); if matches { return std::result::Result::Ok(ExistingTransactionMatch::Purged); } return std::result::Result::Err(conflict("raw_acquisition_purged_conflict")); } let stored_result = decode_raw_transaction_row(network, row); let stored = match stored_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_acquisition_active_shape")), std::result::Result::Err(error) => return std::result::Result::Err(error), }; if raw_transactions_equal(&stored, incoming) { return std::result::Result::Ok(ExistingTransactionMatch::Active); } return std::result::Result::Err(conflict("raw_acquisition_content_conflict")); } fn raw_transactions_equal(left: &ksp_store_api::RawTransaction, right: &ksp_store_api::RawTransaction) -> bool { return left.reference() == right.reference() && left.slot() == right.slot() && left.block_time() == right.block_time() && left.payload().format_id() == right.payload().format_id() && left.payload().format_version() == right.payload().format_version() && left.payload().content_hash() == right.payload().content_hash() && left.payload().bytes() == right.payload().bytes(); } async fn rehydrate_transaction( sql_transaction: &deadpool_postgres::Transaction<'_>, raw_transaction: &ksp_store_api::RawTransaction, ) -> std::result::Result<(), crate::PostgresBackendError> { let signature = raw_transaction.reference().signature(); let signature_bytes: &[u8] = signature.as_bytes(); let block_time = match raw_transaction.block_time() { 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_rehydrate_block_time")), }, std::option::Option::None => std::option::Option::None, }; let payload_bytes = raw_transaction.payload().bytes(); let update_result = sql_transaction.execute(REHYDRATE_TRANSACTION_SQL, &[&signature_bytes, &block_time, &payload_bytes]).await; return match update_result { std::result::Result::Ok(1) => std::result::Result::Ok(()), std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_rehydrate_cardinality")), std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_rehydrate_update")), }; } async fn persist_observation_row( sql_transaction: &deadpool_postgres::Transaction<'_>, network: &ksp_store_api::RawNetworkId, observation: &ksp_store_api::RawTransactionObservation, ) -> std::result::Result { let provenance = observation.provenance(); let observation_key = observation.observation_key(); let observation_key_bytes: &[u8] = observation_key.as_bytes(); let signature = observation.transaction().signature(); let signature_bytes: &[u8] = signature.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_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_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_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 insert_result = sql_transaction .query_opt( INSERT_OBSERVATION_SQL, &[ &observation_key_bytes, &signature_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, ], ) .await; let inserted = match insert_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_observation_insert")), }; if inserted { return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::Inserted); } let existing_result = sql_transaction.query_opt(LOCK_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_observation_conflict_missing")), std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_conflict_query")), }; let physical = match raw_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_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_observation_content_conflict")); } fn ensure_acquisition_inputs( network: &ksp_store_api::RawNetworkId, raw_transaction: &ksp_store_api::RawTransaction, observation: &ksp_store_api::RawTransactionObservation, ) -> std::result::Result<(), crate::PostgresBackendError> { let transaction_network_result = ensure_network(network, raw_transaction.reference(), "raw_acquisition_transaction_network"); if let std::result::Result::Err(error) = transaction_network_result { return std::result::Result::Err(error); } let observation_network_result = ensure_network(network, observation.transaction(), "raw_acquisition_observation_network"); if let std::result::Result::Err(error) = observation_network_result { return std::result::Result::Err(error); } if observation.transaction() != raw_transaction.reference() { return std::result::Result::Err(conflict("raw_acquisition_reference_mismatch")); } return std::result::Result::Ok(()); } fn ensure_retention_transition_inputs( network: &ksp_store_api::RawNetworkId, transition: &ksp_store_api::RawTransactionRetentionTransition, ) -> std::result::Result<(), crate::PostgresBackendError> { let network_result = ensure_network(network, transition.reference(), "raw_retention_transition_network"); if let std::result::Result::Err(error) = network_result { return std::result::Result::Err(error); } if transition.expected() == ksp_store_api::RawRetentionState::Compacted || transition.target() == ksp_store_api::RawRetentionState::Compacted { return std::result::Result::Err(retention_compaction_unsupported("raw_retention_compaction")); } return std::result::Result::Ok(()); } 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_observation_origin_encode")), }; } fn conflict(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::Conflict, phase); } fn retention_compaction_unsupported(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::RetentionCompactionUnsupported, 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); } fn raw_list_db_row(row: &tokio_postgres::Row) -> std::result::Result { let signature = match row.try_get::<_, std::vec::Vec>("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 { let signature = match row.try_get::<_, std::vec::Vec>("signature") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_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_decode")), }; let block_time_unix_millis = match row.try_get::<_, std::option::Option>("block_time_unix_millis") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let format_id = match row.try_get::<_, std::string::String>("format_id") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let format_version = match row.try_get::<_, i64>("format_version") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let content_hash = match row.try_get::<_, std::vec::Vec>("content_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let payload = match row.try_get::<_, std::option::Option>>("payload") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let retention_state = match row.try_get::<_, std::string::String>("retention_state") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; let archive_payload = match row.try_get::<_, std::option::Option>>("archive_payload") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_decode")), }; return std::result::Result::Ok(RawTransactionDbRow { archive_payload, block_time_unix_millis, content_hash, format_id, format_version, payload, retention_state, signature, slot_text, }); } fn raw_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result { macro_rules! required { ($name:literal, $ty:ty) => { match row.try_get::<_, $ty>($name) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_observation_decode")), } }; } return std::result::Result::Ok(RawObservationDbRow { acquisition_method: required!("acquisition_method", std::string::String), capture_session_id: required!("capture_session_id", std::option::Option), commitment: required!("commitment", std::option::Option), endpoint_id: required!("endpoint_id", std::option::Option), filter_id: required!("filter_id", std::option::Option), observation_key: required!("observation_key", std::vec::Vec), observed_at_unix_millis: required!("observed_at_unix_millis", std::option::Option), origin: required!("origin", std::string::String), protocol: required!("protocol", std::string::String), provider: required!("provider", std::string::String), received_at_unix_millis: required!("received_at_unix_millis", i64), source_payload_hash: required!("source_payload_hash", std::option::Option>), source_payload_size_bytes: required!("source_payload_size_bytes", std::option::Option), transaction_signature: required!("transaction_signature", std::vec::Vec), }); } fn raw_tombstone_db_row(row: &tokio_postgres::Row) -> std::result::Result { let signature = match row.try_get::<_, std::vec::Vec>("signature") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_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_tombstone_decode")), }; let block_time_unix_millis = match row.try_get::<_, std::option::Option>("block_time_unix_millis") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_decode")), }; let format_id = match row.try_get::<_, std::string::String>("format_id") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_decode")), }; let format_version = match row.try_get::<_, i64>("format_version") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_decode")), }; let content_hash = match row.try_get::<_, std::vec::Vec>("content_hash") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_decode")), }; let retention_state = match row.try_get::<_, std::string::String>("retention_state") { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_decode")), }; return std::result::Result::Ok(RawTombstoneDbRow { block_time_unix_millis, content_hash, format_id, format_version, retention_state, signature, slot_text, }); } 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, ) -> std::result::Result, 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_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_block_time") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_id = match decode_format_id(row.format_id) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_version = match decode_u32_i64(row.format_version, "raw_transaction_format_version") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let content_hash = match fixed_bytes::<32>(row.content_hash) { std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), std::result::Result::Err(error) => return std::result::Result::Err(error), }; let retention_state = match decode_retention_state(row.retention_state.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let payload_bytes = match retention_state { ksp_store_api::RawRetentionState::Full => { if row.archive_payload.is_some() { return std::result::Result::Err(data_invalid("raw_transaction_full_archive")); } match row.payload { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_transaction_full_payload")), } }, ksp_store_api::RawRetentionState::Archived => { if row.payload.is_some() { return std::result::Result::Err(data_invalid("raw_transaction_archived_hot_payload")); } match row.archive_payload { std::option::Option::Some(value) => value, std::option::Option::None => return std::result::Result::Err(data_invalid("raw_transaction_archive_payload")), } }, ksp_store_api::RawRetentionState::Purged => { if row.payload.is_some() || row.archive_payload.is_some() || block_time.is_some() { return std::result::Result::Err(data_invalid("raw_transaction_purged_shape")); } return std::result::Result::Ok(std::option::Option::None); }, _ => return std::result::Result::Err(data_invalid("raw_transaction_retention_state")), }; let payload = match ksp_store_api::RawPayload::try_new(format_id, format_version, payload_bytes.into_boxed_slice(), content_hash) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_payload")), }; let reference = ksp_store_api::RawTransactionReference::new(network.clone(), signature); return std::result::Result::Ok(std::option::Option::Some(ksp_store_api::RawTransaction::new(reference, slot, block_time, payload))); } fn decode_raw_observation_row( network: &ksp_store_api::RawNetworkId, row: RawObservationDbRow, ) -> std::result::Result { let observation_key = match fixed_bytes::<32>(row.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 signature = match fixed_bytes::<64>(row.transaction_signature) { std::result::Result::Ok(value) => ksp_store_api::RawTransactionSignature::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_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_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_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) { 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_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_observation_source_size")), } }, std::option::Option::None => provenance, }; let transaction = ksp_store_api::RawTransactionReference::new(network.clone(), signature); return std::result::Result::Ok(ksp_store_api::RawTransactionObservation::new(observation_key, transaction, provenance)); } fn decode_raw_tombstone_row( network: &ksp_store_api::RawNetworkId, row: RawTombstoneDbRow, ) -> std::result::Result, crate::PostgresBackendError> { let state = match decode_retention_state(row.retention_state.as_str()) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; if state != ksp_store_api::RawRetentionState::Purged { return std::result::Result::Ok(std::option::Option::None); } if row.block_time_unix_millis.is_some() { return std::result::Result::Err(data_invalid("raw_tombstone_block_time")); } 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_tombstone_slot") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_id = match decode_format_id(row.format_id) { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let format_version = match decode_u32_i64(row.format_version, "raw_tombstone_format_version") { std::result::Result::Ok(value) => value, std::result::Result::Err(error) => return std::result::Result::Err(error), }; let content_hash = match fixed_bytes::<32>(row.content_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::RawTransactionReference::new(network.clone(), signature); let tombstone = match ksp_store_api::RawTransactionTombstone::try_new(reference, slot, format_id, format_version, content_hash) { std::result::Result::Ok(value) => value, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_tombstone_model")), }; return std::result::Result::Ok(std::option::Option::Some(tombstone)); } fn decode_format_id(value: std::string::String) -> std::result::Result { return match ksp_store_api::RawFormatId::new(value) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_format_id")), }; } fn decode_origin(value: &str) -> std::result::Result { return match value { "backfill" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Backfill), "import" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Import), "live" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Live), "repair" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Repair), "replay" => std::result::Result::Ok(ksp_store_api::RawAcquisitionOrigin::Replay), _ => std::result::Result::Err(data_invalid("raw_observation_origin")), }; } fn decode_provenance_code(value: std::string::String) -> std::result::Result { return match ksp_store_api::RawProvenanceCode::new(value) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_provenance_code")), }; } fn decode_retention_state(value: &str) -> std::result::Result { return match value { "full" => std::result::Result::Ok(ksp_store_api::RawRetentionState::Full), "archived" => std::result::Result::Ok(ksp_store_api::RawRetentionState::Archived), "purged" => std::result::Result::Ok(ksp_store_api::RawRetentionState::Purged), _ => std::result::Result::Err(data_invalid("raw_retention_state")), }; } fn decode_optional_timestamp( value: std::option::Option, phase: &'static str, ) -> std::result::Result, crate::PostgresBackendError> { return match value { std::option::Option::Some(inner) => match decode_timestamp_i64(inner, phase) { std::result::Result::Ok(decoded) => std::result::Result::Ok(std::option::Option::Some(decoded)), std::result::Result::Err(error) => std::result::Result::Err(error), }, std::option::Option::None => std::result::Result::Ok(std::option::Option::None), }; } fn decode_timestamp_i64(value: i64, phase: &'static str) -> std::result::Result { let unsigned = match u64::try_from(value) { std::result::Result::Ok(decoded) => decoded, std::result::Result::Err(_) => return std::result::Result::Err(data_invalid(phase)), }; return match ksp_store_api::RawTimestamp::from_unix_millis(unsigned) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn decode_u32_i64(value: i64, phase: &'static str) -> std::result::Result { return match u32::try_from(value) { std::result::Result::Ok(decoded) if decoded > 0 => std::result::Result::Ok(decoded), _ => std::result::Result::Err(data_invalid(phase)), }; } fn decode_u64_decimal(value: &str, phase: &'static str) -> std::result::Result { return match value.parse::() { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn decode_u64_i64(value: i64, phase: &'static str) -> std::result::Result { return match u64::try_from(value) { std::result::Result::Ok(decoded) => std::result::Result::Ok(decoded), std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)), }; } fn fixed_bytes(value: std::vec::Vec) -> 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("raw_fixed_bytes")), }; } fn ensure_network( network: &ksp_store_api::RawNetworkId, reference: &ksp_store_api::RawTransactionReference, 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); } #[cfg(test)] #[path = "../unit_tests/raw_transaction.rs"] mod tests;