0.3.16-pre.004
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/lib.rs
|
||||
// version: 27
|
||||
// version: 28
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
#![forbid(unsafe_code)]
|
||||
@@ -38,6 +38,8 @@
|
||||
//! families without changing the physical schema or read-by-key contracts.
|
||||
//! `0.3.16-pre.003` layers additive V003 transaction-variant resources over the
|
||||
//! frozen V000-V002 registry without changing any prior migration resource bytes.
|
||||
//! `0.3.16-pre.004` wires acquisition writes into the V003 variant ledger with
|
||||
//! lazy canonical bootstrap, exact native-variant reuse and exact observation linkage.
|
||||
//!
|
||||
//! This crate depends on `ksp-store-api` and never on `ksp-store-lib`. The
|
||||
//! common facade consumes only this crate's narrow backend bridge and never
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/src/raw_transaction.rs
|
||||
// version: 10
|
||||
// version: 11
|
||||
|
||||
pub(crate) mod cursor;
|
||||
|
||||
@@ -13,6 +13,9 @@ const GET_TRANSACTION_SQL: &str = "SELECT transaction_row.signature, transaction
|
||||
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 INSERT_TRANSACTION_VARIANT_LINK_SQL: &str = "INSERT INTO ksp_raw_transaction_observation_variants (observation_key, transaction_signature, variant_id, linked_at_unix_millis) VALUES ($1, $2, $3::TEXT::NUMERIC, $4) RETURNING observation_key";
|
||||
const INSERT_TRANSACTION_VARIANT_SELECTOR_SQL: &str = "INSERT INTO ksp_raw_transaction_canonical_selectors (transaction_signature, canonical_variant_id, canonical_revision, updated_at_unix_millis) VALUES ($1, $2::TEXT::NUMERIC, 1, $3)";
|
||||
const INSERT_TRANSACTION_VARIANT_SQL: &str = "INSERT INTO ksp_raw_transaction_variants (transaction_signature, variant_id, origin_kind, slot, block_time_unix_millis, format_id, format_version, content_hash, payload, retention_state, created_at_unix_millis) VALUES ($1, $2::TEXT::NUMERIC, 'native', $3::TEXT::NUMERIC, $4, $5, $6, $7, $8, $9, $10)";
|
||||
const INSPECT_OBSERVATIONS_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_transaction_observations WHERE ($1::BYTEA IS NULL OR transaction_signature = $1)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_transaction_observations) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.observation_key IS NOT NULL AS page_present, page.observation_key, page.transaction_signature, page.provider, page.protocol, page.acquisition_method, page.origin, page.received_at_unix_millis, page.capture_session_id, page.commitment, page.endpoint_id, page.filter_id, page.observed_at_unix_millis, page.source_payload_hash, page.source_payload_size_bytes FROM counts LEFT JOIN LATERAL (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 ($1::BYTEA IS NULL OR transaction_signature = $1) ORDER BY received_at_unix_millis ASC, observation_key ASC LIMIT $2 OFFSET $3) AS page ON TRUE";
|
||||
const INSPECT_OBSERVATIONS_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_transaction_observations WHERE ($1::BYTEA IS NULL OR transaction_signature = $1)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_transaction_observations) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.observation_key IS NOT NULL AS page_present, page.observation_key, page.transaction_signature, page.provider, page.protocol, page.acquisition_method, page.origin, page.received_at_unix_millis, page.capture_session_id, page.commitment, page.endpoint_id, page.filter_id, page.observed_at_unix_millis, page.source_payload_hash, page.source_payload_size_bytes FROM counts LEFT JOIN LATERAL (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 ($1::BYTEA IS NULL OR transaction_signature = $1) ORDER BY received_at_unix_millis DESC, observation_key DESC LIMIT $2 OFFSET $3) AS page ON TRUE";
|
||||
const INSPECT_TRANSACTIONS_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_transactions WHERE ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::TEXT IS NULL AND $2::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_transactions) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.signature IS NOT NULL AS page_present, page.signature, page.slot_text, page.block_time_unix_millis, page.format_id, page.format_version, page.content_hash, page.retention_state, page.payload_size_bytes, page.hot_payload_present, page.archive_payload_present FROM counts LEFT JOIN LATERAL (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.retention_state, CASE WHEN transaction_row.retention_state = 'full' THEN OCTET_LENGTH(transaction_row.payload)::BIGINT WHEN transaction_row.retention_state = 'archived' THEN OCTET_LENGTH(archive_row.payload)::BIGINT ELSE NULL END AS payload_size_bytes, transaction_row.payload IS NOT NULL AS hot_payload_present, archive_row.payload IS NOT NULL AS archive_payload_present 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 ($1::TEXT IS NULL OR transaction_row.slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR transaction_row.slot <= $2::TEXT::NUMERIC) ORDER BY transaction_row.slot ASC, transaction_row.signature ASC LIMIT $3 OFFSET $4) AS page ON TRUE";
|
||||
@@ -23,7 +26,6 @@ const LOCK_OBSERVATION_SQL: &str = "SELECT observation_key, transaction_signatur
|
||||
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 RAW_TRANSACTION_CONTENT_CONFLICT_META_FIELDS: [&str; 15] = [
|
||||
"err",
|
||||
"status",
|
||||
@@ -44,9 +46,17 @@ const RAW_TRANSACTION_CONTENT_CONFLICT_META_FIELDS: [&str; 15] = [
|
||||
const RAW_TRANSACTION_CONTENT_CONFLICT_PAYLOAD_FIELDS: [&str; 4] = ["transaction", "meta", "version", "transactionIndex"];
|
||||
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_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";
|
||||
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_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";
|
||||
|
||||
struct RawLogMessagesContentConflictDiagnostic {
|
||||
available: bool,
|
||||
@@ -157,6 +167,7 @@ struct RawRetentionDbRow {
|
||||
retention_state: std::string::String,
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
struct RawTransactionDbRow {
|
||||
archive_payload: std::option::Option<std::vec::Vec<u8>>,
|
||||
block_time_unix_millis: std::option::Option<i64>,
|
||||
@@ -656,6 +667,10 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
if let std::result::Result::Err(error) = input_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let linked_at = match raw_variant_timestamp(observation.provenance().received_at(), "raw_acquisition_variant_timestamp") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let client_result = pool.get().await;
|
||||
let mut client = match client_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
@@ -671,15 +686,20 @@ 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 = if inserted {
|
||||
ksp_store_api::RawEntityWriteOutcome::Inserted
|
||||
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 canonical_variant_result = ensure_v003_transaction_variant_bootstrap(&sql_transaction, &locked, linked_at).await;
|
||||
let canonical_variant_id = match canonical_variant_result {
|
||||
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)
|
||||
} 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,
|
||||
@@ -691,14 +711,26 @@ pub(crate) async fn persist_raw_transaction_acquisition(
|
||||
},
|
||||
};
|
||||
match comparison {
|
||||
ExistingTransactionMatch::Active | ExistingTransactionMatch::ActiveIncomingTruncatedLogs => ksp_store_api::RawEntityWriteOutcome::AlreadyPresent,
|
||||
ExistingTransactionMatch::Active => (ksp_store_api::RawEntityWriteOutcome::AlreadyPresent, canonical_variant_id),
|
||||
ExistingTransactionMatch::ActiveIncomingTruncatedLogs => {
|
||||
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)
|
||||
},
|
||||
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
|
||||
let variant_result = rehydrate_canonical_transaction_variant(&sql_transaction, &raw_transaction, canonical_variant_id).await;
|
||||
if let std::result::Result::Err(error) = variant_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
(ksp_store_api::RawEntityWriteOutcome::Rehydrated, canonical_variant_id)
|
||||
} else {
|
||||
let commit_result = sql_transaction.commit().await;
|
||||
if commit_result.is_err() {
|
||||
@@ -717,6 +749,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 link_result = persist_transaction_variant_link(&sql_transaction, &observation, observed_variant_id, observation_outcome, linked_at).await;
|
||||
if let std::result::Result::Err(error) = link_result {
|
||||
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"));
|
||||
@@ -734,6 +770,10 @@ pub(crate) async fn record_raw_transaction_observation(
|
||||
if let std::result::Result::Err(error) = network_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let linked_at = match raw_variant_timestamp(observation.provenance().received_at(), "raw_observation_variant_timestamp") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let client_result = pool.get().await;
|
||||
let mut client = match client_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
@@ -744,21 +784,15 @@ pub(crate) async fn record_raw_transaction_observation(
|
||||
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 {
|
||||
let locked_result = load_locked_transaction_row(&sql_transaction, observation.transaction()).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(reference_not_found("raw_observation_transaction")),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_observation_lock_transaction")),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
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")),
|
||||
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),
|
||||
};
|
||||
if state == ksp_store_api::RawRetentionState::Purged {
|
||||
let commit_result = sql_transaction.commit().await;
|
||||
@@ -767,11 +801,20 @@ pub(crate) async fn record_raw_transaction_observation(
|
||||
}
|
||||
return std::result::Result::Ok(ksp_store_api::RawObservationWriteOutcome::NotRecorded);
|
||||
}
|
||||
let canonical_variant_result = ensure_v003_transaction_variant_bootstrap(&sql_transaction, &locked, linked_at).await;
|
||||
let canonical_variant_id = match canonical_variant_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
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 link_result = persist_transaction_variant_link(&sql_transaction, &observation, canonical_variant_id, outcome, linked_at).await;
|
||||
if let std::result::Result::Err(error) = link_result {
|
||||
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"));
|
||||
@@ -1024,6 +1067,261 @@ fn validate_retained_payload_bytes(payload: &[u8], phase: &'static str) -> std::
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
async fn ensure_v003_transaction_variant_bootstrap(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
row: &RawTransactionDbRow,
|
||||
created_at_unix_millis: i64,
|
||||
) -> std::result::Result<u64, crate::PostgresBackendError> {
|
||||
let signature_bytes: &[u8] = row.signature.as_slice();
|
||||
let selector_result = sql_transaction.query_opt(SELECT_CANONICAL_VARIANT_ID_SQL, &[&signature_bytes]).await;
|
||||
match selector_result {
|
||||
std::result::Result::Ok(std::option::Option::Some(selector_row)) => {
|
||||
let 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_selector_id")),
|
||||
};
|
||||
return decode_u64_decimal(variant_id_text.as_str(), "raw_variant_selector_id");
|
||||
},
|
||||
std::result::Result::Ok(std::option::Option::None) => {},
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_selector_query")),
|
||||
}
|
||||
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 variant_payload: std::option::Option<&[u8]> = match retention_state {
|
||||
ksp_store_api::RawRetentionState::Full => {
|
||||
if row.archive_payload.is_some() {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_bootstrap_full_archive"));
|
||||
}
|
||||
match row.payload.as_ref() {
|
||||
std::option::Option::Some(value) => std::option::Option::Some(value.as_slice()),
|
||||
std::option::Option::None => return std::result::Result::Err(data_invalid("raw_variant_bootstrap_full_payload")),
|
||||
}
|
||||
},
|
||||
ksp_store_api::RawRetentionState::Archived => {
|
||||
if row.payload.is_some() || row.archive_payload.is_none() {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_bootstrap_archived_shape"));
|
||||
}
|
||||
std::option::Option::None
|
||||
},
|
||||
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_variant_bootstrap_purged_shape"));
|
||||
}
|
||||
std::option::Option::None
|
||||
},
|
||||
_ => return std::result::Result::Err(data_invalid("raw_variant_bootstrap_retention_state")),
|
||||
};
|
||||
let variant_id_text = "1";
|
||||
let content_hash_bytes: &[u8] = row.content_hash.as_slice();
|
||||
let insert_result = sql_transaction
|
||||
.execute(
|
||||
INSERT_TRANSACTION_VARIANT_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&variant_id_text,
|
||||
&row.slot_text.as_str(),
|
||||
&row.block_time_unix_millis,
|
||||
&row.format_id.as_str(),
|
||||
&row.format_version,
|
||||
&content_hash_bytes,
|
||||
&variant_payload,
|
||||
&row.retention_state.as_str(),
|
||||
&created_at_unix_millis,
|
||||
],
|
||||
)
|
||||
.await;
|
||||
match insert_result {
|
||||
std::result::Result::Ok(1) => {},
|
||||
std::result::Result::Ok(_) => return std::result::Result::Err(data_invalid("raw_variant_bootstrap_insert_count")),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_bootstrap_insert")),
|
||||
}
|
||||
let selector_insert_result =
|
||||
sql_transaction.execute(INSERT_TRANSACTION_VARIANT_SELECTOR_SQL, &[&signature_bytes, &variant_id_text, &created_at_unix_millis]).await;
|
||||
return match selector_insert_result {
|
||||
std::result::Result::Ok(1) => std::result::Result::Ok(1),
|
||||
std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_variant_selector_insert_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_selector_insert")),
|
||||
};
|
||||
}
|
||||
|
||||
async fn persist_or_reuse_native_transaction_variant(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
raw_transaction: &ksp_store_api::RawTransaction,
|
||||
created_at_unix_millis: i64,
|
||||
) -> std::result::Result<u64, 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_variant_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 exact_result = sql_transaction
|
||||
.query_opt(
|
||||
SELECT_EXACT_NATIVE_VARIANT_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&slot_text.as_str(),
|
||||
&block_time,
|
||||
&raw_transaction.payload().format_id().as_str(),
|
||||
&format_version,
|
||||
&content_hash_bytes,
|
||||
&payload_bytes,
|
||||
],
|
||||
)
|
||||
.await;
|
||||
match exact_result {
|
||||
std::result::Result::Ok(std::option::Option::Some(row)) => {
|
||||
let variant_id_text = match row.try_get::<_, std::string::String>("variant_id_text") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_exact_id")),
|
||||
};
|
||||
return decode_u64_decimal(variant_id_text.as_str(), "raw_variant_exact_id");
|
||||
},
|
||||
std::result::Result::Ok(std::option::Option::None) => {},
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_exact_query")),
|
||||
}
|
||||
let max_result = sql_transaction.query_one(SELECT_MAX_VARIANT_ID_SQL, &[&signature_bytes]).await;
|
||||
let max_row = match max_result {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_max_query")),
|
||||
};
|
||||
let max_variant_id_text = match max_row.try_get::<_, std::string::String>("variant_id_text") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_max_id")),
|
||||
};
|
||||
let variant_id = match next_transaction_variant_id(max_variant_id_text.as_str()) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let variant_id_text = variant_id.to_string();
|
||||
let variant_payload = std::option::Option::Some(payload_bytes);
|
||||
let insert_result = sql_transaction
|
||||
.execute(
|
||||
INSERT_TRANSACTION_VARIANT_SQL,
|
||||
&[
|
||||
&signature_bytes,
|
||||
&variant_id_text.as_str(),
|
||||
&slot_text.as_str(),
|
||||
&block_time,
|
||||
&raw_transaction.payload().format_id().as_str(),
|
||||
&format_version,
|
||||
&content_hash_bytes,
|
||||
&variant_payload,
|
||||
&"full",
|
||||
&created_at_unix_millis,
|
||||
],
|
||||
)
|
||||
.await;
|
||||
return match insert_result {
|
||||
std::result::Result::Ok(1) => std::result::Result::Ok(variant_id),
|
||||
std::result::Result::Ok(_) => std::result::Result::Err(data_invalid("raw_variant_insert_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_insert")),
|
||||
};
|
||||
}
|
||||
|
||||
async fn persist_transaction_variant_link(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
observation: &ksp_store_api::RawTransactionObservation,
|
||||
variant_id: u64,
|
||||
observation_outcome: ksp_store_api::RawObservationWriteOutcome,
|
||||
linked_at_unix_millis: i64,
|
||||
) -> std::result::Result<(), crate::PostgresBackendError> {
|
||||
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 variant_id_text = variant_id.to_string();
|
||||
if observation_outcome == ksp_store_api::RawObservationWriteOutcome::Inserted {
|
||||
let insert_result = sql_transaction
|
||||
.query_opt(INSERT_TRANSACTION_VARIANT_LINK_SQL, &[&observation_key_bytes, &signature_bytes, &variant_id_text.as_str(), &linked_at_unix_millis])
|
||||
.await;
|
||||
return match insert_result {
|
||||
std::result::Result::Ok(std::option::Option::Some(_)) => std::result::Result::Ok(()),
|
||||
std::result::Result::Ok(std::option::Option::None) => std::result::Result::Err(data_invalid("raw_variant_link_insert_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_link_insert")),
|
||||
};
|
||||
}
|
||||
if observation_outcome != ksp_store_api::RawObservationWriteOutcome::AlreadyPresent {
|
||||
return std::result::Result::Err(data_invalid("raw_variant_link_observation_outcome"));
|
||||
}
|
||||
let existing_result = sql_transaction.query_opt(SELECT_TRANSACTION_VARIANT_LINK_SQL, &[&observation_key_bytes]).await;
|
||||
let existing = 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::Ok(()),
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(write_failed("raw_variant_link_query")),
|
||||
};
|
||||
let existing_signature = match existing.try_get::<_, std::vec::Vec<u8>>("transaction_signature") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_link_signature")),
|
||||
};
|
||||
let existing_variant_id_text = match existing.try_get::<_, std::string::String>("variant_id_text") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_variant_link_id")),
|
||||
};
|
||||
let existing_variant_id = match decode_u64_decimal(existing_variant_id_text.as_str(), "raw_variant_link_id") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
if existing_signature.as_slice() != signature_bytes || existing_variant_id != variant_id {
|
||||
return std::result::Result::Err(conflict("raw_variant_link_conflict"));
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
async fn rehydrate_canonical_transaction_variant(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
raw_transaction: &ksp_store_api::RawTransaction,
|
||||
canonical_variant_id: u64,
|
||||
) -> std::result::Result<(), crate::PostgresBackendError> {
|
||||
let signature = raw_transaction.reference().signature();
|
||||
let signature_bytes: &[u8] = signature.as_bytes();
|
||||
let variant_id_text = canonical_variant_id.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_variant_rehydrate_block_time")),
|
||||
},
|
||||
std::option::Option::None => std::option::Option::None,
|
||||
};
|
||||
let payload_bytes = raw_transaction.payload().bytes();
|
||||
let update_result = sql_transaction
|
||||
.execute(UPDATE_REHYDRATED_TRANSACTION_VARIANT_SQL, &[&signature_bytes, &variant_id_text.as_str(), &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_variant_rehydrate_count")),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(write_failed("raw_variant_rehydrate")),
|
||||
};
|
||||
}
|
||||
|
||||
fn next_transaction_variant_id(current_max: &str) -> std::result::Result<u64, crate::PostgresBackendError> {
|
||||
let current = match decode_u64_decimal(current_max, "raw_variant_max_id") {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
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_id_exhausted")),
|
||||
};
|
||||
}
|
||||
|
||||
fn raw_variant_timestamp(timestamp: ksp_store_api::RawTimestamp, phase: &'static str) -> std::result::Result<i64, crate::PostgresBackendError> {
|
||||
return match i64::try_from(timestamp.unix_millis()) {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(data_invalid(phase)),
|
||||
};
|
||||
}
|
||||
|
||||
async fn insert_canonical_transaction(
|
||||
sql_transaction: &deadpool_postgres::Transaction<'_>,
|
||||
raw_transaction: &ksp_store_api::RawTransaction,
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs
|
||||
// version: 4
|
||||
// version: 5
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -15,7 +15,11 @@
|
||||
const LIVE_BOOTSTRAP_SQL: &str = include_str!("../migrations/v000_bootstrap/tables/001_ksp_store_schema_migrations.sql");
|
||||
const LIVE_BROKEN_CHECKSUM_A: &str = "0000000000000000000000000000000000000000000000000000000000000000";
|
||||
const LIVE_BROKEN_CHECKSUM_B: &str = "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff";
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_observations;
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_conflicts;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observation_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_canonical_selectors;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observations;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_archive_payloads;
|
||||
DROP TABLE IF EXISTS ksp_raw_transactions;
|
||||
DROP TABLE IF EXISTS ksp_store_identity;
|
||||
@@ -28,7 +32,11 @@ const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
|
||||
'ksp_store_identity',
|
||||
'ksp_raw_transactions',
|
||||
'ksp_raw_transaction_observations',
|
||||
'ksp_raw_transaction_archive_payloads'
|
||||
'ksp_raw_transaction_archive_payloads',
|
||||
'ksp_raw_transaction_variants',
|
||||
'ksp_raw_transaction_canonical_selectors',
|
||||
'ksp_raw_transaction_observation_variants',
|
||||
'ksp_raw_transaction_conflicts'
|
||||
)
|
||||
AND table_type = 'BASE TABLE'
|
||||
)"#;
|
||||
@@ -153,7 +161,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
|
||||
}
|
||||
*owns_schema = true;
|
||||
let initial_health = initial.health().await;
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(2) || initial_health.pending_migration_count() != 0 {
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(3) || initial_health.pending_migration_count() != 0 {
|
||||
return std::result::Result::Err(LiveFailure::new("initial_health"));
|
||||
}
|
||||
let initial_close = close_backend(initial).await;
|
||||
@@ -167,7 +175,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
|
||||
};
|
||||
let idempotent_health = idempotent.health().await;
|
||||
if !idempotent_health.is_ready()
|
||||
|| idempotent_health.migration_version() != std::option::Option::Some(2)
|
||||
|| idempotent_health.migration_version() != std::option::Option::Some(3)
|
||||
|| idempotent_health.pending_migration_count() != 0
|
||||
{
|
||||
return std::result::Result::Err(LiveFailure::new("idempotent_health"));
|
||||
@@ -250,7 +258,7 @@ async fn run_foundation_scenario(admin: &mut tokio_postgres::Client, uri: &str,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let final_health = final_backend.health().await;
|
||||
if !final_health.is_ready() || final_health.migration_version() != std::option::Option::Some(2) || final_health.pending_migration_count() != 0 {
|
||||
if !final_health.is_ready() || final_health.migration_version() != std::option::Option::Some(3) || final_health.pending_migration_count() != 0 {
|
||||
return std::result::Result::Err(LiveFailure::new("final_health"));
|
||||
}
|
||||
return close_backend(final_backend).await;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -20,7 +20,11 @@ const LIVE_ACCOUNT_STATE_EXISTS_SQL: &str =
|
||||
"SELECT EXISTS (SELECT 1 FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3)";
|
||||
const LIVE_CANCEL_WAIT: std::time::Duration = std::time::Duration::from_millis(300);
|
||||
const LIVE_LOCK_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key FROM ksp_raw_account_observations WHERE observation_key = $1 FOR UPDATE";
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_account_observations;
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_conflicts;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observation_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_canonical_selectors;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_account_observations;
|
||||
DROP TABLE IF EXISTS ksp_raw_account_states;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observations;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_archive_payloads;
|
||||
@@ -36,6 +40,10 @@ const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
|
||||
'ksp_raw_transactions',
|
||||
'ksp_raw_transaction_observations',
|
||||
'ksp_raw_transaction_archive_payloads',
|
||||
'ksp_raw_transaction_variants',
|
||||
'ksp_raw_transaction_canonical_selectors',
|
||||
'ksp_raw_transaction_observation_variants',
|
||||
'ksp_raw_transaction_conflicts',
|
||||
'ksp_raw_account_states',
|
||||
'ksp_raw_account_observations'
|
||||
)
|
||||
@@ -152,7 +160,7 @@ async fn run_raw_account_scenario(admin: &mut tokio_postgres::Client, uri: &str,
|
||||
};
|
||||
*owns_schema = true;
|
||||
let initial_health = initial.health().await;
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(2) || initial_health.pending_migration_count() != 0 {
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(3) || initial_health.pending_migration_count() != 0 {
|
||||
return std::result::Result::Err(LiveFailure::new("initial_health"));
|
||||
}
|
||||
let wrong_network_result = open_backend_result(uri, "testnet", true, true).await;
|
||||
@@ -211,7 +219,7 @@ async fn run_raw_account_scenario(admin: &mut tokio_postgres::Client, uri: &str,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let reopened_health = reopened.health().await;
|
||||
if !reopened_health.is_ready() || reopened_health.migration_version() != std::option::Option::Some(2) {
|
||||
if !reopened_health.is_ready() || reopened_health.migration_version() != std::option::Option::Some(3) {
|
||||
return std::result::Result::Err(LiveFailure::new("reopen_health"));
|
||||
}
|
||||
let reference = match account_reference(10, 100, 10) {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs
|
||||
// version: 7
|
||||
// version: 9
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
@@ -18,7 +18,11 @@ const LIVE_INDEX_EXISTS_SQL: &str = r#"SELECT EXISTS (
|
||||
AND indexname = 'ix_ksp_raw_transactions_slot_signature'
|
||||
)"#;
|
||||
const LIVE_LOCK_OBSERVATION_SQL: &str = "SELECT observation_key FROM ksp_raw_transaction_observations WHERE observation_key = $1 FOR UPDATE";
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_account_observations;
|
||||
const LIVE_MANAGED_SCHEMA_DROP_SQL: &str = r#"DROP TABLE IF EXISTS ksp_raw_transaction_conflicts;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observation_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_canonical_selectors;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_variants;
|
||||
DROP TABLE IF EXISTS ksp_raw_account_observations;
|
||||
DROP TABLE IF EXISTS ksp_raw_account_states;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_observations;
|
||||
DROP TABLE IF EXISTS ksp_raw_transaction_archive_payloads;
|
||||
@@ -34,6 +38,10 @@ const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
|
||||
'ksp_raw_transactions',
|
||||
'ksp_raw_transaction_observations',
|
||||
'ksp_raw_transaction_archive_payloads',
|
||||
'ksp_raw_transaction_variants',
|
||||
'ksp_raw_transaction_canonical_selectors',
|
||||
'ksp_raw_transaction_observation_variants',
|
||||
'ksp_raw_transaction_conflicts',
|
||||
'ksp_raw_account_states',
|
||||
'ksp_raw_account_observations'
|
||||
)
|
||||
@@ -42,6 +50,7 @@ const LIVE_MANAGED_SCHEMA_EXISTS_SQL: &str = r#"SELECT EXISTS (
|
||||
const LIVE_MAX_URI_BYTES: usize = 4_096;
|
||||
const LIVE_OBSERVATION_EXISTS_SQL: &str = "SELECT EXISTS (SELECT 1 FROM ksp_raw_transaction_observations WHERE observation_key = $1)";
|
||||
const LIVE_TRANSACTION_EXISTS_SQL: &str = "SELECT EXISTS (SELECT 1 FROM ksp_raw_transactions WHERE signature = $1)";
|
||||
const LIVE_VARIANT_IDEMPOTENCE_SQL: &str = "SELECT (SELECT COUNT(*) = 1 FROM ksp_raw_transaction_variants WHERE transaction_signature = $1) AND (SELECT COUNT(*) = 1 FROM ksp_raw_transaction_canonical_selectors WHERE transaction_signature = $1) AND (SELECT COUNT(*) = 1 FROM ksp_raw_transaction_observation_variants WHERE transaction_signature = $1) AND EXISTS (SELECT 1 FROM ksp_raw_transaction_observation_variants AS mapping INNER JOIN ksp_raw_transaction_canonical_selectors AS selector ON selector.transaction_signature = mapping.transaction_signature AND selector.canonical_variant_id = mapping.variant_id WHERE mapping.observation_key = $2 AND mapping.transaction_signature = $1)";
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
struct LiveFailure {
|
||||
@@ -157,7 +166,7 @@ async fn run_raw_transaction_scenario(admin: &mut tokio_postgres::Client, uri: &
|
||||
};
|
||||
*owns_schema = true;
|
||||
let initial_health = initial.health().await;
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(2) || initial_health.pending_migration_count() != 0 {
|
||||
if !initial_health.is_ready() || initial_health.migration_version() != std::option::Option::Some(3) || initial_health.pending_migration_count() != 0 {
|
||||
return std::result::Result::Err(LiveFailure::new("initial_health"));
|
||||
}
|
||||
let wrong_network_result = open_backend_result(uri, "testnet", true, true).await;
|
||||
@@ -197,7 +206,7 @@ async fn run_raw_transaction_scenario(admin: &mut tokio_postgres::Client, uri: &
|
||||
if let std::result::Result::Err(error) = backend_close {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
let identical_result = prove_concurrent_identical_insert(uri).await;
|
||||
let identical_result = prove_concurrent_identical_insert(admin, uri).await;
|
||||
if let std::result::Result::Err(error) = identical_result {
|
||||
return std::result::Result::Err(error);
|
||||
}
|
||||
@@ -219,7 +228,7 @@ async fn run_raw_transaction_scenario(admin: &mut tokio_postgres::Client, uri: &
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let reopened_health = reopened.health().await;
|
||||
if !reopened_health.is_ready() || reopened_health.migration_version() != std::option::Option::Some(2) {
|
||||
if !reopened_health.is_ready() || reopened_health.migration_version() != std::option::Option::Some(3) {
|
||||
return std::result::Result::Err(LiveFailure::new("reopen_health"));
|
||||
}
|
||||
let reference_result = raw_reference(10);
|
||||
@@ -361,7 +370,7 @@ async fn prove_additional_observation_and_atomic_rollback(backend: &ksp_store_po
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
async fn prove_concurrent_identical_insert(uri: &str) -> std::result::Result<(), LiveFailure> {
|
||||
async fn prove_concurrent_identical_insert(admin: &tokio_postgres::Client, uri: &str) -> std::result::Result<(), LiveFailure> {
|
||||
let first = tokio::spawn(persist_once(uri.to_owned(), 20, 2_000, 20, 20));
|
||||
let second = tokio::spawn(persist_once(uri.to_owned(), 20, 2_000, 20, 20));
|
||||
let first_result = match joined_persist(first.await, "concurrent_identical_first") {
|
||||
@@ -378,6 +387,17 @@ async fn prove_concurrent_identical_insert(uri: &str) -> std::result::Result<(),
|
||||
if inserted != 1 || already != 1 {
|
||||
return std::result::Result::Err(LiveFailure::new("concurrent_identical_outcome"));
|
||||
}
|
||||
let reference = match raw_reference(20) {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
};
|
||||
let key = ksp_store_api::RawObservationKey::new([20; 32]);
|
||||
let projection_result = variant_projection_is_exact(admin, &reference, &key).await;
|
||||
match projection_result {
|
||||
std::result::Result::Ok(true) => {},
|
||||
std::result::Result::Ok(false) => return std::result::Result::Err(LiveFailure::new("concurrent_identical_variant_projection")),
|
||||
std::result::Result::Err(error) => return std::result::Result::Err(error),
|
||||
}
|
||||
return std::result::Result::Ok(());
|
||||
}
|
||||
|
||||
@@ -1256,3 +1276,21 @@ async fn observation_exists(client: &tokio_postgres::Client, key: &ksp_store_api
|
||||
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("observation_exists_decode")),
|
||||
};
|
||||
}
|
||||
|
||||
async fn variant_projection_is_exact(
|
||||
client: &tokio_postgres::Client,
|
||||
reference: &ksp_store_api::RawTransactionReference,
|
||||
key: &ksp_store_api::RawObservationKey,
|
||||
) -> std::result::Result<bool, LiveFailure> {
|
||||
let signature = reference.signature();
|
||||
let signature_bytes: &[u8] = signature.as_bytes();
|
||||
let key_bytes: &[u8] = key.as_bytes();
|
||||
let row = match client.query_one(LIVE_VARIANT_IDEMPOTENCE_SQL, &[&signature_bytes, &key_bytes]).await {
|
||||
std::result::Result::Ok(value) => value,
|
||||
std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("variant_projection_probe")),
|
||||
};
|
||||
return match row.try_get::<usize, bool>(0) {
|
||||
std::result::Result::Ok(value) => std::result::Result::Ok(value),
|
||||
std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("variant_projection_decode")),
|
||||
};
|
||||
}
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
// file: crates/ksp-store-postgres-lib/tests/v003_variant_persistence.rs
|
||||
// version: 1
|
||||
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
#![forbid(unsafe_code)]
|
||||
|
||||
//! Static completeness canaries for the first V003 variant-aware PostgreSQL write slice.
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_acquisition_orders_identity_lock_bootstrap_observation_and_variant_link_atomically() {
|
||||
let source = include_str!("../src/raw_transaction.rs");
|
||||
let start = source.find("pub(crate) async fn persist_raw_transaction_acquisition(");
|
||||
let end = source.find("pub(crate) async fn record_raw_transaction_observation(");
|
||||
let (start, end) = match (start, end) {
|
||||
(std::option::Option::Some(start), std::option::Option::Some(end)) if start < end => (start, end),
|
||||
_ => panic!("pre.004 acquisition function markers are missing"),
|
||||
};
|
||||
let body = &source[start..end];
|
||||
let lock = body.find("load_locked_transaction_row");
|
||||
let bootstrap = body.find("ensure_v003_transaction_variant_bootstrap");
|
||||
let compare = body.find("compare_existing_transaction");
|
||||
let observation = body.find("persist_observation_row");
|
||||
let link = body.find("persist_transaction_variant_link");
|
||||
assert!(matches!((lock, bootstrap, compare), (Some(lock), Some(bootstrap), Some(compare)) if lock < bootstrap && bootstrap < compare));
|
||||
assert!(matches!((observation, link), (Some(observation), Some(link)) if observation < link));
|
||||
assert!(body.contains("sql_transaction.commit().await"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_native_variant_reuse_is_exact_and_hash_is_only_a_prefilter() {
|
||||
let source = include_str!("../src/raw_transaction.rs");
|
||||
let exact_sql = source.find("const SELECT_EXACT_NATIVE_VARIANT_SQL");
|
||||
let exact_helper = source.find("async fn persist_or_reuse_native_transaction_variant(");
|
||||
assert!(exact_sql.is_some());
|
||||
assert!(exact_helper.is_some());
|
||||
assert!(source.contains("content_hash = $6 AND payload = $7"));
|
||||
assert!(source.contains("block_time_unix_millis IS NOT DISTINCT FROM $3"));
|
||||
assert!(source.contains("origin_kind = 'native'"));
|
||||
assert!(!source.contains("UNIQUE (transaction_signature, content_hash)"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_legacy_observation_without_mapping_is_not_fabricated() {
|
||||
let source = include_str!("../src/raw_transaction.rs");
|
||||
let start = source.find("async fn persist_transaction_variant_link(");
|
||||
let end = source.find("async fn rehydrate_canonical_transaction_variant(");
|
||||
let (start, end) = match (start, end) {
|
||||
(std::option::Option::Some(start), std::option::Option::Some(end)) if start < end => (start, end),
|
||||
_ => panic!("pre.004 observation-variant link markers are missing"),
|
||||
};
|
||||
let body = &source[start..end];
|
||||
assert!(body.contains("RawObservationWriteOutcome::Inserted"));
|
||||
assert!(body.contains("RawObservationWriteOutcome::AlreadyPresent"));
|
||||
assert!(body.contains("Ok(std::option::Option::None) => return std::result::Result::Ok(())"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_keeps_comparator_and_durable_conflict_scope_deferred() {
|
||||
let source = include_str!("../src/raw_transaction.rs");
|
||||
assert!(!source.contains("RawTransactionVariantRelation::"));
|
||||
assert!(!source.contains("INSERT INTO ksp_raw_transaction_conflicts"));
|
||||
assert!(!source.contains("UPDATE ksp_raw_transaction_canonical_selectors SET canonical_variant_id"));
|
||||
return;
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs
|
||||
// version: 10
|
||||
// version: 11
|
||||
|
||||
fn network() -> ksp_store_api::RawNetworkId {
|
||||
return match ksp_store_api::RawNetworkId::new("devnet") {
|
||||
@@ -817,3 +817,33 @@ fn pre_007_retained_payload_shape_rejects_empty_and_oversized_bytes() {
|
||||
assert_eq!(oversized.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::DataInvalid));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_variant_id_allocation_is_monotonic_and_bounded() {
|
||||
assert_eq!(super::next_transaction_variant_id("0"), std::result::Result::Ok(1));
|
||||
assert_eq!(super::next_transaction_variant_id("1"), std::result::Result::Ok(2));
|
||||
let exhausted = super::next_transaction_variant_id(u64::MAX.to_string().as_str()).err();
|
||||
assert_eq!(exhausted.as_ref().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::DataInvalid));
|
||||
assert_eq!(exhausted.as_ref().map(|value| return value.phase()), std::option::Option::Some("raw_variant_id_exhausted"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_variant_sql_uses_exact_content_and_identity_locking() {
|
||||
assert!(super::LOCK_TRANSACTION_SQL.contains("FOR UPDATE"));
|
||||
assert!(super::SELECT_EXACT_NATIVE_VARIANT_SQL.contains("content_hash = $6"));
|
||||
assert!(super::SELECT_EXACT_NATIVE_VARIANT_SQL.contains("payload = $7"));
|
||||
assert!(super::SELECT_EXACT_NATIVE_VARIANT_SQL.contains("origin_kind = 'native'"));
|
||||
assert!(super::SELECT_MAX_VARIANT_ID_SQL.contains("MAX(variant_id)"));
|
||||
assert!(!super::INSERT_TRANSACTION_VARIANT_SQL.contains("ON CONFLICT"));
|
||||
assert!(!super::INSERT_TRANSACTION_VARIANT_SQL.contains("UNIQUE"));
|
||||
return;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn v0_3_16_pre_004_variant_link_contract_keeps_legacy_observations_unfabricated() {
|
||||
assert!(super::INSERT_TRANSACTION_VARIANT_LINK_SQL.contains("ksp_raw_transaction_observation_variants"));
|
||||
assert!(super::SELECT_TRANSACTION_VARIANT_LINK_SQL.contains("WHERE observation_key = $1"));
|
||||
assert!(!super::INSERT_TRANSACTION_VARIANT_LINK_SQL.contains("ON CONFLICT"));
|
||||
return;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user