diff --git a/Cargo.toml b/Cargo.toml index 9dd4ac2..c94766b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ resolver = "3" members = ["crates/ksp-app-backfill-desk", "crates/ksp-app-config-desk", "crates/ksp-app-raw-transaction-ingest-desk", "crates/ksp-app-solprices-desk", "crates/ksp-app-store-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-interface-lib", "crates/ksp-job-api", "crates/ksp-job-backfill-lib", "crates/ksp-logging-lib", "crates/ksp-offchain-transport-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-program-api", "crates/ksp-raw-transaction-lib", "crates/ksp-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib", "crates/ksp-worker-api", "crates/ksp-worker-raw-transaction-ingest-lib"] [workspace.package] -version = "0.3.16-pre.3.fix.2" +version = "0.3.16-pre.4" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-store-postgres-lib/src/lib.rs b/crates/ksp-store-postgres-lib/src/lib.rs index 8e3f862..224a5b4 100644 --- a/crates/ksp-store-postgres-lib/src/lib.rs +++ b/crates/ksp-store-postgres-lib/src/lib.rs @@ -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 diff --git a/crates/ksp-store-postgres-lib/src/raw_transaction.rs b/crates/ksp-store-postgres-lib/src/raw_transaction.rs index b77adc2..6e2741a 100644 --- a/crates/ksp-store-postgres-lib/src/raw_transaction.rs +++ b/crates/ksp-store-postgres-lib/src/raw_transaction.rs @@ -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>, block_time_unix_millis: std::option::Option, @@ -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 { + 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 { + 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>("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 { + 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 { + 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, diff --git a/crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs b/crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs index c4362da..9c53da9 100644 --- a/crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs +++ b/crates/ksp-store-postgres-lib/tests/postgres_foundation_live.rs @@ -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; diff --git a/crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs b/crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs index fbac2ff..8aa862c 100644 --- a/crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs +++ b/crates/ksp-store-postgres-lib/tests/postgres_raw_account_live.rs @@ -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) { diff --git a/crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs b/crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs index 2a0439e..7b5c016 100644 --- a/crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs +++ b/crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs @@ -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 { + 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::(0) { + std::result::Result::Ok(value) => std::result::Result::Ok(value), + std::result::Result::Err(_) => std::result::Result::Err(LiveFailure::new("variant_projection_decode")), + }; +} diff --git a/crates/ksp-store-postgres-lib/tests/v003_variant_persistence.rs b/crates/ksp-store-postgres-lib/tests/v003_variant_persistence.rs new file mode 100644 index 0000000..2fc2f55 --- /dev/null +++ b/crates/ksp-store-postgres-lib/tests/v003_variant_persistence.rs @@ -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; +} diff --git a/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs b/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs index 08dd226..a4fc746 100644 --- a/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs +++ b/crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs @@ -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; +} diff --git a/deltas/0.3.16/pre.004.md b/deltas/0.3.16/pre.004.md new file mode 100644 index 0000000..6a1c0db --- /dev/null +++ b/deltas/0.3.16/pre.004.md @@ -0,0 +1,127 @@ + + + +# Delta `0.3.16-pre.004` — persistance PostgreSQL V003 des variantes RAW + +## Base requise + +```text +0.3.16-pre.003-fix.002 appliquée +workspace.package.version = 0.3.16-pre.3.fix.2 +``` + +Le gate opérateur de la base est confirmé propre : audits Rust/export/KSP/Markdown, `cargo check --workspace`, Clippy `-D warnings`, 82 tests unitaires `ksp-store-postgres-lib` et toutes les suites d'intégration ciblées passent. Les trois preuves PostgreSQL live restent opt-in. + +## Version + +Cette tranche modifie le runtime PostgreSQL : + +```text +workspace.package.version = 0.3.16-pre.4 +``` + +## Portée + +Cette tranche branche la persistance RAW existante sur les surfaces V003 créées par `pre.003`, sans introduire encore le comparateur backend-neutral de `pre.005`, la promotion canonique de `pre.006` ni le dossier de conflit durable de `pre.007`. + +Le backend PostgreSQL fournit désormais : + +- un bootstrap paresseux V003 sous le verrou transactionnel de l'identité V001 ; +- une variante bootstrap `variant_id = 1` représentant l'état canonique réellement disponible lorsqu'aucun selector V003 n'existe encore ; +- un selector canonique initial de revision `1` créé atomiquement avec cette variante ; +- la réutilisation d'une variante native uniquement après comparaison exacte de `slot`, `block_time`, format, version, `content_hash` et payload ; +- `content_hash` comme préfiltre de recherche, jamais comme preuve d'égalité ; +- une allocation monotone et bornée du prochain `variant_id` sous le verrou de l'identité ; +- le rattachement atomique de toute nouvelle observation persistée à la variante réellement reçue ; +- le rattachement des observations supplémentaires sans payload au selector canonique courant ; +- la conservation explicite des observations legacy déjà présentes mais sans mapping V003 : aucune association historique n'est fabriquée ; +- la réhydratation de la variante canonique bootstrap lorsqu'un tombstone legacy V001 est réhydraté avec preuve d'identité/hash existante. + +Le cas de compatibilité étroit déjà admis en `0.3.15`, `ActiveIncomingTruncatedLogs`, matérialise maintenant l'entrant tronqué comme variante native distincte et rattache son observation à cette variante sans modifier le canonique. + +## Atomicité et concurrence + +L'acquisition reste une seule transaction PostgreSQL : + +```text +BEGIN +insert/collision V001 +lock identité V001 FOR UPDATE +bootstrap selector/variant V003 si absent +comparaison canonique existante +réutilisation/insertion exacte de variante native si nécessaire +insertion/idempotence observation +rattachement observation -> variante +COMMIT +``` + +Le verrou de l'identité sérialise l'allocation de `variant_id` et empêche deux writers concurrents de créer deux variantes pour un même contenu exact. + +Une preuve PostgreSQL live existante est étendue : après deux acquisitions identiques concurrentes, elle exige exactement une variante, un selector canonique, un mapping d'observation et l'égalité entre la variante mappée et le selector. + +## Rétention + +Cette tranche ne met pas encore le ledger V003 sous une politique de rétention indépendante. + +Une variante native déjà matérialisée peut donc conserver ses bytes `full` même si la projection V001 canonique passe ensuite `archived` ou `purged`. Ce comportement est volontaire : la rétention, les pins, les purge guards et le rollback local des variantes sont réservés à `0.3.17-pre.003`. + +Pour une identité legacy bootstrapée alors que V001 est déjà `archived` ou `purged`, la variante bootstrap reflète uniquement l'état réellement disponible et ne prétend pas reconstruire des bytes perdus. + +## Hors scope maintenu + +Cette tranche n'ajoute pas : + +- `RawTransactionVariantRelation` dans le backend PostgreSQL ; +- dominance `CompatibleLessComplete` / `CompatibleMoreComplete` générique ; +- promotion ou remplacement du selector canonique ; +- écriture dans `ksp_raw_transaction_conflicts` ; +- dossier de conflit durable non terminal ; +- Store retry Worker ; +- inspection/résolution Store Desk ; +- rétention propre aux variantes. + +Aucune ressource SQL V003 et aucun fichier historique `migration.rs` / `schema.rs` ne sont modifiés. + +## Tests et canaris + +Les tests unitaires ajoutés verrouillent : + +- allocation `variant_id` monotone et bornée ; +- recherche exacte incluant le payload et ne reposant pas sur le hash seul ; +- absence de `ON CONFLICT` masquant une collision d'identité de variante ; +- absence de fabrication de mapping pour une observation legacy. + +Le canari d'intégration `v003_variant_persistence.rs` verrouille : + +- l'ordre lock -> bootstrap -> comparaison ; +- l'ordre observation -> mapping -> commit ; +- la comparaison exacte avant réutilisation ; +- l'absence de promotion canonique et de conflit durable avant leurs tranches dédiées. + +Les trois tests PostgreSQL live existants sont alignés sur `migration_version = 3` et nettoient désormais aussi les quatre tables V003. + +## Validation exécutée lors de la génération + +Les scripts Python réels du checkout complet sont exécutés sur l'état final du delta avant archivage. + +Le toolchain Rust n'est pas disponible dans l'environnement de génération ; les gates Cargo restent opérateur. + +## Validation opérateur demandée + +```bash +cargo fmt --all +cargo fmt --all -- --check + +python3 scripts/audit_rust_workspace_rules.py +python3 scripts/audit_markdown_tables.py README.md RULES.md ROADMAP.md CHANGELOG.md docs prompts crates deltas + +cargo check --workspace +cargo clippy --workspace --all-targets --all-features -- -D warnings +cargo test -p ksp-store-postgres-lib --all-targets --all-features +``` + +La preuve PostgreSQL réelle reste opt-in et peut être exécutée séparément sur une base dédiée lorsque souhaité. + +## Suite + +Après gate propre : `0.3.16-pre.005` — comparateur partagé backend-neutral `Exact | CompatibleLessComplete | CompatibleMoreComplete | Conflict | Incomparable`, avec canaris `logMessages` bidirectionnels et politique fail-closed sur les autres champs.