|
|
|
|
@@ -1,10 +1,18 @@
|
|
|
|
|
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
|
|
|
|
|
// version: 1
|
|
|
|
|
// version: 2
|
|
|
|
|
|
|
|
|
|
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_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 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 RawObservationDbRow {
|
|
|
|
|
acquisition_method: std::string::String,
|
|
|
|
|
@@ -227,6 +235,426 @@ pub(crate) async fn get_raw_transaction_tombstone(
|
|
|
|
|
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<ksp_store_api::RawAcquisitionWriteOutcome, crate::PostgresBackendError> {
|
|
|
|
|
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<ksp_store_api::RawObservationWriteOutcome, crate::PostgresBackendError> {
|
|
|
|
|
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);
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
|
|
|
|
enum ExistingTransactionMatch {
|
|
|
|
|
Active,
|
|
|
|
|
Purged,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
async fn insert_canonical_transaction(
|
|
|
|
|
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
|
|
|
|
raw_transaction: &ksp_store_api::RawTransaction,
|
|
|
|
|
) -> std::result::Result<bool, crate::PostgresBackendError> {
|
|
|
|
|
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<std::option::Option<RawTransactionDbRow>, 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<u8>>("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<ExistingTransactionMatch, 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 {
|
|
|
|
|
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<ksp_store_api::RawObservationWriteOutcome, crate::PostgresBackendError> {
|
|
|
|
|
let provenance = observation.provenance();
|
|
|
|
|
let observation_key = observation.observation_key();
|
|
|
|
|
let observation_key_bytes: &[u8] = observation_key.as_bytes();
|
|
|
|
|
let 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 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 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_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,
|
|
|
|
|
|