0.3.16-pre.006
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
|
||||
// version: 12
|
||||
// version: 13
|
||||
|
||||
pub(crate) mod cursor;
|
||||
|
||||
@@ -48,6 +48,7 @@ const REHYDRATE_TRANSACTION_SQL: &str =
|
||||
"UPDATE ksp_raw_transactions SET block_time_unix_millis = $2, payload = $3, retention_state = 'full' WHERE signature = $1";
|
||||
const SELECT_CANONICAL_VARIANT_ID_SQL: &str =
|
||||
"SELECT canonical_variant_id::TEXT AS canonical_variant_id_text FROM ksp_raw_transaction_canonical_selectors WHERE transaction_signature = $1";
|
||||
const SELECT_CANONICAL_VARIANT_SELECTOR_SQL: &str = "SELECT canonical_variant_id::TEXT AS canonical_variant_id_text, canonical_revision::TEXT AS canonical_revision_text FROM ksp_raw_transaction_canonical_selectors WHERE transaction_signature = $1";
|
||||
const SELECT_EXACT_NATIVE_VARIANT_SQL: &str = "SELECT variant_id::TEXT AS variant_id_text FROM ksp_raw_transaction_variants WHERE transaction_signature = $1 AND origin_kind = 'native' AND slot = $2::TEXT::NUMERIC AND block_time_unix_millis IS NOT DISTINCT FROM $3 AND format_id = $4 AND format_version = $5 AND content_hash = $6 AND payload = $7 AND retention_state = 'full' ORDER BY variant_id ASC LIMIT 1";
|
||||
const SELECT_MAX_VARIANT_ID_SQL: &str =
|
||||
"SELECT COALESCE(MAX(variant_id), 0)::TEXT AS variant_id_text FROM ksp_raw_transaction_variants WHERE transaction_signature = $1";
|
||||
@@ -55,8 +56,11 @@ const SELECT_TRANSACTION_VARIANT_LINK_SQL: &str =
|
||||
"SELECT transaction_signature, variant_id::TEXT AS variant_id_text FROM ksp_raw_transaction_observation_variants WHERE observation_key = $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_CANONICAL_TRANSACTION_PROJECTION_SQL: &str = "UPDATE ksp_raw_transactions SET slot = $2::TEXT::NUMERIC, block_time_unix_millis = $3, format_id = $4, format_version = $5, content_hash = $6, payload = $7, retention_state = 'full' WHERE signature = $1";
|
||||
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'";
|
||||
const UPDATE_REHYDRATED_TRANSACTION_VARIANT_SQL: &str = "UPDATE ksp_raw_transaction_variants SET block_time_unix_millis = $3, payload = $4, retention_state = 'full' WHERE transaction_signature = $1 AND variant_id = $2::TEXT::NUMERIC";
|
||||
const UPDATE_TRANSACTION_VARIANT_SELECTOR_SQL: &str = "UPDATE ksp_raw_transaction_canonical_selectors SET canonical_variant_id = $4::TEXT::NUMERIC, canonical_revision = $5::TEXT::NUMERIC, updated_at_unix_millis = $6 WHERE transaction_signature = $1 AND canonical_variant_id = $2::TEXT::NUMERIC AND canonical_revision = $3::TEXT::NUMERIC";
|
||||
const UPDATE_TRANSACTION_VARIANT_TO_FULL_SQL: &str = "UPDATE ksp_raw_transaction_variants SET payload = $8, retention_state = 'full' WHERE transaction_signature = $1 AND variant_id = $2::TEXT::NUMERIC AND slot = $3::TEXT::NUMERIC AND block_time_unix_millis IS NOT DISTINCT FROM $4 AND format_id = $5 AND format_version = $6 AND content_hash = $7 AND ((retention_state = 'archived' AND payload IS NULL) OR (retention_state = 'full' AND payload = $8))";
|
||||
|
||||
struct RawLogMessagesContentConflictDiagnostic {
|
||||
available: bool,
|
||||
@@ -697,10 +701,10 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let (entity_outcome, observed_variant_id) = if inserted {
|
||||
(ksp_store_api::RawEntityWriteOutcome::Inserted, canonical_variant_id)
|
||||
let (entity_outcome, observed_variant_id, canonical_promotion) = if inserted {
|
||||
(ksp_store_api::RawEntityWriteOutcome::Inserted, canonical_variant_id, std::option::Option::None)
|
||||
} else {
|
||||
let comparison_result = compare_existing_transaction(network, locked, &raw_transaction);
|
||||
let comparison_result = compare_existing_transaction(network, locked.clone(), &raw_transaction);
|
||||
let comparison = match comparison_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => {
|
||||
@@ -711,14 +715,28 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
},
|
||||
};
|
||||
match comparison {
|
||||
ExistingTransactionMatch::Active => (ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, canonical_variant_id),
|
||||
ExistingTransactionMatch::Active => (ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, canonical_variant_id, std::option::Option::None),
|
||||
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,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
(ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, variant_id)
|
||||
(ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, variant_id, std::option::Option::None)
|
||||
},
|
||||
ExistingTransactionMatch::ActiveCompatibleMoreComplete => {
|
||||
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,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let promotion_result =
|
||||
promote_more_complete_transaction_variant(&sql_transaction, &locked, &raw_transaction, canonical_variant_id, variant_id, linked_at).await;
|
||||
let promotion = match promotion_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
(ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, variant_id, std::option::Option::Some(promotion))
|
||||
},
|
||||
ExistingTransactionMatch::Purged => {
|
||||
if mode == ksp_store_api::RawTransactionAcquisitionMode::ForceRehydrate {
|
||||
@@ -730,7 +748,7 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
if let std::result::Result::Err(error) = variant_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
(ksp_store_api::RawEntityWriteOutcome::Rehydrated, canonical_variant_id)
|
||||
(ksp_store_api::RawEntityWriteOutcome::Rehydrated, canonical_variant_id, std::option::Option::None)
|
||||
} else {
|
||||
let commit_result = sql_transaction.commit().await;
|
||||
if commit_result.is_err() {
|
||||
@@ -757,6 +775,9 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
if commit_result.is_err() {
|
||||
return std::result::Result::Err(write_failed("raw_acquisition_commit"));
|
||||
}
|
||||
if let std::option::Option::Some(transition) = canonical_promotion {
|
||||
log_raw_transaction_canonical_promotion_committed(network, raw_transaction.slot(), transition);
|
||||
}
|
||||
return std::result::Result::Ok(ksp_store_api::RawAcquisitionWriteOutcome::new(entity_outcome, observation_outcome));
|
||||
}
|
||||
|
||||
@@ -893,10 +914,19 @@ pub(crate) async fn transition_raw_transaction_retention(
|
||||
return commit_retention_outcome(sql_transaction, ksp_store_api::RawRetentionWriteOutcome::Applied).await;
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
struct RawCanonicalPromotionTransition {
|
||||
from_revision: u64,
|
||||
from_variant_id: u64,
|
||||
to_revision: u64,
|
||||
to_variant_id: u64,
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum ExistingTransactionMatch {
|
||||
Active,
|
||||
ActiveCompatibleLessComplete,
|
||||
ActiveCompatibleMoreComplete,
|
||||
Purged,
|
||||
}
|
||||
|
||||
@@ -1229,6 +1259,182 @@ async fn persist_or_reuse_native_transaction_variant(
|
||||
};
|
||||
}
|
||||
|
||||
async fn promote_more_complete_transaction_variant(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
locked: &RawTransactionDbRow,
|
||||
incoming: &ksp_store_api::RawTransaction,
|
||||
expected_canonical_variant_id: u64,
|
||||
promoted_variant_id: u64,
|
||||
updated_at_unix_millis: i64,
|
||||
) -> std::result::Result<RawCanonicalPromotionTransition, crate::PostgresBackendError> {
|
||||
if promoted_variant_id == expected_canonical_variant_id {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_identity"));
|
||||
}
|
||||
let signature = incoming.reference().signature();
|
||||
let signature_bytes: &[u8] = signature.as_bytes();
|
||||
if locked.signature.as_slice() != signature_bytes {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_reference"));
|
||||
}
|
||||
let selector_result = sql_transaction.query_opt(SELECT_CANONICAL_VARIANT_SELECTOR_SQL, &[&signature_bytes]).await;
|
||||
let selector_row = match selector_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_variant_promotion_selector_missing")),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_promotion_selector_query")),
|
||||
};
|
||||
let canonical_variant_id_text = match selector_row.try_get::<_, std::string::String>("canonical_variant_id_text") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_promotion_selector_id")),
|
||||
};
|
||||
let current_canonical_variant_id = match decode_u64_decimal(canonical_variant_id_text.as_str(), "raw_variant_promotion_selector_id") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
if current_canonical_variant_id != expected_canonical_variant_id {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_selector_changed"));
|
||||
}
|
||||
let canonical_revision_text = match selector_row.try_get::<_, std::string::String>("canonical_revision_text") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_promotion_revision")),
|
||||
};
|
||||
let current_revision = match decode_u64_decimal(canonical_revision_text.as_str(), "raw_variant_promotion_revision") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let next_revision = match next_transaction_canonical_revision(current_revision) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let preserve_result = preserve_previous_canonical_variant(sql_transaction, locked, expected_canonical_variant_id).await;
|
||||
if let std::result::Result::Err(error) = preserve_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let archive_cleanup_result = sql_transaction.execute(DELETE_ARCHIVE_PAYLOAD_SQL, &[&signature_bytes]).await;
|
||||
match archive_cleanup_result {
|
||||
std::result::Result::Ok(count) if count <= 1 => {},
|
||||
std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_variant_promotion_archive_cleanup_count")),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_promotion_archive_cleanup")),
|
||||
}
|
||||
let slot_text = incoming.slot().to_string();
|
||||
let block_time = match incoming.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_variant_promotion_block_time")),
|
||||
},
|
||||
std::option::Option::None => std::option::Option::None,
|
||||
};
|
||||
let format_version = i64::from(incoming.payload().format_version());
|
||||
let content_hash = incoming.payload().content_hash();
|
||||
let content_hash_bytes: &[u8] = content_hash.as_bytes();
|
||||
let payload_bytes = incoming.payload().bytes();
|
||||
let projection_result = sql_transaction
|
||||
.execute(
|
||||
UPDATE_CANONICAL_TRANSACTION_PROJECTION_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&slot_text.as_str(),
|
||||
&block_time,
|
||||
&incoming.payload().format_id().as_str(),
|
||||
&format_version,
|
||||
&content_hash_bytes,
|
||||
&payload_bytes,
|
||||
],
|
||||
)
|
||||
.await;
|
||||
match projection_result {
|
||||
std::result::Result::Ok(1) => {},
|
||||
std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_variant_promotion_projection_count")),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_promotion_projection")),
|
||||
}
|
||||
let expected_variant_id_text = expected_canonical_variant_id.to_string();
|
||||
let expected_revision_text = current_revision.to_string();
|
||||
let promoted_variant_id_text = promoted_variant_id.to_string();
|
||||
let next_revision_text = next_revision.to_string();
|
||||
let selector_update_result = sql_transaction
|
||||
.execute(
|
||||
UPDATE_TRANSACTION_VARIANT_SELECTOR_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&expected_variant_id_text.as_str(),
|
||||
&expected_revision_text.as_str(),
|
||||
&promoted_variant_id_text.as_str(),
|
||||
&next_revision_text.as_str(),
|
||||
&updated_at_unix_millis,
|
||||
],
|
||||
)
|
||||
.await;
|
||||
return match selector_update_result {
|
||||
std::result::Result::Ok(1) => std::result::Result::Ok(RawCanonicalPromotionTransition {
|
||||
from_revision: current_revision,
|
||||
from_variant_id: expected_canonical_variant_id,
|
||||
to_revision: next_revision,
|
||||
to_variant_id: promoted_variant_id,
|
||||
}),
|
||||
std::result::Result::Ok(0) => std::result::Result::Err(data_invalid("raw_variant_promotion_stale_selector")),
|
||||
std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_variant_promotion_selector_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_promotion_selector_update")),
|
||||
};
|
||||
}
|
||||
|
||||
async fn preserve_previous_canonical_variant(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
locked: &RawTransactionDbRow,
|
||||
canonical_variant_id: u64,
|
||||
) -> std::result::Result<(), crate::PostgresBackendError> {
|
||||
let state = match decode_retention_state(locked.retention_state.as_str()) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let preserved_payload = if state == ksp_store_api::RawRetentionState::Full {
|
||||
if locked.archive_payload.is_some() {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_full_shape"));
|
||||
}
|
||||
match locked.payload.as_ref() {
|
||||
std::option::Option::Some(value) => value.as_slice(),
|
||||
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_variant_promotion_full_payload")),
|
||||
}
|
||||
} else if state == ksp_store_api::RawRetentionState::Archived {
|
||||
if locked.payload.is_some() {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_archived_shape"));
|
||||
}
|
||||
match locked.archive_payload.as_ref() {
|
||||
std::option::Option::Some(value) => value.as_slice(),
|
||||
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_variant_promotion_archived_payload")),
|
||||
}
|
||||
} else {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_promotion_retention_state"));
|
||||
};
|
||||
let signature_bytes: &[u8] = locked.signature.as_slice();
|
||||
let canonical_variant_id_text = canonical_variant_id.to_string();
|
||||
let content_hash_bytes: &[u8] = locked.content_hash.as_slice();
|
||||
let update_result = sql_transaction
|
||||
.execute(
|
||||
UPDATE_TRANSACTION_VARIANT_TO_FULL_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&canonical_variant_id_text.as_str(),
|
||||
&locked.slot_text.as_str(),
|
||||
&locked.block_time_unix_millis,
|
||||
&locked.format_id.as_str(),
|
||||
&locked.format_version,
|
||||
&content_hash_bytes,
|
||||
&preserved_payload,
|
||||
],
|
||||
)
|
||||
.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_variant_promotion_preserve_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_promotion_preserve")),
|
||||
};
|
||||
}
|
||||
|
||||
fn next_transaction_canonical_revision(current: u64) -> std::result::Result<u64, crate::PostgresBackendError> {
|
||||
return match current.checked_add(1) {
|
||||
std::option::Option::Some(value) if value != 0 => std::result::Result::Ok(value),
|
||||
_ => std::result::Result::Err(data_invalid("raw_variant_canonical_revision_exhausted")),
|
||||
};
|
||||
}
|
||||
|
||||
async fn persist_transaction_variant_link(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
observation: &ksp_store_api::RawTransactionObservation,
|
||||
@@ -1449,9 +1655,8 @@ fn compare_existing_transaction(
|
||||
log_raw_transaction_compatible_truncated_log_messages(network, &stored, incoming);
|
||||
std::result::Result::Ok(ExistingTransactionMatch::ActiveCompatibleLessComplete)
|
||||
},
|
||||
ksp_store_api::RawTransactionVariantRelation::CompatibleMoreComplete
|
||||
| ksp_store_api::RawTransactionVariantRelation::Conflict
|
||||
| ksp_store_api::RawTransactionVariantRelation::Incomparable => {
|
||||
ksp_store_api::RawTransactionVariantRelation::CompatibleMoreComplete => std::result::Result::Ok(ExistingTransactionMatch::ActiveCompatibleMoreComplete),
|
||||
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"))
|
||||
},
|
||||
@@ -1705,6 +1910,21 @@ fn log_raw_transaction_compatible_truncated_log_messages(
|
||||
return;
|
||||
}
|
||||
|
||||
fn log_raw_transaction_canonical_promotion_committed(network: &ksp_store_api::RawNetworkId, slot: u64, transition: RawCanonicalPromotionTransition) {
|
||||
ksp_logging_lib::info!(
|
||||
target: crate::TRACING_TARGET,
|
||||
domain = "store.raw_transaction.canonical_promotion",
|
||||
network = network.as_str(),
|
||||
slot = slot,
|
||||
from_variant_id = transition.from_variant_id,
|
||||
to_variant_id = transition.to_variant_id,
|
||||
from_revision = transition.from_revision,
|
||||
to_revision = transition.to_revision,
|
||||
"PostgreSQL Store committed a proved RAW transaction canonical promotion"
|
||||
);
|
||||
return;
|
||||
}
|
||||
|
||||
fn log_raw_transaction_content_conflict(
|
||||
network: &ksp_store_api::RawNetworkId,
|
||||
stored: &ksp_store_api::RawTransaction,
|
||||
|
||||
Reference in New Issue
Block a user