diff --git a/Cargo.toml b/Cargo.toml index 4b6a7b5..62656f9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 449 +# version: 450 [workspace] resolver = "3" members = ["crates/ksp-app-backfill-desk", "crates/ksp-app-config-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-store-api", "crates/ksp-store-lib", "crates/ksp-store-postgres-lib", "crates/ksp-wallet-lib"] [workspace.package] -version = "0.3.8-pre.3.fix.1" +version = "0.3.8-pre.4" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-store-lib/src/store.rs b/crates/ksp-store-lib/src/store.rs index 2d519c4..cf4d014 100644 --- a/crates/ksp-store-lib/src/store.rs +++ b/crates/ksp-store-lib/src/store.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-lib/src/store.rs -// version: 7 +// version: 8 /// Opaque common Store runtime facade. /// @@ -324,6 +324,36 @@ impl ksp_store_api::RawTransactionRead for Store { } } +impl ksp_store_api::RawTransactionInspectionRead for Store { + fn inspect_raw_transactions<'a>( + &'a self, + query: &'a ksp_store_api::RawTransactionInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + let network_check = validate_operation_network(&self.network, query.network(), self.backend_kind); + if let std::result::Result::Err(error) = network_check { + return std::boxed::Box::pin(async move { + return std::result::Result::Err(error); + }); + } + return std::boxed::Box::pin(async move { + #[cfg(feature = "postgres")] + { + return match &self.runtime { + StoreRuntime::Postgres(backend) => { + let result = backend.inspect_raw_transactions(query).await; + result.map_err(|error| return map_postgres_error(error, self.backend_kind, self.network.as_str())) + }, + }; + } + #[cfg(not(feature = "postgres"))] + { + let _ = query; + return std::result::Result::Err(unavailable_runtime_error(self.backend_kind)); + } + }); + } +} + impl ksp_store_api::RawTransactionWrite for Store { fn persist_raw_transaction_acquisition<'a>( &'a self, diff --git a/crates/ksp-store-lib/tests/dependency_boundary.rs b/crates/ksp-store-lib/tests/dependency_boundary.rs index 16c02a0..8f891a8 100644 --- a/crates/ksp-store-lib/tests/dependency_boundary.rs +++ b/crates/ksp-store-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-lib/tests/dependency_boundary.rs -// version: 8 +// version: 9 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -60,13 +60,14 @@ fn pre_005_facade_exposes_no_physical_postgres_types_or_environment_bypass() { } #[test] -fn pre_008_facade_dispatches_exact_ten_raw_capabilities_without_physical_leak() { +fn v0_3_8_pre_004_facade_dispatches_eleven_raw_capabilities_without_physical_leak() { let store = include_str!("../src/store.rs"); for required in [ "impl ksp_store_api::RawAccountObservationRead for Store", "impl ksp_store_api::RawAccountObservationWrite for Store", "impl ksp_store_api::RawAccountStateRead for Store", "impl ksp_store_api::RawAccountStateWrite for Store", + "impl ksp_store_api::RawTransactionInspectionRead for Store", "impl ksp_store_api::RawTransactionObservationRead for Store", "impl ksp_store_api::RawTransactionObservationWrite for Store", "impl ksp_store_api::RawTransactionRead for Store", @@ -79,7 +80,7 @@ fn pre_008_facade_dispatches_exact_ten_raw_capabilities_without_physical_leak() ] { assert!(store.contains(required), "missing pre.008 Store capability dispatch contract: {required}"); } - assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 10); + assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 11); for forbidden in ["tokio_postgres::", "deadpool_postgres::", "CREATE TABLE", "INSERT INTO", "UPDATE ksp_", "DELETE FROM"] { assert!(!store.contains(forbidden), "pre.008 facade leaked physical backend material: {forbidden}"); } diff --git a/crates/ksp-store-lib/tests/hardening_completeness.rs b/crates/ksp-store-lib/tests/hardening_completeness.rs index 458884a..9542211 100644 --- a/crates/ksp-store-lib/tests/hardening_completeness.rs +++ b/crates/ksp-store-lib/tests/hardening_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-lib/tests/hardening_completeness.rs -// version: 7 +// version: 8 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -292,13 +292,14 @@ fn pre_009_facade_production_sources_keep_config_env_physical_sql_and_backend_ha } #[test] -fn pre_010_facade_raw_capability_inventory_is_exactly_ten() { +fn v0_3_8_pre_004_facade_raw_capability_inventory_is_exactly_eleven() { let store = include_str!("../src/store.rs"); let capability_impls = [ "impl ksp_store_api::RawAccountObservationRead for Store", "impl ksp_store_api::RawAccountObservationWrite for Store", "impl ksp_store_api::RawAccountStateRead for Store", "impl ksp_store_api::RawAccountStateWrite for Store", + "impl ksp_store_api::RawTransactionInspectionRead for Store", "impl ksp_store_api::RawTransactionObservationRead for Store", "impl ksp_store_api::RawTransactionObservationWrite for Store", "impl ksp_store_api::RawTransactionRead for Store", @@ -309,20 +310,22 @@ fn pre_010_facade_raw_capability_inventory_is_exactly_ten() { for implementation in capability_impls { assert_eq!(store.matches(implementation).count(), 1, "unexpected Store capability implementation inventory: {implementation}"); } - assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 10); - assert_eq!(store.matches("validate_operation_network(").count(), 14); + assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 11); + assert_eq!(store.matches("validate_operation_network(").count(), 15); return; } #[test] -fn pre_010_facade_and_backend_raw_capability_sets_match_exactly_without_account_retention() { +fn v0_3_8_pre_004_facade_and_backend_capability_sets_match_with_transaction_inspection_only() { let store = include_str!("../src/store.rs"); let backend = include_str!("../../ksp-store-postgres-lib/src/runtime.rs"); let store_traits = raw_capability_trait_names(store, " for Store"); let backend_traits = raw_capability_trait_names(backend, " for PostgresBackend"); - assert_eq!(store_traits.len(), 10); - assert_eq!(backend_traits.len(), 10); + assert_eq!(store_traits.len(), 11); + assert_eq!(backend_traits.len(), 11); assert_eq!(store_traits, backend_traits); + assert!(store_traits.contains(&"RawTransactionInspectionRead")); + assert!(!store_traits.contains(&"RawAccountStateInspectionRead")); for forbidden in ["RawAccountRetentionRead", "RawAccountRetentionWrite", "RawAccountDelete", "RawAccountCompaction"] { assert!(!store_traits.contains(&forbidden), "unexpected account capability added to Store: {forbidden}"); assert!(!backend_traits.contains(&forbidden), "unexpected account capability added to PostgreSQL backend: {forbidden}"); diff --git a/crates/ksp-store-lib/tests/public_api.rs b/crates/ksp-store-lib/tests/public_api.rs index b8310c5..892fea7 100644 --- a/crates/ksp-store-lib/tests/public_api.rs +++ b/crates/ksp-store-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-lib/tests/public_api.rs -// version: 9 +// version: 10 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -83,6 +83,7 @@ where + ksp_store_lib::RawAccountObservationWrite + ksp_store_lib::RawAccountStateRead + ksp_store_lib::RawAccountStateWrite + + ksp_store_lib::RawTransactionInspectionRead + ksp_store_lib::RawTransactionObservationRead + ksp_store_lib::RawTransactionObservationWrite + ksp_store_lib::RawTransactionRead @@ -95,7 +96,7 @@ where } #[test] -fn pre_008_store_facade_implements_exact_raw_capability_set_10_of_10() { +fn v0_3_8_pre_004_store_facade_implements_raw_transaction_inspection_as_capability_11_of_11() { assert_raw_capabilities::(); return; } diff --git a/crates/ksp-store-postgres-lib/src/error.rs b/crates/ksp-store-postgres-lib/src/error.rs index 70cba4f..efc3c88 100644 --- a/crates/ksp-store-postgres-lib/src/error.rs +++ b/crates/ksp-store-postgres-lib/src/error.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/error.rs -// version: 8 +// version: 9 /// Stable KSP error code reserved for PostgreSQL retention transitions that require unsupported physical compaction. pub const ERROR_CODE_POSTGRES_RETENTION_COMPACTION_UNSUPPORTED: ksp_store_api::ErrorCode = @@ -23,7 +23,7 @@ pub enum PostgresBackendErrorKind { DataInvalid, /// PostgreSQL migration/bootstrap execution failed without exposing server text or SQL. MigrationFailed, - /// The requested RAW page size cannot be represented by PostgreSQL LIMIT plus the continuation probe row. + /// The requested RAW page size cannot be represented by the physical PostgreSQL LIMIT domain. PageLimitUnsupported, /// Applied PostgreSQL migration history diverges from the embedded immutable KSP history. MigrationMismatch, diff --git a/crates/ksp-store-postgres-lib/src/lib.rs b/crates/ksp-store-postgres-lib/src/lib.rs index 55784f8..a287f54 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: 21 +// version: 22 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -33,6 +33,8 @@ //! deterministic account keyset pagination with the fixed `KSPA` cursor. `0.3.4-pre.008` //! implements the four `RawAccount*` capabilities directly on `PostgresBackend`, completing //! the backend RAW capability inventory at ten without exposing physical PostgreSQL types. +//! `0.3.8-pre.004` adds the payload-free random-access RawTransaction inspection +//! capability with exact counts while preserving keyset traversal unchanged. //! //! 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 @@ -102,6 +104,8 @@ pub(crate) use self::raw_transaction::get_raw_transaction_observation; pub(crate) use self::raw_transaction::get_raw_transaction_retention_state; /// Private RAW transaction tombstone reader consumed by the physical backend runtime. pub(crate) use self::raw_transaction::get_raw_transaction_tombstone; +/// Private payload-free RAW transaction inspection reader consumed by the physical backend runtime. +pub(crate) use self::raw_transaction::inspect_raw_transactions; /// Private RAW transaction list reader consumed by the physical backend runtime. pub(crate) use self::raw_transaction::list_raw_transactions; /// Private atomic RAW transaction acquisition writer consumed by the physical backend runtime. diff --git a/crates/ksp-store-postgres-lib/src/raw_transaction.rs b/crates/ksp-store-postgres-lib/src/raw_transaction.rs index a0a3393..a5653f8 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: 4 +// version: 5 pub(crate) mod cursor; @@ -12,6 +12,8 @@ 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 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"; +const INSPECT_TRANSACTIONS_DESC_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 DESC, transaction_row.signature DESC LIMIT $3 OFFSET $4) AS page ON TRUE"; const LIST_TRANSACTIONS_ASC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, signature) > ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot ASC, signature ASC LIMIT $5"; const LIST_TRANSACTIONS_DESC_SQL: &str = "SELECT signature, slot::text AS slot_text FROM ksp_raw_transactions WHERE retention_state <> 'purged' AND ($1::TEXT IS NULL OR slot >= $1::TEXT::NUMERIC) AND ($2::TEXT IS NULL OR slot <= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR (slot, signature) < ($3::TEXT::NUMERIC, $4::BYTEA)) ORDER BY slot DESC, signature DESC LIMIT $5"; const LOCK_OBSERVATION_SQL: &str = "SELECT observation_key, transaction_signature, provider, protocol, acquisition_method, origin, received_at_unix_millis, capture_session_id, commitment, endpoint_id, filter_id, observed_at_unix_millis, source_payload_hash, source_payload_size_bytes FROM ksp_raw_transaction_observations WHERE observation_key = $1 FOR UPDATE"; @@ -30,6 +32,22 @@ struct RawListDbRow { slot_text: std::string::String, } +struct RawInspectionDbRow { + archive_payload_present: std::option::Option, + block_time_unix_millis: std::option::Option, + content_hash: std::option::Option>, + filtered_count_text: std::string::String, + format_id: std::option::Option, + format_version: std::option::Option, + hot_payload_present: std::option::Option, + page_present: bool, + payload_size_bytes: std::option::Option, + retention_state: std::option::Option, + signature: std::option::Option>, + slot_text: std::option::Option, + total_count_text: std::string::String, +} + struct RawObservationDbRow { acquisition_method: std::string::String, capture_session_id: std::option::Option, @@ -192,6 +210,103 @@ pub(crate) async fn list_raw_transactions( return std::result::Result::Ok(ksp_store_api::RawPage::new(items, next_cursor)); } +/// Inspects one payload-free random-access RAW transaction window with exact counts. +pub(crate) async fn inspect_raw_transactions( + pool: &deadpool_postgres::Pool, + network: &ksp_store_api::RawNetworkId, + query: &ksp_store_api::RawTransactionInspectionQuery, +) -> std::result::Result, crate::PostgresBackendError> { + if query.network() != network { + return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WrongNetwork, "raw_transaction_inspection_network")); + } + let sql = match query.direction() { + ksp_store_api::RawSortDirection::Ascending => INSPECT_TRANSACTIONS_ASC_SQL, + ksp_store_api::RawSortDirection::Descending => INSPECT_TRANSACTIONS_DESC_SQL, + _ => { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::QueryInvalid, + "raw_transaction_inspection_direction", + )); + }, + }; + let (sql_limit, sql_offset) = match raw_transaction_inspection_sql_window(query.page()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let client_result = pool.get().await; + let client = match client_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(crate::map_pool_error(error)), + }; + let slots = query.slots(); + let start_text = slots.start_inclusive().map(|value| return value.to_string()); + let end_text = slots.end_inclusive().map(|value| return value.to_string()); + let rows_result = client.query(sql, &[&start_text, &end_text, &sql_limit, &sql_offset]).await; + let rows = match rows_result { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::ReadFailed, "raw_transaction_inspection_query")); + }, + }; + if rows.is_empty() { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_cardinality")); + } + let row_count = rows.len(); + let mut total_items = std::option::Option::None; + let mut filtered_items = std::option::Option::None; + let mut items = std::vec::Vec::new(); + for row in rows { + let physical = match raw_inspection_db_row(&row) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let row_total = match decode_u64_decimal(physical.total_count_text.as_str(), "raw_transaction_inspection_total") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let row_filtered = match decode_u64_decimal(physical.filtered_count_text.as_str(), "raw_transaction_inspection_filtered") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + match (total_items, filtered_items) { + (std::option::Option::None, std::option::Option::None) => { + total_items = std::option::Option::Some(row_total); + filtered_items = std::option::Option::Some(row_filtered); + }, + (std::option::Option::Some(total), std::option::Option::Some(filtered)) if total == row_total && filtered == row_filtered => {}, + _ => return std::result::Result::Err(data_invalid("raw_transaction_inspection_counts")), + } + if physical.page_present { + let summary = match decode_raw_inspection_summary(network, physical) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + items.push(summary); + } else if row_count != 1 || !raw_inspection_empty_page_is_clean(&physical) { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_page")); + } + } + let total_items = match total_items { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(data_invalid("raw_transaction_inspection_total")), + }; + let filtered_items = match filtered_items { + std::option::Option::Some(value) => value, + std::option::Option::None => return std::result::Result::Err(data_invalid("raw_transaction_inspection_filtered")), + }; + let item_count = match u64::try_from(items.len()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_page")), + }; + if item_count > query.page().limit().get() { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_page")); + } + return match ksp_store_api::RawInspectionPage::try_new(items, total_items, filtered_items) { + std::result::Result::Ok(value) => std::result::Result::Ok(value), + std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_transaction_inspection_page")), + }; +} + /// Reads one RAW transaction observation from the physical PostgreSQL backend. pub(crate) async fn get_raw_transaction_observation( pool: &deadpool_postgres::Pool, @@ -1009,6 +1124,221 @@ fn write_failed(phase: &'static str) -> crate::PostgresBackendError { return crate::PostgresBackendError::new(crate::PostgresBackendErrorKind::WriteFailed, phase); } +fn raw_inspection_empty_page_is_clean(row: &RawInspectionDbRow) -> bool { + return !row.page_present + && row.archive_payload_present.is_none() + && row.block_time_unix_millis.is_none() + && row.content_hash.is_none() + && row.format_id.is_none() + && row.format_version.is_none() + && row.hot_payload_present.is_none() + && row.payload_size_bytes.is_none() + && row.retention_state.is_none() + && row.signature.is_none() + && row.slot_text.is_none(); +} + +fn raw_transaction_inspection_sql_window(page: ksp_store_api::RawInspectionPageRequest) -> std::result::Result<(i64, i64), crate::PostgresBackendError> { + let sql_limit = match i64::try_from(page.limit().get()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::PageLimitUnsupported, + "raw_transaction_inspection_limit", + )); + }, + }; + let sql_offset = match i64::try_from(page.offset()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::QueryInvalid, + "raw_transaction_inspection_offset", + )); + }, + }; + return std::result::Result::Ok((sql_limit, sql_offset)); +} + +fn raw_inspection_db_row(row: &tokio_postgres::Row) -> std::result::Result { + let archive_payload_present = match row.try_get::<_, std::option::Option>("archive_payload_present") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let block_time_unix_millis = match row.try_get::<_, std::option::Option>("block_time_unix_millis") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let content_hash = match row.try_get::<_, std::option::Option>>("content_hash") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let filtered_count_text = match row.try_get::<_, std::string::String>("filtered_count_text") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let format_id = match row.try_get::<_, std::option::Option>("format_id") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let format_version = match row.try_get::<_, std::option::Option>("format_version") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let hot_payload_present = match row.try_get::<_, std::option::Option>("hot_payload_present") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let page_present = match row.try_get::<_, bool>("page_present") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let payload_size_bytes = match row.try_get::<_, std::option::Option>("payload_size_bytes") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let retention_state = match row.try_get::<_, std::option::Option>("retention_state") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let signature = match row.try_get::<_, std::option::Option>>("signature") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let slot_text = match row.try_get::<_, std::option::Option>("slot_text") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + let total_count_text = match row.try_get::<_, std::string::String>("total_count_text") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_decode")), + }; + return std::result::Result::Ok(RawInspectionDbRow { + archive_payload_present, + block_time_unix_millis, + content_hash, + filtered_count_text, + format_id, + format_version, + hot_payload_present, + page_present, + payload_size_bytes, + retention_state, + signature, + slot_text, + total_count_text, + }); +} + +fn decode_raw_inspection_summary( + network: &ksp_store_api::RawNetworkId, + row: RawInspectionDbRow, +) -> std::result::Result { + if !row.page_present { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_page")); + } + let signature = match inspection_required(row.signature, "raw_transaction_inspection_signature") { + std::result::Result::Ok(value) => match fixed_bytes::<64>(value) { + std::result::Result::Ok(bytes) => ksp_store_api::RawTransactionSignature::new(bytes), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let slot_text = match inspection_required(row.slot_text, "raw_transaction_inspection_slot") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let slot = match decode_u64_decimal(slot_text.as_str(), "raw_transaction_inspection_slot") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let block_time = match decode_optional_timestamp(row.block_time_unix_millis, "raw_transaction_inspection_block_time") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_id_raw = match inspection_required(row.format_id, "raw_transaction_inspection_format_id") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_id = match decode_format_id(format_id_raw) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_version_raw = match inspection_required(row.format_version, "raw_transaction_inspection_format_version") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let format_version = match decode_u32_i64(format_version_raw, "raw_transaction_inspection_format_version") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let content_hash_raw = match inspection_required(row.content_hash, "raw_transaction_inspection_content_hash") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let content_hash = match fixed_bytes::<32>(content_hash_raw) { + std::result::Result::Ok(value) => ksp_store_api::RawContentHash::new(value), + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let retention_raw = match inspection_required(row.retention_state, "raw_transaction_inspection_retention") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let retention_state = match decode_retention_state(retention_raw.as_str()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let hot_payload_present = match inspection_required(row.hot_payload_present, "raw_transaction_inspection_hot_payload") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let archive_payload_present = match inspection_required(row.archive_payload_present, "raw_transaction_inspection_archive_payload") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let payload_size_bytes = match row.payload_size_bytes { + std::option::Option::Some(value) => match u64::try_from(value) { + std::result::Result::Ok(decoded) => std::option::Option::Some(decoded), + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_inspection_payload_size")), + }, + std::option::Option::None => std::option::Option::None, + }; + match retention_state { + ksp_store_api::RawRetentionState::Full if !hot_payload_present || archive_payload_present => { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_full_shape")); + }, + ksp_store_api::RawRetentionState::Archived if hot_payload_present || !archive_payload_present => { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_archived_shape")); + }, + ksp_store_api::RawRetentionState::Purged if hot_payload_present || archive_payload_present || block_time.is_some() => { + return std::result::Result::Err(data_invalid("raw_transaction_inspection_purged_shape")); + }, + ksp_store_api::RawRetentionState::Full | ksp_store_api::RawRetentionState::Archived | ksp_store_api::RawRetentionState::Purged => {}, + _ => return std::result::Result::Err(data_invalid("raw_transaction_inspection_retention")), + } + let reference = ksp_store_api::RawTransactionReference::new(network.clone(), signature); + return match ksp_store_api::RawTransactionSummary::try_new( + reference, + slot, + block_time, + format_id, + format_version, + content_hash, + payload_size_bytes, + retention_state, + ) { + std::result::Result::Ok(value) => std::result::Result::Ok(value), + std::result::Result::Err(_) => std::result::Result::Err(data_invalid("raw_transaction_inspection_summary")), + }; +} + +fn inspection_required(value: std::option::Option, phase: &'static str) -> std::result::Result { + return match value { + std::option::Option::Some(inner) => std::result::Result::Ok(inner), + std::option::Option::None => std::result::Result::Err(data_invalid(phase)), + }; +} + fn raw_list_db_row(row: &tokio_postgres::Row) -> std::result::Result { let signature = match row.try_get::<_, std::vec::Vec>("signature") { std::result::Result::Ok(value) => value, diff --git a/crates/ksp-store-postgres-lib/src/runtime.rs b/crates/ksp-store-postgres-lib/src/runtime.rs index 90bc7ab..db6697d 100644 --- a/crates/ksp-store-postgres-lib/src/runtime.rs +++ b/crates/ksp-store-postgres-lib/src/runtime.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/runtime.rs -// version: 15 +// version: 16 const APPLICATION_NAME: &str = "ksp-store"; const MAX_CONNECTION_URI_BYTES: usize = 4_096; @@ -389,6 +389,14 @@ impl PostgresBackend { return crate::list_raw_transactions(&self.pool, &self.network, query).await; } + /// Inspects one payload-free random-access RAW transaction window with exact logical counts. + pub async fn inspect_raw_transactions( + &self, + query: &ksp_store_api::RawTransactionInspectionQuery, + ) -> std::result::Result, crate::PostgresBackendError> { + return crate::inspect_raw_transactions(&self.pool, &self.network, query).await; + } + /// Reads one persisted RAW transaction observation by producer-owned idempotence key. pub async fn get_raw_transaction_observation( &self, @@ -550,6 +558,18 @@ impl ksp_store_api::RawTransactionRead for PostgresBackend { } } +impl ksp_store_api::RawTransactionInspectionRead for PostgresBackend { + fn inspect_raw_transactions<'a>( + &'a self, + query: &'a ksp_store_api::RawTransactionInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + return std::boxed::Box::pin(async move { + let result = PostgresBackend::inspect_raw_transactions(self, query).await; + return result.map_err(map_capability_error); + }); + } +} + impl ksp_store_api::RawTransactionWrite for PostgresBackend { fn persist_raw_transaction_acquisition<'a>( &'a self, diff --git a/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs b/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs index 07dc9a6..f9fef2d 100644 --- a/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs +++ b/crates/ksp-store-postgres-lib/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/dependency_boundary.rs -// version: 24 +// version: 25 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -408,13 +408,14 @@ fn pre_007_raw_retention_is_atomic_compare_and_transition_without_fake_compactio } #[test] -fn pre_008_backend_trait_implementations_cover_exact_ten_raw_capabilities_in_runtime_bridge() { +fn v0_3_8_pre_004_backend_trait_implementations_cover_exact_eleven_raw_capabilities_in_runtime_bridge() { let runtime = include_str!("../src/runtime.rs"); for implementation in [ "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", "impl ksp_store_api::RawAccountStateRead for PostgresBackend", "impl ksp_store_api::RawAccountStateWrite for PostgresBackend", + "impl ksp_store_api::RawTransactionInspectionRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationWrite for PostgresBackend", "impl ksp_store_api::RawTransactionRead for PostgresBackend", @@ -424,6 +425,6 @@ fn pre_008_backend_trait_implementations_cover_exact_ten_raw_capabilities_in_run ] { assert_eq!(runtime.matches(implementation).count(), 1, "unexpected PostgreSQL RAW capability inventory: {implementation}"); } - assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 10); + assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 11); return; } diff --git a/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs b/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs index b56343f..b39c30f 100644 --- a/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs +++ b/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/hardening_completeness.rs -// version: 17 +// version: 18 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -281,6 +281,7 @@ fn pre_009_live_raw_transaction_proof_is_opt_in_isolated_and_secret_safe() { "prove_concurrent_identical_insert", "prove_concurrent_divergent_insert", "prove_pagination", + "prove_inspection", "prove_retention_and_rehydrate", "prove_retention_races", "prove_cancellation_rollback", @@ -298,13 +299,14 @@ fn pre_009_live_raw_transaction_proof_is_opt_in_isolated_and_secret_safe() { } #[test] -fn pre_010_raw_capability_implementation_inventory_is_exactly_ten() { +fn v0_3_8_pre_004_raw_capability_implementation_inventory_is_exactly_eleven() { let runtime = include_str!("../src/runtime.rs"); let capability_impls = [ "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", "impl ksp_store_api::RawAccountStateRead for PostgresBackend", "impl ksp_store_api::RawAccountStateWrite for PostgresBackend", + "impl ksp_store_api::RawTransactionInspectionRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationWrite for PostgresBackend", "impl ksp_store_api::RawTransactionRead for PostgresBackend", @@ -315,7 +317,7 @@ fn pre_010_raw_capability_implementation_inventory_is_exactly_ten() { for implementation in capability_impls { assert_eq!(runtime.matches(implementation).count(), 1, "unexpected PostgreSQL capability implementation inventory: {implementation}"); } - assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 10); + assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 11); let migration = include_str!("../src/migration.rs"); assert!(migration.contains("raw_account_state")); assert!(migration.contains("crate::V002_RESOURCES")); @@ -324,25 +326,77 @@ fn pre_010_raw_capability_implementation_inventory_is_exactly_ten() { } #[test] -fn pre_010_raw_transaction_private_sql_keeps_keyset_navigation_and_bounded_statement_surface() { +fn v0_3_8_pre_004_raw_transaction_keyset_sql_remains_offset_free_and_unchanged_in_role() { let source = include_str!("../src/raw_transaction.rs"); - for required in [ - "ORDER BY slot ASC, signature ASC", - "ORDER BY slot DESC, signature DESC", - "LIMIT $5", - "FOR UPDATE", - "ON CONFLICT (signature) DO NOTHING", - "ON CONFLICT (observation_key) DO NOTHING", - "ksp_raw_transaction_archive_payloads", - ] { + let ascending = source.lines().find(|line| line.starts_with("const LIST_TRANSACTIONS_ASC_SQL")); + let ascending = match ascending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing canonical ascending keyset SQL"), + }; + let descending = source.lines().find(|line| line.starts_with("const LIST_TRANSACTIONS_DESC_SQL")); + let descending = match descending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing canonical descending keyset SQL"), + }; + for statement in [ascending, descending] { + assert!(statement.contains("retention_state <> 'purged'")); + assert!(statement.contains("LIMIT $5")); + assert!(!statement.contains(" OFFSET ")); + } + assert!(ascending.contains("(slot, signature) >")); + assert!(ascending.contains("ORDER BY slot ASC, signature ASC")); + assert!(descending.contains("(slot, signature) <")); + assert!(descending.contains("ORDER BY slot DESC, signature DESC")); + for required in ["FOR UPDATE", "ON CONFLICT (signature) DO NOTHING", "ON CONFLICT (observation_key) DO NOTHING", "ksp_raw_transaction_archive_payloads"] { assert!(source.contains(required), "required hardened RawTransaction SQL contract missing: {required}"); } - for forbidden in [" OFFSET ", "SELECT *", "ON CONFLICT DO UPDATE", "processing_state", "batch_size", "priority"] { + for forbidden in ["SELECT *", "ON CONFLICT DO UPDATE", "processing_state", "batch_size", "priority"] { assert!(!source.contains(forbidden), "forbidden RawTransaction scope/policy SQL detected: {forbidden}"); } return; } +#[test] +fn v0_3_8_pre_004_transaction_inspection_sql_is_single_statement_payload_free_counted_and_random_access() { + let source = include_str!("../src/raw_transaction.rs"); + let ascending = source.lines().find(|line| line.starts_with("const INSPECT_TRANSACTIONS_ASC_SQL")); + let ascending = match ascending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing ascending inspection SQL"), + }; + let descending = source.lines().find(|line| line.starts_with("const INSPECT_TRANSACTIONS_DESC_SQL")); + let descending = match descending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing descending inspection SQL"), + }; + for statement in [ascending, descending] { + for required in [ + "COUNT(*)::TEXT AS filtered_count_text", + "COUNT(*)::TEXT FROM ksp_raw_transactions", + "LEFT JOIN LATERAL", + "LEFT JOIN ksp_raw_transaction_archive_payloads", + "OCTET_LENGTH(transaction_row.payload)::BIGINT", + "OCTET_LENGTH(archive_row.payload)::BIGINT", + "LIMIT $3 OFFSET $4", + "payload_size_bytes", + "hot_payload_present", + "archive_payload_present", + ] { + assert!(statement.contains(required), "missing transaction inspection SQL contract: {required}"); + } + assert!(!statement.contains("transaction_row.payload, transaction_row.retention_state")); + assert!(!statement.contains("archive_row.payload AS")); + assert!(!statement.contains("retention_state <> 'purged'")); + } + assert!(ascending.contains("ORDER BY transaction_row.slot ASC, transaction_row.signature ASC")); + assert!(descending.contains("ORDER BY transaction_row.slot DESC, transaction_row.signature DESC")); + assert_eq!(source.matches(" OFFSET ").count(), 2); + assert!(source.contains("raw_transaction_inspection_sql_window(query.page())")); + assert!(source.contains("raw_transaction_inspection_offset")); + assert!(source.contains("raw_transaction_inspection_limit")); + return; +} + #[test] fn pre_010_raw_account_private_sql_is_non_destructive_keyset_and_family_local() { let source = include_str!("../src/raw_account.rs"); 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 c4943f7..cbda70d 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: 5 +// version: 6 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -183,6 +183,7 @@ async fn run_raw_transaction_scenario(admin: &mut tokio_postgres::Client, uri: & prove_atomic_insert_and_reads(&backend).await, prove_additional_observation_and_atomic_rollback(&backend).await, prove_pagination(&backend).await, + prove_inspection(&backend).await, prove_retention_and_rehydrate(&backend).await, ] { if let std::result::Result::Err(error) = proof { @@ -476,6 +477,57 @@ async fn prove_pagination(backend: &ksp_store_postgres_lib::PostgresBackend) -> return std::result::Result::Ok(()); } +async fn prove_inspection(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { + let ascending_query = + match inspection_query(ksp_store_api::RawSortDirection::Ascending, 1, 2, std::option::Option::Some(9_000), std::option::Option::Some(9_002)) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let ascending = match backend.inspect_raw_transactions(&ascending_query).await { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("inspection_ascending_query")), + }; + if ascending.filtered_items() != 5 || ascending.total_items() < ascending.filtered_items() || ascending.items().len() != 2 { + return std::result::Result::Err(LiveFailure::new("inspection_ascending_counts")); + } + if ascending.items()[0].reference().signature().as_bytes()[0] != 51 || ascending.items()[1].reference().signature().as_bytes()[0] != 52 { + return std::result::Result::Err(LiveFailure::new("inspection_ascending_order")); + } + for item in ascending.items() { + if item.payload_size_bytes() != std::option::Option::Some(4) || item.retention_state() != ksp_store_api::RawRetentionState::Full { + return std::result::Result::Err(LiveFailure::new("inspection_ascending_summary")); + } + } + let descending_query = + match inspection_query(ksp_store_api::RawSortDirection::Descending, 1, 2, std::option::Option::Some(9_000), std::option::Option::Some(9_002)) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let descending = match backend.inspect_raw_transactions(&descending_query).await { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("inspection_descending_query")), + }; + if descending.filtered_items() != 5 || descending.items().len() != 2 { + return std::result::Result::Err(LiveFailure::new("inspection_descending_counts")); + } + if descending.items()[0].reference().signature().as_bytes()[0] != 53 || descending.items()[1].reference().signature().as_bytes()[0] != 52 { + return std::result::Result::Err(LiveFailure::new("inspection_descending_order")); + } + let empty_query = + match inspection_query(ksp_store_api::RawSortDirection::Ascending, 99, 25, std::option::Option::Some(9_000), std::option::Option::Some(9_002)) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let empty = match backend.inspect_raw_transactions(&empty_query).await { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("inspection_empty_query")), + }; + if empty.filtered_items() != 5 || !empty.items().is_empty() { + return std::result::Result::Err(LiveFailure::new("inspection_empty_page")); + } + return std::result::Result::Ok(()); +} + async fn prove_retention_and_rehydrate(backend: &ksp_store_postgres_lib::PostgresBackend) -> std::result::Result<(), LiveFailure> { let transaction = match raw_transaction(70, 7_000, 70) { std::result::Result::Ok(value) => value, @@ -750,6 +802,29 @@ async fn collect_pages( return std::result::Result::Ok(output); } +fn inspection_query( + direction: ksp_store_api::RawSortDirection, + offset: u64, + limit: u64, + start: std::option::Option, + end: std::option::Option, +) -> std::result::Result { + let network = match network("devnet") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let slots = match ksp_store_api::RawSlotRange::new(start, end) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("inspection_slot_range")), + }; + let limit = match ksp_store_api::RawPageLimit::new(limit) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(LiveFailure::new("inspection_page_limit")), + }; + let page = ksp_store_api::RawInspectionPageRequest::new(offset, limit); + return std::result::Result::Ok(ksp_store_api::RawTransactionInspectionQuery::new(network, slots, direction, page)); +} + fn page_query( direction: ksp_store_api::RawSortDirection, cursor: std::option::Option, diff --git a/crates/ksp-store-postgres-lib/tests/public_api.rs b/crates/ksp-store-postgres-lib/tests/public_api.rs index c436022..edb7bf3 100644 --- a/crates/ksp-store-postgres-lib/tests/public_api.rs +++ b/crates/ksp-store-postgres-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/tests/public_api.rs -// version: 14 +// version: 15 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -117,8 +117,9 @@ fn pre_005_raw_write_bridge_uses_only_backend_independent_models_and_outcomes() } #[test] -fn pre_006_raw_list_bridge_uses_backend_independent_query_page_and_reference_models() { +fn v0_3_8_pre_004_raw_list_and_inspection_bridges_use_only_backend_independent_models() { let _list = ksp_store_postgres_lib::PostgresBackend::list_raw_transactions; + let _inspection = ksp_store_postgres_lib::PostgresBackend::inspect_raw_transactions; return; } @@ -134,6 +135,7 @@ where + ksp_store_api::RawAccountObservationWrite + ksp_store_api::RawAccountStateRead + ksp_store_api::RawAccountStateWrite + + ksp_store_api::RawTransactionInspectionRead + ksp_store_api::RawTransactionObservationRead + ksp_store_api::RawTransactionObservationWrite + ksp_store_api::RawTransactionRead @@ -146,7 +148,7 @@ where } #[test] -fn pre_008_postgres_backend_implements_exact_raw_capability_set_10_of_10() { +fn v0_3_8_pre_004_postgres_backend_implements_transaction_inspection_as_capability_11_of_11() { assert_raw_capabilities::(); 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 ba9bc04..2a0cf32 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: 4 +// version: 5 fn network() -> ksp_store_api::RawNetworkId { return match ksp_store_api::RawNetworkId::new("devnet") { @@ -169,6 +169,105 @@ fn pre_004_retention_and_tombstone_decoding_is_exact() { return; } +fn inspection_row( + retention_state: &str, + hot_payload_present: bool, + archive_payload_present: bool, + payload_size_bytes: std::option::Option, + block_time_unix_millis: std::option::Option, +) -> super::RawInspectionDbRow { + return super::RawInspectionDbRow { + archive_payload_present: std::option::Option::Some(archive_payload_present), + block_time_unix_millis, + content_hash: std::option::Option::Some(vec![7; 32]), + filtered_count_text: "3".to_owned(), + format_id: std::option::Option::Some("ksp.raw.transaction".to_owned()), + format_version: std::option::Option::Some(1), + hot_payload_present: std::option::Option::Some(hot_payload_present), + page_present: true, + payload_size_bytes, + retention_state: std::option::Option::Some(retention_state.to_owned()), + signature: std::option::Option::Some(vec![9; 64]), + slot_text: std::option::Option::Some(u64::MAX.to_string()), + total_count_text: "4".to_owned(), + }; +} + +#[test] +fn v0_3_8_pre_004_inspection_window_rejects_unrepresentable_limit_and_offset_before_io() { + let normal_limit = match ksp_store_api::RawPageLimit::new(100) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid inspection page limit rejected: {error:?}"), + }; + let normal = super::raw_transaction_inspection_sql_window(ksp_store_api::RawInspectionPageRequest::new(25, normal_limit)); + assert_eq!(normal, std::result::Result::Ok((100, 25))); + let huge_limit = match ksp_store_api::RawPageLimit::new(u64::MAX) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("API unexpectedly rejected backend-physical limit fixture: {error:?}"), + }; + let limit_error = super::raw_transaction_inspection_sql_window(ksp_store_api::RawInspectionPageRequest::new(0, huge_limit)); + assert_eq!( + limit_error.err().map(|value| return (value.kind(), value.phase())), + std::option::Option::Some((crate::PostgresBackendErrorKind::PageLimitUnsupported, "raw_transaction_inspection_limit")), + ); + let offset_error = super::raw_transaction_inspection_sql_window(ksp_store_api::RawInspectionPageRequest::new(u64::MAX, normal_limit)); + assert_eq!( + offset_error.err().map(|value| return (value.kind(), value.phase())), + std::option::Option::Some((crate::PostgresBackendErrorKind::QueryInvalid, "raw_transaction_inspection_offset")), + ); + return; +} + +#[test] +fn v0_3_8_pre_004_inspection_summary_decoding_handles_full_archived_and_purged_without_payload_bytes() { + let network = network(); + let full = super::decode_raw_inspection_summary( + &network, + inspection_row("full", true, false, std::option::Option::Some(4), std::option::Option::Some(1_700_000_000_000)), + ); + let full = match full { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid full inspection summary rejected: {error:?}"), + }; + assert_eq!(full.slot(), u64::MAX); + assert_eq!(full.payload_size_bytes(), std::option::Option::Some(4)); + assert_eq!(full.retention_state(), ksp_store_api::RawRetentionState::Full); + let archived = super::decode_raw_inspection_summary( + &network, + inspection_row("archived", false, true, std::option::Option::Some(5), std::option::Option::Some(1_700_000_000_000)), + ); + let archived = match archived { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid archived inspection summary rejected: {error:?}"), + }; + assert_eq!(archived.payload_size_bytes(), std::option::Option::Some(5)); + assert_eq!(archived.retention_state(), ksp_store_api::RawRetentionState::Archived); + let purged = super::decode_raw_inspection_summary(&network, inspection_row("purged", false, false, std::option::Option::None, std::option::Option::None)); + let purged = match purged { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => panic!("valid purged inspection summary rejected: {error:?}"), + }; + assert_eq!(purged.payload_size_bytes(), std::option::Option::None); + assert_eq!(purged.block_time(), std::option::Option::None); + assert_eq!(purged.retention_state(), ksp_store_api::RawRetentionState::Purged); + return; +} + +#[test] +fn v0_3_8_pre_004_inspection_summary_rejects_physical_retention_shape_residue() { + let network = network(); + for malformed in [ + inspection_row("full", true, true, std::option::Option::Some(4), std::option::Option::Some(1)), + inspection_row("archived", true, true, std::option::Option::Some(4), std::option::Option::Some(1)), + inspection_row("purged", false, true, std::option::Option::None, std::option::Option::None), + inspection_row("purged", false, false, std::option::Option::None, std::option::Option::Some(1)), + ] { + let result = super::decode_raw_inspection_summary(&network, malformed); + assert_eq!(result.err().map(|value| return value.kind()), std::option::Option::Some(crate::PostgresBackendErrorKind::DataInvalid)); + } + return; +} + #[test] fn pre_004_wrong_network_is_rejected_by_the_private_pre_io_guard() { let backend_network = network(); diff --git a/deltas/0.3.8/pre.004.md b/deltas/0.3.8/pre.004.md new file mode 100644 index 0000000..b03d515 --- /dev/null +++ b/deltas/0.3.8/pre.004.md @@ -0,0 +1,198 @@ + + + +# Delta `0.3.8-pre.004` — inspection PostgreSQL RawTransaction et dispatch Store + +## Base requise + +Base directe attendue : + +```text +0.3.8-pre.003-fix.001 +workspace.package.version = 0.3.8-pre.3.fix.1 +``` + +Le gate opérateur fourni sur cette base est entièrement vert : audits Rust/Markdown, `cargo check --workspace`, `cargo clippy --workspace --all-targets`, `ksp-store-api`, `ksp-store-lib`, `ksp-store-lib --no-default-features` et `ksp-app-store-desk`. + +La livraison est : + +```text +0.3.8-pre.004 +workspace.package.version = 0.3.8-pre.4 +commit = v0.3.8-pre.004 +tag = aucun +``` + +## Objectif + +Implémenter le premier vertical slice physique de la primitive backend-neutral créée en `pre.003` : + +```text +RawTransactionInspectionRead + ksp-store-postgres-lib + -> ksp-store-lib::Store +``` + +`RawAccountStateInspectionRead` reste volontairement non implémenté jusqu'à `pre.005`. Aucun SQL de migration, aucune UI Store Desk et aucun branchement DataTables `serverSide` ne sont introduits. + +## SQL d'inspection + +Deux statements privés existent, un par direction canonique. Chaque appel exécute une seule instruction SQL comprenant : + +```text +filtered_count +conditional total_count +LATERAL page +ORDER BY slot, signature +LIMIT/OFFSET +``` + +Cette forme garantit que counts et page sont issus du même snapshot de statement PostgreSQL. + +Lorsque le range slot est absent, `filtered_count` est réutilisé comme `total_count`; aucun second `COUNT(*)` n'est exécuté. Avec un filtre slot, `total_items` reste le total exact de la famille transaction dans le Store/network et `filtered_items` représente le sous-ensemble filtré. + +Les tombstones `Purged` sont inclus dans l'inspection afin que l'outil opérateur puisse voir l'état de rétention durable. Le listage cursor/keyset historique reste inchangé et continue d'exclure `Purged`. + +## Pas de transfert payload / pas de N+1 + +Le statement de page ne sélectionne jamais : + +```text +transaction_row.payload +archive_row.payload +``` + +Il projette seulement : + +```text +OCTET_LENGTH(transaction_row.payload) +OCTET_LENGTH(archive_row.payload) +hot_payload_present +archive_payload_present +``` + +Le summary final est donc construit en une seule requête, sans chargement individuel par ligne et sans gros bytes IPC/runtime. Les marqueurs de présence permettent en outre de rejeter une forme PostgreSQL incohérente (`Full` avec archive résiduelle, `Archived` avec hot payload, `Purged` avec payload ou block time résiduel). + +## Limites physiques et offset + +`RawPageLimit` reste sans plafond politique côté API. Le backend convertit avant I/O : + +```text +limit -> i64 PostgreSQL +offset -> i64 PostgreSQL +``` + +Un limit non représentable renvoie `PageLimitUnsupported`; un offset non représentable renvoie `QueryInvalid`. Ces validations sont effectuées avant `pool.get()`. + +Une page vide causée par un offset supérieur à `filtered_items` reste valide et retourne les counts exacts avec `items = []`. + +## Façade Store + +`PostgresBackend` et `Store` implémentent désormais `RawTransactionInspectionRead`. Le dispatch Store applique d'abord le même guard network que les autres opérations network-scoped. + +L'inventaire physique/facade passe donc temporairement de 10 à 11 capabilities RAW implémentées : + +```text +10 historiques ++ RawTransactionInspectionRead += 11 +``` + +Les deux inventaires restent strictement symétriques et `RawAccountStateInspectionRead` est explicitement absent jusqu'à `pre.005`. + +## Cursor/keyset préservé + +Les statements historiques `LIST_TRANSACTIONS_ASC_SQL` / `LIST_TRANSACTIONS_DESC_SQL` restent contrôlés séparément : + +```text +keyset (slot, signature) +LIMIT + continuation probe +aucun OFFSET +retention_state <> 'purged' +``` + +L'exception `OFFSET` est donc locale au chemin inspection random-access et ne modifie pas le contrat machine/replay/backfill. + +## Tests et canaris + +Le delta ajoute/actualise les preuves suivantes : + +- conversion limit/offset hostile avant I/O ; +- décodage summary `Full` / `Archived` / `Purged` sans bytes ; +- rejet des résidus physiques incompatibles avec l'état de rétention ; +- canari SQL séparant explicitement keyset sans OFFSET et inspection avec OFFSET ; +- preuve que l'inspection utilise counts, LATERAL page et `OCTET_LENGTH` sans sélectionner les payloads ; +- inventaires exacts 11/11 côté backend et façade ; +- live proof PostgreSQL opt-in : filtre slot, counts, asc/desc, offset et page vide profonde. + +## Fichiers ajoutés + +```text +deltas/0.3.8/pre.004.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-store-lib/src/store.rs +crates/ksp-store-lib/tests/dependency_boundary.rs +crates/ksp-store-lib/tests/hardening_completeness.rs +crates/ksp-store-lib/tests/public_api.rs +crates/ksp-store-postgres-lib/src/error.rs +crates/ksp-store-postgres-lib/src/lib.rs +crates/ksp-store-postgres-lib/src/raw_transaction.rs +crates/ksp-store-postgres-lib/src/runtime.rs +crates/ksp-store-postgres-lib/tests/dependency_boundary.rs +crates/ksp-store-postgres-lib/tests/hardening_completeness.rs +crates/ksp-store-postgres-lib/tests/postgres_raw_transaction_live.rs +crates/ksp-store-postgres-lib/tests/public_api.rs +crates/ksp-store-postgres-lib/unit_tests/raw_transaction.rs +docs/plans/029-V0_3_8_STORE_DESK_PLAN.md +docs/validation/025-V0_3_8_STORE_DESK.md +``` + +## Fichiers supprimés + +```text +aucun +``` + +## Validations d'assemblage + +Contrôles exécutés avant emballage : + +```text +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +Markdown table audit: clean (314 table(s), 696 file(s)) +pre.004 structural contract audit: clean +workspace.package.version: 0.3.8-pre.4 +Store Raw capability implementations: 11 +PostgreSQL Raw capability implementations: 11 +RawAccountStateInspectionRead implementations: 0 +RawTransaction keyset OFFSET occurrences: 0 +RawTransaction inspection OFFSET occurrences: 2 +delta scope: 1 ajout / 16 modifications / 0 suppression +``` + +Le sandbox d'assemblage ne dispose pas de `cargo`/`rustfmt`; aucun gate Cargo de `pre.004` n'est donc revendiqué localement. Le gate opérateur requis est : + +```bash +cargo fmt --all +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 +cargo test -p ksp-store-postgres-lib +cargo test -p ksp-store-lib +cargo check -p ksp-store-lib --no-default-features +cargo test -p ksp-store-api +``` + +Le live proof reste opt-in et doit être exécuté uniquement sur une base PostgreSQL dédiée vide si une preuve réelle est souhaitée : + +```bash +cargo test -p ksp-store-postgres-lib --test postgres_raw_transaction_live -- --ignored --nocapture +``` diff --git a/docs/plans/029-V0_3_8_STORE_DESK_PLAN.md b/docs/plans/029-V0_3_8_STORE_DESK_PLAN.md index cd19c70..5bb3630 100644 --- a/docs/plans/029-V0_3_8_STORE_DESK_PLAN.md +++ b/docs/plans/029-V0_3_8_STORE_DESK_PLAN.md @@ -1,5 +1,5 @@ - + # Plan v0.3.8 — Store Desk V1 RAW @@ -568,7 +568,7 @@ Ajouter dans `ksp-store-api` les types page offset/count, queries Tx/Account, su ### pre.004 — inspection PostgreSQL RawTransaction -Implémenter summary query, counts exacts, offset/limit checked et trait Tx dans `ksp-store-postgres-lib`, puis dispatch façade Store. Conserver intégralement le listage keyset existant. Tests no-payload/N+1, counts filtres, offset hostile et conformance. +Implémenter summary query, counts exacts, offset/limit checked et trait Tx dans `ksp-store-postgres-lib`, puis dispatch façade Store. Conserver intégralement le listage keyset existant. L'inspection utilise une seule instruction SQL par requête afin que `total_items`, `filtered_items` et la page partagent le même snapshot de statement. Le SQL ne sélectionne jamais les payload bytes : il projette uniquement `OCTET_LENGTH` du payload actif/archivé et les marqueurs de présence nécessaires au contrôle de cohérence. Les tombstones `Purged` restent inspectables. Tests no-payload/N+1, counts filtres, offset hostile, pages vides profondes et conformance. ### pre.005 — inspection PostgreSQL RawAccountState diff --git a/docs/validation/025-V0_3_8_STORE_DESK.md b/docs/validation/025-V0_3_8_STORE_DESK.md index 2999c45..d46ff5f 100644 --- a/docs/validation/025-V0_3_8_STORE_DESK.md +++ b/docs/validation/025-V0_3_8_STORE_DESK.md @@ -1,5 +1,5 @@ - + # Validation v0.3.8 — Store Desk V1 RAW @@ -90,7 +90,7 @@ Le journal opérateur fourni pour v0.3.7 rapporte un gate complet. Il reste une ## 8. Gates futures - [X] `pre.002` scaffold KSP strict + DataTables skeleton ; -- [ ] `pre.003` contrats Store inspection offset/count/summaries ; +- [X] `pre.003` contrats Store inspection offset/count/summaries ; - [ ] `pre.004` PostgreSQL + façade inspection RawTransaction ; - [ ] `pre.005` PostgreSQL + façade inspection RawAccountState ; - [ ] `pre.006` composite Logging+Store + Overview/health ; @@ -294,4 +294,42 @@ RawTransactionSummary Le correctif `pre.003-fix.001` met uniquement à jour ce canari exact de `91` à `99` exports. Aucun contrat, comportement runtime, backend PostgreSQL, SQL, migration, pagination ou code Store Desk n'est modifié. La case `pre.003` reste ouverte jusqu'au rejeu vert du gate opérateur sur `0.3.8-pre.3.fix.1`. +## 17. Gate opérateur `pre.003-fix.001` + +Le rejeu opérateur sur `0.3.8-pre.3.fix.1` est entièrement vert : + +```text +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +Markdown table audit: clean (314 table(s), 695 file(s)) +cargo check --workspace: PASS +cargo clippy --workspace --all-targets: PASS +cargo test -p ksp-store-api: PASS +cargo test -p ksp-store-lib: PASS +cargo check -p ksp-store-lib --no-default-features: PASS +cargo test -p ksp-app-store-desk: PASS +``` + +Le canari d'inventaire façade corrigé est vert. `pre.003` est fermé et `pre.004` peut introduire uniquement le vertical slice `RawTransactionInspection` PostgreSQL + façade. + +## 18. `pre.004` — inspection PostgreSQL RawTransaction + façade Store + +- [X] `workspace.package.version = 0.3.8-pre.4` ; +- [X] `PostgresBackend` implémente `RawTransactionInspectionRead` sans implémenter encore `RawAccountStateInspectionRead` ; +- [X] `Store` implémente le même trait et garde le guard network avant dispatch ; +- [X] la query d'inspection PostgreSQL est une instruction unique par direction, donc counts et page partagent le même snapshot SQL ; +- [X] `total_items` compte toute la famille transaction du Store/network, y compris les tombstones `Purged` ; +- [X] `filtered_items` applique uniquement le range slot optionnel ; +- [X] sans filtre slot, le SQL réutilise `filtered_count` comme `total_count` au lieu d'un second count ; +- [X] la page random-access utilise `LIMIT/OFFSET` uniquement dans le chemin d'inspection ; +- [X] limit non représentable par PostgreSQL et offset > `i64::MAX` sont rejetés avant acquisition du pool ; +- [X] les deux queries cursor/keyset historiques restent sans `OFFSET` et excluent toujours `Purged` ; +- [X] les summaries ne sélectionnent ni `transaction_row.payload` ni `archive_row.payload` ; seules les tailles `OCTET_LENGTH` sont projetées ; +- [X] les marqueurs hot/archive payload sont projetés comme booléens pour rejeter les formes physiques incohérentes sans transférer les bytes ; +- [X] les summaries `Full`, `Archived` et `Purged` sont décodés avec les invariants physiques correspondants ; +- [X] une page vide obtenue par offset profond conserve les counts exacts ; +- [X] le live proof PostgreSQL opt-in couvre counts filtrés, ordre asc/desc, offset et page vide ; +- [X] inventaires exacts Store/PostgreSQL passent de 10 à 11 implementations RAW et restent symétriques ; +- [X] aucune migration, table, index, capability Account inspection, UI Desk ou protocole DataTables n'est ajouté dans cette tranche.