0.3.16-pre.005

This commit is contained in:
2026-09-21 19:31:35 +02:00
parent 19712aa7a6
commit 8b0fbf831e
14 changed files with 895 additions and 132 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
// version: 11
// version: 12
pub(crate) mod cursor;
@@ -712,7 +712,7 @@ pub(crate) async fn persist_raw_transaction_acquisition(
};
match comparison {
ExistingTransactionMatch::Active => (ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, canonical_variant_id),
ExistingTransactionMatch::ActiveIncomingTruncatedLogs => {
ExistingTransactionMatch::ActiveCompatibleLessComplete => {
let variant_result = persist_or_reuse_native_transaction_variant(&sql_transaction, &raw_transaction, linked_at).await;
let variant_id = match variant_result {
std::result::Result::Ok(value) => value,
@@ -896,7 +896,7 @@ pub(crate) async fn transition_raw_transaction_retention(
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ExistingTransactionMatch {
Active,
ActiveIncomingTruncatedLogs,
ActiveCompatibleLessComplete,
Purged,
}
@@ -1439,87 +1439,24 @@ fn compare_existing_transaction(
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);
}
if raw_transaction_incoming_truncated_log_messages_compatible(&stored, incoming) {
log_raw_transaction_compatible_truncated_log_messages(network, &stored, incoming);
return std::result::Result::Ok(ExistingTransactionMatch::ActiveIncomingTruncatedLogs);
}
log_raw_transaction_content_conflict(network, &stored, incoming);
return std::result::Result::Err(conflict("raw_acquisition_content_conflict"));
}
fn raw_transaction_incoming_truncated_log_messages_compatible(stored: &ksp_store_api::RawTransaction, incoming: &ksp_store_api::RawTransaction) -> bool {
if stored.reference() != incoming.reference()
|| stored.slot() != incoming.slot()
|| stored.block_time() != incoming.block_time()
|| stored.payload().format_id() != incoming.payload().format_id()
|| stored.payload().format_version() != incoming.payload().format_version()
{
return false;
}
let stored_payload = serde_json::from_slice::<serde_json::Value>(stored.payload().bytes());
let incoming_payload = serde_json::from_slice::<serde_json::Value>(incoming.payload().bytes());
let (stored_payload, incoming_payload) = match (stored_payload, incoming_payload) {
(std::result::Result::Ok(serde_json::Value::Object(stored)), std::result::Result::Ok(serde_json::Value::Object(incoming))) => (stored, incoming),
_ => return false,
let comparison = match ksp_store_api::compare_raw_transaction_variants(&stored, incoming) {
std::result::Result::Ok(value) => value,
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_acquisition_variant_comparison")),
};
if stored_payload.get("transaction") != incoming_payload.get("transaction")
|| stored_payload.get("version") != incoming_payload.get("version")
|| stored_payload.get("transactionIndex") != incoming_payload.get("transactionIndex")
|| json_object_other_fields_mismatch(&stored_payload, &incoming_payload, RAW_TRANSACTION_CONTENT_CONFLICT_PAYLOAD_FIELDS.as_slice())
{
return false;
}
let (stored_meta, incoming_meta) = match (stored_payload.get("meta"), incoming_payload.get("meta")) {
(std::option::Option::Some(serde_json::Value::Object(stored)), std::option::Option::Some(serde_json::Value::Object(incoming))) => (stored, incoming),
_ => return false,
};
if raw_meta_other_than_log_messages_mismatch(stored_meta, incoming_meta) {
return false;
}
let (stored_logs, incoming_logs) = match (stored_meta.get("logMessages"), incoming_meta.get("logMessages")) {
(std::option::Option::Some(serde_json::Value::Array(stored)), std::option::Option::Some(serde_json::Value::Array(incoming))) => {
(stored.as_slice(), incoming.as_slice())
return match comparison.relation() {
ksp_store_api::RawTransactionVariantRelation::Exact => std::result::Result::Ok(ExistingTransactionMatch::Active),
ksp_store_api::RawTransactionVariantRelation::CompatibleLessComplete => {
log_raw_transaction_compatible_truncated_log_messages(network, &stored, incoming);
std::result::Result::Ok(ExistingTransactionMatch::ActiveCompatibleLessComplete)
},
_ => return false,
ksp_store_api::RawTransactionVariantRelation::CompatibleMoreComplete
| ksp_store_api::RawTransactionVariantRelation::Conflict
| ksp_store_api::RawTransactionVariantRelation::Incomparable => {
log_raw_transaction_content_conflict(network, &stored, incoming);
std::result::Result::Err(conflict("raw_acquisition_content_conflict"))
},
_ => std::result::Result::Err(data_invalid("raw_acquisition_variant_relation")),
};
return raw_log_messages_incoming_truncated_compatible(stored_logs, incoming_logs);
}
fn raw_meta_other_than_log_messages_mismatch(
stored: &serde_json::Map<std::string::String, serde_json::Value>,
incoming: &serde_json::Map<std::string::String, serde_json::Value>,
) -> bool {
let stored_without_logs = stored.iter().filter(|(key, _)| return key.as_str() != "logMessages").collect::<std::collections::BTreeMap<_, _>>();
let incoming_without_logs = incoming.iter().filter(|(key, _)| return key.as_str() != "logMessages").collect::<std::collections::BTreeMap<_, _>>();
return stored_without_logs != incoming_without_logs;
}
fn raw_log_messages_incoming_truncated_compatible(stored: &[serde_json::Value], incoming: &[serde_json::Value]) -> bool {
if raw_log_messages_exact_truncation_marker_count(stored) != 0 {
return false;
}
let marker_index = match incoming.iter().position(raw_log_message_is_exact_truncation_marker) {
std::option::Option::Some(value) => value,
std::option::Option::None => return false,
};
if raw_log_messages_exact_truncation_marker_count(incoming) != 1 || marker_index >= stored.len() {
return false;
}
if incoming[..marker_index] != stored[..marker_index] {
return false;
}
return true;
}
fn raw_log_messages_exact_truncation_marker_count(values: &[serde_json::Value]) -> usize {
return values.iter().filter(|value| return raw_log_message_is_exact_truncation_marker(value)).count();
}
fn raw_log_message_is_exact_truncation_marker(value: &serde_json::Value) -> bool {
return matches!(value, serde_json::Value::String(line) if line == "Log truncated");
}
fn raw_transaction_content_conflict_diagnostic(
@@ -1863,16 +1800,6 @@ async fn log_raw_transaction_content_conflict_provenance(
return;
}
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,