From 32c4b67541c17568a750c566373763289e5e2b33 Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Thu, 3 Sep 2026 20:44:21 +0200 Subject: [PATCH] v0.3.8-pre.009 --- Cargo.toml | 4 +- .../src/capability/raw_account.rs | 11 +- .../src/capability/raw_transaction.rs | 11 +- crates/ksp-store-api/src/lib.rs | 14 +- .../ksp-store-api/src/model/raw_inspection.rs | 198 ++++++++++++- .../tests/dependency_boundary.rs | 4 +- .../ksp-store-api/tests/external_backend.rs | 30 +- crates/ksp-store-api/tests/public_api.rs | 6 +- .../tests/release_completeness.rs | 10 +- .../unit_tests/model/raw_inspection.rs | 117 +++++++- crates/ksp-store-lib/src/lib.rs | 14 +- crates/ksp-store-lib/src/store.rs | 62 +++- .../tests/dependency_boundary.rs | 8 +- .../tests/hardening_completeness.rs | 24 +- crates/ksp-store-lib/tests/public_api.rs | 8 +- crates/ksp-store-postgres-lib/src/lib.rs | 8 +- .../ksp-store-postgres-lib/src/raw_account.rs | 279 +++++++++++++++++- .../src/raw_transaction.rs | 244 ++++++++++++++- crates/ksp-store-postgres-lib/src/runtime.rs | 42 ++- .../tests/dependency_boundary.rs | 11 +- .../tests/hardening_completeness.rs | 75 ++++- .../tests/public_api.rs | 10 +- deltas/0.3.8/pre.009.md | 114 +++++++ docs/plans/029-V0_3_8_STORE_DESK_PLAN.md | 6 +- docs/validation/025-V0_3_8_STORE_DESK.md | 53 +++- 25 files changed, 1319 insertions(+), 44 deletions(-) create mode 100644 deltas/0.3.8/pre.009.md diff --git a/Cargo.toml b/Cargo.toml index d56b4ba..23662b5 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 463 +# version: 464 [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.8.fix.1" +version = "0.3.8-pre.9" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-store-api/src/capability/raw_account.rs b/crates/ksp-store-api/src/capability/raw_account.rs index e4815bc..d457c4b 100644 --- a/crates/ksp-store-api/src/capability/raw_account.rs +++ b/crates/ksp-store-api/src/capability/raw_account.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/src/capability/raw_account.rs -// version: 3 +// version: 4 /// Read capability for complete canonical RAW account states. /// @@ -62,6 +62,15 @@ pub trait RawAccountObservationRead: std::marker::Send + std::marker::Sync { ) -> crate::StoreApiFuture<'a, crate::Result>>; } +/// Read capability for random-access RAW account-observation inspection. +pub trait RawAccountObservationInspectionRead: std::marker::Send + std::marker::Sync { + /// Inspects one account-observation window with exact logical counts. + fn inspect_raw_account_observations<'a>( + &'a self, + query: &'a crate::RawAccountObservationInspectionQuery, + ) -> crate::StoreApiFuture<'a, crate::Result>>; +} + /// Write capability for an additional observation of an already persisted RAW account state. /// /// This capability allows repeated HTTP/WS/gRPC acquisitions to be retained diff --git a/crates/ksp-store-api/src/capability/raw_transaction.rs b/crates/ksp-store-api/src/capability/raw_transaction.rs index 9e9d629..09217f1 100644 --- a/crates/ksp-store-api/src/capability/raw_transaction.rs +++ b/crates/ksp-store-api/src/capability/raw_transaction.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/src/capability/raw_transaction.rs -// version: 3 +// version: 4 /// Read capability for canonical RAW transactions. /// @@ -63,6 +63,15 @@ pub trait RawTransactionObservationRead: std::marker::Send + std::marker::Sync { ) -> crate::StoreApiFuture<'a, crate::Result>>; } +/// Read capability for random-access RAW transaction-observation inspection. +pub trait RawTransactionObservationInspectionRead: std::marker::Send + std::marker::Sync { + /// Inspects one transaction-observation window with exact logical counts. + fn inspect_raw_transaction_observations<'a>( + &'a self, + query: &'a crate::RawTransactionObservationInspectionQuery, + ) -> crate::StoreApiFuture<'a, crate::Result>>; +} + /// Write capability for an additional observation of an already persisted RAW transaction. /// /// This capability exists so repeated acquisitions can be recorded without diff --git a/crates/ksp-store-api/src/lib.rs b/crates/ksp-store-api/src/lib.rs index 78543ab..0eff331 100644 --- a/crates/ksp-store-api/src/lib.rs +++ b/crates/ksp-store-api/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/src/lib.rs -// version: 6 +// version: 7 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -22,6 +22,8 @@ mod model; /// Boxed async operation returned by object-safe Store capability contracts. pub use self::capability::StoreApiFuture; +/// Read capability for random-access RAW account-observation inspection. +pub use self::capability::raw_account::RawAccountObservationInspectionRead; /// Read capability for persisted RAW account-state observations. pub use self::capability::raw_account::RawAccountObservationRead; /// Write capability for additional observations of already persisted RAW account states. @@ -38,6 +40,8 @@ pub use self::capability::raw_retention::RawTransactionRetentionRead; pub use self::capability::raw_retention::RawTransactionRetentionWrite; /// Read capability for payload-free random-access RAW transaction inspection. pub use self::capability::raw_transaction::RawTransactionInspectionRead; +/// Read capability for random-access RAW transaction-observation inspection. +pub use self::capability::raw_transaction::RawTransactionObservationInspectionRead; /// Read capability for persisted RAW transaction observations. pub use self::capability::raw_transaction::RawTransactionObservationRead; /// Write capability for additional observations of already persisted RAW transactions. @@ -64,6 +68,10 @@ pub use self::model::raw_account::RawAccountObservation; pub use self::model::raw_account::RawAccountState; /// Durable backend-independent identity of one canonical RAW account state. pub use self::model::raw_account::RawAccountStateReference; +/// Backend-independent random-access inspection query for RAW account observations. +pub use self::model::raw_inspection::RawAccountObservationInspectionQuery; +/// Safe observation summary for one canonical RAW account-state acquisition. +pub use self::model::raw_inspection::RawAccountObservationSummary; /// Backend-independent random-access inspection query for RAW account states. pub use self::model::raw_inspection::RawAccountStateInspectionQuery; /// Data-free summary of one canonical RAW account state for operator inspection. @@ -74,6 +82,10 @@ pub use self::model::raw_inspection::RawInspectionPage; pub use self::model::raw_inspection::RawInspectionPageRequest; /// Backend-independent random-access inspection query for RAW transactions. pub use self::model::raw_inspection::RawTransactionInspectionQuery; +/// Backend-independent random-access inspection query for RAW transaction observations. +pub use self::model::raw_inspection::RawTransactionObservationInspectionQuery; +/// Safe observation summary for one canonical RAW transaction acquisition. +pub use self::model::raw_inspection::RawTransactionObservationSummary; /// Payload-free summary of one canonical RAW transaction for operator inspection. pub use self::model::raw_inspection::RawTransactionSummary; /// Combined outcome of one atomic canonical RAW entity plus observation acquisition. diff --git a/crates/ksp-store-api/src/model/raw_inspection.rs b/crates/ksp-store-api/src/model/raw_inspection.rs index 857a192..837e58a 100644 --- a/crates/ksp-store-api/src/model/raw_inspection.rs +++ b/crates/ksp-store-api/src/model/raw_inspection.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/src/model/raw_inspection.rs -// version: 1 +// version: 2 /// Random-access page request dedicated to bounded interactive RAW inspection. /// @@ -87,6 +87,202 @@ impl RawInspectionPage { } } +/// Backend-independent random-access inspection query for RAW transaction observations. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RawTransactionObservationInspectionQuery { + direction: crate::RawSortDirection, + network: crate::RawNetworkId, + page: crate::RawInspectionPageRequest, + transaction: std::option::Option, +} + +impl RawTransactionObservationInspectionQuery { + /// Creates one transaction-observation inspection query after enforcing network coherence. + pub fn try_new( + network: crate::RawNetworkId, + transaction: std::option::Option, + direction: crate::RawSortDirection, + page: crate::RawInspectionPageRequest, + ) -> crate::Result { + if let std::option::Option::Some(reference) = transaction.as_ref() { + if reference.network() != &network { + return std::result::Result::Err(raw_model_error("transaction")); + } + } + return std::result::Result::Ok(Self { direction, network, page, transaction }); + } + + /// Returns the requested deterministic traversal direction. + #[must_use] + pub const fn direction(&self) -> crate::RawSortDirection { + return self.direction; + } + + /// Returns the mandatory logical network scope. + #[must_use] + pub fn network(&self) -> &crate::RawNetworkId { + return &self.network; + } + + /// Returns the random-access inspection page request. + #[must_use] + pub const fn page(&self) -> crate::RawInspectionPageRequest { + return self.page; + } + + /// Returns the optional exact transaction identity used to filter observations. + #[must_use] + pub fn transaction(&self) -> std::option::Option<&crate::RawTransactionReference> { + return self.transaction.as_ref(); + } +} + +/// Backend-independent random-access inspection query for RAW account observations. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RawAccountObservationInspectionQuery { + account: std::option::Option, + direction: crate::RawSortDirection, + network: crate::RawNetworkId, + page: crate::RawInspectionPageRequest, +} + +impl RawAccountObservationInspectionQuery { + /// Creates one account-observation inspection query after enforcing network coherence. + pub fn try_new( + network: crate::RawNetworkId, + account: std::option::Option, + direction: crate::RawSortDirection, + page: crate::RawInspectionPageRequest, + ) -> crate::Result { + if let std::option::Option::Some(reference) = account.as_ref() { + if reference.network() != &network { + return std::result::Result::Err(raw_model_error("account")); + } + } + return std::result::Result::Ok(Self { account, direction, network, page }); + } + + /// Returns the optional exact account-state identity used to filter observations. + #[must_use] + pub fn account(&self) -> std::option::Option<&crate::RawAccountStateReference> { + return self.account.as_ref(); + } + + /// Returns the requested deterministic traversal direction. + #[must_use] + pub const fn direction(&self) -> crate::RawSortDirection { + return self.direction; + } + + /// Returns the mandatory logical network scope. + #[must_use] + pub fn network(&self) -> &crate::RawNetworkId { + return &self.network; + } + + /// Returns the random-access inspection page request. + #[must_use] + pub const fn page(&self) -> crate::RawInspectionPageRequest { + return self.page; + } +} + +/// Safe observation summary for one canonical RAW transaction acquisition. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RawTransactionObservationSummary { + observation_key: crate::RawObservationKey, + provenance: crate::RawAcquisitionProvenance, + transaction: crate::RawTransactionReference, +} + +impl RawTransactionObservationSummary { + /// Creates one safe transaction-observation summary without RAW payload bytes. + #[must_use] + pub fn new(observation_key: crate::RawObservationKey, transaction: crate::RawTransactionReference, provenance: crate::RawAcquisitionProvenance) -> Self { + return Self { observation_key, provenance, transaction }; + } + + /// Returns the producer-owned deterministic observation key. + #[must_use] + pub const fn observation_key(&self) -> crate::RawObservationKey { + return self.observation_key; + } + + /// Returns safe source-independent acquisition provenance. + #[must_use] + pub fn provenance(&self) -> &crate::RawAcquisitionProvenance { + return &self.provenance; + } + + /// Returns the durable transaction identity observed by this acquisition. + #[must_use] + pub fn transaction(&self) -> &crate::RawTransactionReference { + return &self.transaction; + } +} + +/// Safe observation summary for one canonical RAW account-state acquisition. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct RawAccountObservationSummary { + account: crate::RawAccountStateReference, + is_startup: std::option::Option, + observation_key: crate::RawObservationKey, + provenance: crate::RawAcquisitionProvenance, + transaction_signature: std::option::Option, + write_version: std::option::Option, +} + +impl RawAccountObservationSummary { + /// Creates one safe account-observation summary without account data bytes. + #[must_use] + pub fn new( + observation_key: crate::RawObservationKey, + account: crate::RawAccountStateReference, + provenance: crate::RawAcquisitionProvenance, + is_startup: std::option::Option, + transaction_signature: std::option::Option, + write_version: std::option::Option, + ) -> Self { + return Self { account, is_startup, observation_key, provenance, transaction_signature, write_version }; + } + + /// Returns the durable account-state identity observed by this acquisition. + #[must_use] + pub fn account(&self) -> &crate::RawAccountStateReference { + return &self.account; + } + + /// Returns the optional source-reported startup/replay marker. + #[must_use] + pub const fn is_startup(&self) -> std::option::Option { + return self.is_startup; + } + + /// Returns the producer-owned deterministic observation key. + #[must_use] + pub const fn observation_key(&self) -> crate::RawObservationKey { + return self.observation_key; + } + + /// Returns safe source-independent acquisition provenance. + #[must_use] + pub fn provenance(&self) -> &crate::RawAcquisitionProvenance { + return &self.provenance; + } + + /// Returns the optional transaction signature associated with the observed account write. + #[must_use] + pub const fn transaction_signature(&self) -> std::option::Option { + return self.transaction_signature; + } + + /// Returns the optional source-specific account write version. + #[must_use] + pub const fn write_version(&self) -> std::option::Option { + return self.write_version; + } +} + /// Backend-independent random-access inspection query for RAW transactions. #[derive(Clone, Debug, Eq, PartialEq)] pub struct RawTransactionInspectionQuery { diff --git a/crates/ksp-store-api/tests/dependency_boundary.rs b/crates/ksp-store-api/tests/dependency_boundary.rs index fd3ca6f..a2a8cc9 100644 --- a/crates/ksp-store-api/tests/dependency_boundary.rs +++ b/crates/ksp-store-api/tests/dependency_boundary.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/tests/dependency_boundary.rs -// version: 6 +// version: 7 //! Dependency canaries for the Store API RAW foundation. @@ -113,6 +113,7 @@ fn pre_006_source_boundary_keeps_models_and_capabilities_backend_free() { } assert!(raw_transaction_capability.contains("trait RawTransactionRead")); assert!(raw_transaction_capability.contains("trait RawTransactionInspectionRead")); + assert!(raw_transaction_capability.contains("trait RawTransactionObservationInspectionRead")); assert!(raw_transaction_capability.contains("trait RawTransactionWrite")); assert!(raw_transaction_capability.contains("trait RawTransactionObservationRead")); assert!(raw_transaction_capability.contains("trait RawTransactionObservationWrite")); @@ -120,6 +121,7 @@ fn pre_006_source_boundary_keeps_models_and_capabilities_backend_free() { assert!(raw_retention_capability.contains("trait RawTransactionRetentionWrite")); assert!(raw_account_capability.contains("trait RawAccountStateRead")); assert!(raw_account_capability.contains("trait RawAccountStateInspectionRead")); + assert!(raw_account_capability.contains("trait RawAccountObservationInspectionRead")); assert!(raw_account_capability.contains("trait RawAccountStateWrite")); assert!(raw_account_capability.contains("trait RawAccountObservationRead")); assert!(raw_account_capability.contains("trait RawAccountObservationWrite")); diff --git a/crates/ksp-store-api/tests/external_backend.rs b/crates/ksp-store-api/tests/external_backend.rs index adc38c8..7113ab4 100644 --- a/crates/ksp-store-api/tests/external_backend.rs +++ b/crates/ksp-store-api/tests/external_backend.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/tests/external_backend.rs -// version: 3 +// version: 4 //! External-implementation canary for object-safe Store API capabilities. @@ -58,6 +58,18 @@ impl ksp_store_api::RawTransactionWrite for ExternalMemoryBackend { } } +impl ksp_store_api::RawTransactionObservationInspectionRead for ExternalMemoryBackend { + fn inspect_raw_transaction_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawTransactionObservationInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + let _ = query; + return std::boxed::Box::pin(async { + return ksp_store_api::RawInspectionPage::try_new(std::vec::Vec::new(), 0, 0); + }); + } +} + impl ksp_store_api::RawTransactionObservationRead for ExternalMemoryBackend { fn get_raw_transaction_observation<'a>( &'a self, @@ -133,6 +145,18 @@ impl ksp_store_api::RawAccountStateWrite for ExternalMemoryBackend { } } +impl ksp_store_api::RawAccountObservationInspectionRead for ExternalMemoryBackend { + fn inspect_raw_account_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawAccountObservationInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + let _ = query; + return std::boxed::Box::pin(async { + return ksp_store_api::RawInspectionPage::try_new(std::vec::Vec::new(), 0, 0); + }); + } +} + impl ksp_store_api::RawAccountObservationRead for ExternalMemoryBackend { fn get_raw_account_observation<'a>( &'a self, @@ -197,11 +221,13 @@ fn v0_3_8_pre_003_external_backend_implements_canonical_and_inspection_capabilit let transaction_read: &dyn ksp_store_api::RawTransactionRead = &backend; let transaction_write: &dyn ksp_store_api::RawTransactionWrite = &backend; let transaction_inspection: &dyn ksp_store_api::RawTransactionInspectionRead = &backend; + let transaction_observation_inspection: &dyn ksp_store_api::RawTransactionObservationInspectionRead = &backend; let transaction_observation_read: &dyn ksp_store_api::RawTransactionObservationRead = &backend; let transaction_observation_write: &dyn ksp_store_api::RawTransactionObservationWrite = &backend; let account_read: &dyn ksp_store_api::RawAccountStateRead = &backend; let account_write: &dyn ksp_store_api::RawAccountStateWrite = &backend; let account_inspection: &dyn ksp_store_api::RawAccountStateInspectionRead = &backend; + let account_observation_inspection: &dyn ksp_store_api::RawAccountObservationInspectionRead = &backend; let account_observation_read: &dyn ksp_store_api::RawAccountObservationRead = &backend; let account_observation_write: &dyn ksp_store_api::RawAccountObservationWrite = &backend; let retention_read: &dyn ksp_store_api::RawTransactionRetentionRead = &backend; @@ -209,11 +235,13 @@ fn v0_3_8_pre_003_external_backend_implements_canonical_and_inspection_capabilit let _ = transaction_read; let _ = transaction_write; let _ = transaction_inspection; + let _ = transaction_observation_inspection; let _ = transaction_observation_read; let _ = transaction_observation_write; let _ = account_read; let _ = account_write; let _ = account_inspection; + let _ = account_observation_inspection; let _ = account_observation_read; let _ = account_observation_write; let _ = retention_read; diff --git a/crates/ksp-store-api/tests/public_api.rs b/crates/ksp-store-api/tests/public_api.rs index 2b9effc..45d90b8 100644 --- a/crates/ksp-store-api/tests/public_api.rs +++ b/crates/ksp-store-api/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/tests/public_api.rs -// version: 7 +// version: 8 //! Integration canaries for the public `ksp-store-api` surface. @@ -199,8 +199,12 @@ fn public_v0_3_8_pre_003_inspection_contracts_are_backend_neutral_and_dyn_compat let page = ksp_store_api::RawInspectionPage::::try_new(std::vec![1, 2], 10, 5); assert!(page.is_ok()); let transaction_inspection: std::option::Option<&dyn ksp_store_api::RawTransactionInspectionRead> = std::option::Option::None; + let transaction_observation_inspection: std::option::Option<&dyn ksp_store_api::RawTransactionObservationInspectionRead> = std::option::Option::None; let account_inspection: std::option::Option<&dyn ksp_store_api::RawAccountStateInspectionRead> = std::option::Option::None; + let account_observation_inspection: std::option::Option<&dyn ksp_store_api::RawAccountObservationInspectionRead> = std::option::Option::None; assert!(transaction_inspection.is_none()); + assert!(transaction_observation_inspection.is_none()); assert!(account_inspection.is_none()); + assert!(account_observation_inspection.is_none()); return; } diff --git a/crates/ksp-store-api/tests/release_completeness.rs b/crates/ksp-store-api/tests/release_completeness.rs index d409ee2..ea9fe32 100644 --- a/crates/ksp-store-api/tests/release_completeness.rs +++ b/crates/ksp-store-api/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/tests/release_completeness.rs -// version: 2 +// version: 3 //! Release-level boundary and completeness canaries for the backend-neutral Store API RAW surface. @@ -21,6 +21,7 @@ fn v0_3_8_pre_003_exact_crate_root_export_inventory_is_stable() { "pub use ksp_core_lib::Pubkey;", "pub use ksp_core_lib::Result;", "pub use self::capability::StoreApiFuture;", + "pub use self::capability::raw_account::RawAccountObservationInspectionRead;", "pub use self::capability::raw_account::RawAccountObservationRead;", "pub use self::capability::raw_account::RawAccountObservationWrite;", "pub use self::capability::raw_account::RawAccountStateRead;", @@ -28,6 +29,7 @@ fn v0_3_8_pre_003_exact_crate_root_export_inventory_is_stable() { "pub use self::capability::raw_account::RawAccountStateWrite;", "pub use self::capability::raw_retention::RawTransactionRetentionRead;", "pub use self::capability::raw_retention::RawTransactionRetentionWrite;", + "pub use self::capability::raw_transaction::RawTransactionObservationInspectionRead;", "pub use self::capability::raw_transaction::RawTransactionObservationRead;", "pub use self::capability::raw_transaction::RawTransactionObservationWrite;", "pub use self::capability::raw_transaction::RawTransactionRead;", @@ -42,11 +44,15 @@ fn v0_3_8_pre_003_exact_crate_root_export_inventory_is_stable() { "pub use self::model::raw_account::RawAccountObservation;", "pub use self::model::raw_account::RawAccountState;", "pub use self::model::raw_account::RawAccountStateReference;", + "pub use self::model::raw_inspection::RawAccountObservationInspectionQuery;", + "pub use self::model::raw_inspection::RawAccountObservationSummary;", "pub use self::model::raw_inspection::RawAccountStateInspectionQuery;", "pub use self::model::raw_inspection::RawAccountStateSummary;", "pub use self::model::raw_inspection::RawInspectionPage;", "pub use self::model::raw_inspection::RawInspectionPageRequest;", "pub use self::model::raw_inspection::RawTransactionInspectionQuery;", + "pub use self::model::raw_inspection::RawTransactionObservationInspectionQuery;", + "pub use self::model::raw_inspection::RawTransactionObservationSummary;", "pub use self::model::raw_inspection::RawTransactionSummary;", "pub use self::model::raw_outcome::RawAcquisitionWriteOutcome;", "pub use self::model::raw_outcome::RawEntityWriteOutcome;", @@ -202,11 +208,13 @@ fn v0_3_8_pre_003_capability_inventory_stays_fine_grained_without_runtime_facade } traits.sort_unstable(); let mut expected = std::vec![ + "pub trait RawAccountObservationInspectionRead: std::marker::Send + std::marker::Sync {", "pub trait RawAccountObservationRead: std::marker::Send + std::marker::Sync {", "pub trait RawAccountObservationWrite: std::marker::Send + std::marker::Sync {", "pub trait RawAccountStateRead: std::marker::Send + std::marker::Sync {", "pub trait RawAccountStateInspectionRead: std::marker::Send + std::marker::Sync {", "pub trait RawAccountStateWrite: std::marker::Send + std::marker::Sync {", + "pub trait RawTransactionObservationInspectionRead: std::marker::Send + std::marker::Sync {", "pub trait RawTransactionObservationRead: std::marker::Send + std::marker::Sync {", "pub trait RawTransactionObservationWrite: std::marker::Send + std::marker::Sync {", "pub trait RawTransactionRead: std::marker::Send + std::marker::Sync {", diff --git a/crates/ksp-store-api/unit_tests/model/raw_inspection.rs b/crates/ksp-store-api/unit_tests/model/raw_inspection.rs index f080222..a2e08fa 100644 --- a/crates/ksp-store-api/unit_tests/model/raw_inspection.rs +++ b/crates/ksp-store-api/unit_tests/model/raw_inspection.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-api/unit_tests/model/raw_inspection.rs -// version: 1 +// version: 2 //! Unit tests for backend-neutral RAW inspection contracts. @@ -154,3 +154,118 @@ fn account_summary_exposes_length_only_and_enforces_raw_account_bound() { assert!(crate::RawAccountStateSummary::try_new(reference, 1000, crate::Pubkey::new_from_array([0x43_u8; 32]), false, 9, 16_777_217,).is_err()); return; } + +#[test] +fn observation_inspection_queries_require_network_coherence_and_exact_entity_filters() { + let network = match network() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let other_network = match crate::RawNetworkId::new("devnet".to_owned()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let limit = match crate::RawPageLimit::new(25) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let page = crate::RawInspectionPageRequest::new(125, limit); + let transaction_reference = crate::RawTransactionReference::new(network.clone(), crate::RawTransactionSignature::new([0x51_u8; 64])); + let transaction = crate::RawTransactionObservationInspectionQuery::try_new( + network.clone(), + std::option::Option::Some(transaction_reference.clone()), + crate::RawSortDirection::Descending, + page, + ); + assert!(transaction.is_ok()); + let transaction = match transaction { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(transaction.page().offset(), 125); + assert_eq!(transaction.transaction(), std::option::Option::Some(&transaction_reference)); + assert!( + crate::RawTransactionObservationInspectionQuery::try_new( + other_network.clone(), + std::option::Option::Some(transaction_reference), + crate::RawSortDirection::Ascending, + page, + ) + .is_err() + ); + let account_reference = + crate::RawAccountStateReference::new(network.clone(), crate::Pubkey::new_from_array([0x52_u8; 32]), 99, crate::RawContentHash::new([0x53_u8; 32])); + let account = crate::RawAccountObservationInspectionQuery::try_new( + network, + std::option::Option::Some(account_reference.clone()), + crate::RawSortDirection::Ascending, + page, + ); + assert!(account.is_ok()); + let account = match account { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + assert_eq!(account.account(), std::option::Option::Some(&account_reference)); + let mismatched_account = crate::RawAccountStateReference::new( + other_network.clone(), + crate::Pubkey::new_from_array([0x54_u8; 32]), + 100, + crate::RawContentHash::new([0x55_u8; 32]), + ); + assert!( + crate::RawAccountObservationInspectionQuery::try_new( + other_network, + std::option::Option::Some(account_reference), + crate::RawSortDirection::Ascending, + page, + ) + .is_err() + ); + let _ = mismatched_account; + return; +} + +#[test] +fn observation_summaries_expose_safe_provenance_without_raw_entity_bytes() { + let network = match network() { + std::option::Option::Some(value) => value, + std::option::Option::None => return, + }; + let provider = match crate::RawProvenanceCode::new("publicnode".to_owned()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let protocol = match crate::RawProvenanceCode::new("solana-json-rpc".to_owned()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let method = match crate::RawProvenanceCode::new("getTransaction".to_owned()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let received_at = match crate::RawTimestamp::from_unix_millis(1_750_000_000_000) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return, + }; + let provenance = crate::RawAcquisitionProvenance::new(provider, protocol, method, crate::RawAcquisitionOrigin::Backfill, received_at); + let transaction_reference = crate::RawTransactionReference::new(network.clone(), crate::RawTransactionSignature::new([0x61_u8; 64])); + let transaction = + crate::RawTransactionObservationSummary::new(crate::RawObservationKey::new([0x62_u8; 32]), transaction_reference.clone(), provenance.clone()); + assert_eq!(transaction.transaction(), &transaction_reference); + assert_eq!(transaction.provenance().received_at(), received_at); + let account_reference = + crate::RawAccountStateReference::new(network, crate::Pubkey::new_from_array([0x63_u8; 32]), 700, crate::RawContentHash::new([0x64_u8; 32])); + let account = crate::RawAccountObservationSummary::new( + crate::RawObservationKey::new([0x65_u8; 32]), + account_reference.clone(), + provenance, + std::option::Option::Some(true), + std::option::Option::Some(crate::RawTransactionSignature::new([0x66_u8; 64])), + std::option::Option::Some(u64::MAX), + ); + assert_eq!(account.account(), &account_reference); + assert_eq!(account.is_startup(), std::option::Option::Some(true)); + assert_eq!(account.write_version(), std::option::Option::Some(u64::MAX)); + return; +} diff --git a/crates/ksp-store-lib/src/lib.rs b/crates/ksp-store-lib/src/lib.rs index ba7f9ac..6df93a3 100644 --- a/crates/ksp-store-lib/src/lib.rs +++ b/crates/ksp-store-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-lib/src/lib.rs -// version: 11 +// version: 12 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -121,8 +121,14 @@ pub use ksp_store_api::MAX_RAW_UNIX_MILLIS; pub use ksp_store_api::Pubkey; /// Persistable acquisition observation linked to one complete canonical RAW account state. pub use ksp_store_api::RawAccountObservation; +/// Backend-independent random-access inspection query for RAW account observations. +pub use ksp_store_api::RawAccountObservationInspectionQuery; +/// Read capability for random-access RAW account-observation inspection. +pub use ksp_store_api::RawAccountObservationInspectionRead; /// Read capability for persisted RAW account-state observations. pub use ksp_store_api::RawAccountObservationRead; +/// Safe observation summary for one canonical RAW account-state acquisition. +pub use ksp_store_api::RawAccountObservationSummary; /// Write capability for additional observations of already persisted RAW account states. pub use ksp_store_api::RawAccountObservationWrite; /// Canonical complete N1 RAW account state independent from acquisition transport. @@ -195,8 +201,14 @@ pub use ksp_store_api::RawTransactionInspectionQuery; pub use ksp_store_api::RawTransactionInspectionRead; /// Persistable acquisition observation linked to one canonical RAW transaction. pub use ksp_store_api::RawTransactionObservation; +/// Backend-independent random-access inspection query for RAW transaction observations. +pub use ksp_store_api::RawTransactionObservationInspectionQuery; +/// Read capability for random-access RAW transaction-observation inspection. +pub use ksp_store_api::RawTransactionObservationInspectionRead; /// Read capability for persisted RAW transaction observations. pub use ksp_store_api::RawTransactionObservationRead; +/// Safe observation summary for one canonical RAW transaction acquisition. +pub use ksp_store_api::RawTransactionObservationSummary; /// Write capability for additional observations of already persisted RAW transactions. pub use ksp_store_api::RawTransactionObservationWrite; /// Backend-independent list query for canonical RAW transactions. diff --git a/crates/ksp-store-lib/src/store.rs b/crates/ksp-store-lib/src/store.rs index 0ea8d0f..0fe1b16 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: 9 +// version: 10 /// Opaque common Store runtime facade. /// @@ -116,6 +116,36 @@ impl std::fmt::Debug for Store { } } +impl ksp_store_api::RawAccountObservationInspectionRead for Store { + fn inspect_raw_account_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawAccountObservationInspectionQuery, + ) -> 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_account_observations(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::RawAccountObservationRead for Store { fn get_raw_account_observation<'a>( &'a self, @@ -424,6 +454,36 @@ impl ksp_store_api::RawTransactionWrite for Store { } } +impl ksp_store_api::RawTransactionObservationInspectionRead for Store { + fn inspect_raw_transaction_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawTransactionObservationInspectionQuery, + ) -> 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_transaction_observations(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::RawTransactionObservationRead for Store { fn get_raw_transaction_observation<'a>( &'a self, diff --git a/crates/ksp-store-lib/tests/dependency_boundary.rs b/crates/ksp-store-lib/tests/dependency_boundary.rs index f065ecf..9103557 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: 10 +// version: 11 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -60,15 +60,17 @@ fn pre_005_facade_exposes_no_physical_postgres_types_or_environment_bypass() { } #[test] -fn v0_3_8_pre_005_facade_dispatches_twelve_raw_capabilities_without_physical_leak() { +fn v0_3_8_pre_009_facade_dispatches_fourteen_raw_capabilities_without_physical_leak() { let store = include_str!("../src/store.rs"); for required in [ + "impl ksp_store_api::RawAccountObservationInspectionRead for Store", "impl ksp_store_api::RawAccountObservationRead for Store", "impl ksp_store_api::RawAccountObservationWrite for Store", "impl ksp_store_api::RawAccountStateInspectionRead 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::RawTransactionObservationInspectionRead for Store", "impl ksp_store_api::RawTransactionObservationRead for Store", "impl ksp_store_api::RawTransactionObservationWrite for Store", "impl ksp_store_api::RawTransactionRead for Store", @@ -81,7 +83,7 @@ fn v0_3_8_pre_005_facade_dispatches_twelve_raw_capabilities_without_physical_lea ] { assert!(store.contains(required), "missing pre.008 Store capability dispatch contract: {required}"); } - assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 12); + assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 14); 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 6ef5023..7f25b32 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: 9 +// version: 10 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -155,7 +155,10 @@ fn pre_009_facade_modules_and_crate_root_exports_are_exact() { "PostgresTlsMode", "Pubkey", "RawAccountObservation", + "RawAccountObservationInspectionQuery", + "RawAccountObservationInspectionRead", "RawAccountObservationRead", + "RawAccountObservationSummary", "RawAccountObservationWrite", "RawAccountState", "RawAccountStateInspectionQuery", @@ -192,7 +195,10 @@ fn pre_009_facade_modules_and_crate_root_exports_are_exact() { "RawTransactionInspectionQuery", "RawTransactionInspectionRead", "RawTransactionObservation", + "RawTransactionObservationInspectionQuery", + "RawTransactionObservationInspectionRead", "RawTransactionObservationRead", + "RawTransactionObservationSummary", "RawTransactionObservationWrite", "RawTransactionQuery", "RawTransactionRead", @@ -216,7 +222,7 @@ fn pre_009_facade_modules_and_crate_root_exports_are_exact() { ]; expected.sort_unstable(); assert_eq!(actual.as_slice(), expected.as_slice()); - assert_eq!(actual.len(), 99); + assert_eq!(actual.len(), 105); return; } @@ -292,15 +298,17 @@ fn pre_009_facade_production_sources_keep_config_env_physical_sql_and_backend_ha } #[test] -fn v0_3_8_pre_005_facade_raw_capability_inventory_is_exactly_twelve() { +fn v0_3_8_pre_009_facade_raw_capability_inventory_is_exactly_fourteen() { let store = include_str!("../src/store.rs"); let capability_impls = [ + "impl ksp_store_api::RawAccountObservationInspectionRead for Store", "impl ksp_store_api::RawAccountObservationRead for Store", "impl ksp_store_api::RawAccountObservationWrite for Store", "impl ksp_store_api::RawAccountStateInspectionRead 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::RawTransactionObservationInspectionRead for Store", "impl ksp_store_api::RawTransactionObservationRead for Store", "impl ksp_store_api::RawTransactionObservationWrite for Store", "impl ksp_store_api::RawTransactionRead for Store", @@ -311,8 +319,8 @@ fn v0_3_8_pre_005_facade_raw_capability_inventory_is_exactly_twelve() { 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(), 12); - assert_eq!(store.matches("validate_operation_network(").count(), 16); + assert_eq!(store.matches("impl ksp_store_api::Raw").count(), 14); + assert_eq!(store.matches("validate_operation_network(").count(), 18); return; } @@ -322,11 +330,13 @@ fn v0_3_8_pre_005_facade_and_backend_capability_sets_match_with_both_inspection_ 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(), 12); - assert_eq!(backend_traits.len(), 12); + assert_eq!(store_traits.len(), 14); + assert_eq!(backend_traits.len(), 14); assert_eq!(store_traits, backend_traits); assert!(store_traits.contains(&"RawTransactionInspectionRead")); assert!(store_traits.contains(&"RawAccountStateInspectionRead")); + assert!(store_traits.contains(&"RawTransactionObservationInspectionRead")); + assert!(store_traits.contains(&"RawAccountObservationInspectionRead")); 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 13d647f..ba8e495 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: 11 +// version: 12 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -79,12 +79,14 @@ fn pre_007_health_and_runtime_snapshot_types_are_portable_crate_root_contracts() fn assert_raw_capabilities() where - T: ksp_store_lib::RawAccountObservationRead + T: ksp_store_lib::RawAccountObservationInspectionRead + + ksp_store_lib::RawAccountObservationRead + ksp_store_lib::RawAccountObservationWrite + ksp_store_lib::RawAccountStateInspectionRead + ksp_store_lib::RawAccountStateRead + ksp_store_lib::RawAccountStateWrite + ksp_store_lib::RawTransactionInspectionRead + + ksp_store_lib::RawTransactionObservationInspectionRead + ksp_store_lib::RawTransactionObservationRead + ksp_store_lib::RawTransactionObservationWrite + ksp_store_lib::RawTransactionRead @@ -97,7 +99,7 @@ where } #[test] -fn v0_3_8_pre_005_store_facade_implements_both_inspection_capabilities_as_12_of_12() { +fn v0_3_8_pre_009_store_facade_implements_entity_and_observation_inspection_capabilities_as_14_of_14() { assert_raw_capabilities::(); return; } diff --git a/crates/ksp-store-postgres-lib/src/lib.rs b/crates/ksp-store-postgres-lib/src/lib.rs index e677e91..401bd3e 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: 23 +// version: 24 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -36,6 +36,8 @@ //! `0.3.8-pre.004` adds payload-free random-access RawTransaction inspection; //! `0.3.8-pre.005` adds the corresponding data-free RawAccountState inspection //! with exact counts while preserving both keyset traversal families unchanged. +//! `0.3.8-pre.009` adds safe random-access observation inspection for both RAW +//! families without changing the physical schema or read-by-key contracts. //! //! 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 @@ -85,6 +87,8 @@ pub(crate) use self::raw_account::cursor::raw_account_physical_page_limit; pub(crate) use self::raw_account::get_raw_account_observation; /// Private RAW account state reader consumed by the physical backend runtime. pub(crate) use self::raw_account::get_raw_account_state; +/// Private safe RAW account-observation inspection reader consumed by the physical backend runtime. +pub(crate) use self::raw_account::inspect_raw_account_observations; /// Private data-free RAW account-state inspection reader consumed by the physical backend runtime. pub(crate) use self::raw_account::inspect_raw_account_states; /// Private RAW account-state list reader consumed by the physical backend runtime. @@ -107,6 +111,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 safe RAW transaction-observation inspection reader consumed by the physical backend runtime. +pub(crate) use self::raw_transaction::inspect_raw_transaction_observations; /// 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. diff --git a/crates/ksp-store-postgres-lib/src/raw_account.rs b/crates/ksp-store-postgres-lib/src/raw_account.rs index 44cbcf6..667525d 100644 --- a/crates/ksp-store-postgres-lib/src/raw_account.rs +++ b/crates/ksp-store-postgres-lib/src/raw_account.rs @@ -1,5 +1,5 @@ // file: crates/ksp-store-postgres-lib/src/raw_account.rs -// version: 6 +// version: 7 pub(crate) mod cursor; @@ -7,6 +7,8 @@ const GET_ACCOUNT_OBSERVATION_SQL: &str = "SELECT observation_key, account_pubke const GET_ACCOUNT_STATE_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash, lamports::text AS lamports_text, owner, executable, rent_epoch::text AS rent_epoch_text, data FROM ksp_raw_account_states WHERE pubkey = $1 AND slot = $2::TEXT::NUMERIC AND state_hash = $3"; const INSERT_ACCOUNT_OBSERVATION_SQL: &str = "INSERT INTO ksp_raw_account_observations (observation_key, account_pubkey, account_slot, account_state_hash, 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, is_startup, transaction_signature, write_version) VALUES ($1, $2, $3::TEXT::NUMERIC, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19::TEXT::NUMERIC) ON CONFLICT (observation_key) DO NOTHING RETURNING observation_key"; const INSERT_ACCOUNT_STATE_SQL: &str = "INSERT INTO ksp_raw_account_states (pubkey, slot, state_hash, lamports, owner, executable, rent_epoch, data) VALUES ($1, $2::TEXT::NUMERIC, $3, $4::TEXT::NUMERIC, $5, $6, $7::TEXT::NUMERIC, $8) ON CONFLICT (pubkey, slot, state_hash) DO NOTHING RETURNING pubkey"; +const INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3))), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_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.account_pubkey, page.account_slot_text, page.account_state_hash, 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, page.is_startup, page.transaction_signature, page.write_version_text FROM counts LEFT JOIN LATERAL (SELECT observation_key, account_pubkey, account_slot::TEXT AS account_slot_text, account_state_hash, 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, is_startup, transaction_signature, write_version::TEXT AS write_version_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3)) ORDER BY received_at_unix_millis ASC, observation_key ASC LIMIT $4 OFFSET $5) AS page ON TRUE"; +const INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3))), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_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.account_pubkey, page.account_slot_text, page.account_state_hash, 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, page.is_startup, page.transaction_signature, page.write_version_text FROM counts LEFT JOIN LATERAL (SELECT observation_key, account_pubkey, account_slot::TEXT AS account_slot_text, account_state_hash, 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, is_startup, transaction_signature, write_version::TEXT AS write_version_text FROM ksp_raw_account_observations WHERE ($1::BYTEA IS NULL OR (account_pubkey = $1 AND account_slot = $2::TEXT::NUMERIC AND account_state_hash = $3)) ORDER BY received_at_unix_millis DESC, observation_key DESC LIMIT $4 OFFSET $5) AS page ON TRUE"; const INSPECT_ACCOUNT_STATES_ASC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_states WHERE ($1::BYTEA IS NULL OR pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL AND $2::TEXT IS NULL AND $3::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_states) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.pubkey IS NOT NULL AS page_present, page.pubkey, page.slot_text, page.state_hash, page.lamports_text, page.owner, page.executable, page.rent_epoch_text, page.data_length_bytes FROM counts LEFT JOIN LATERAL (SELECT account_row.pubkey, account_row.slot::TEXT AS slot_text, account_row.state_hash, account_row.lamports::TEXT AS lamports_text, account_row.owner, account_row.executable, account_row.rent_epoch::TEXT AS rent_epoch_text, OCTET_LENGTH(account_row.data)::BIGINT AS data_length_bytes FROM ksp_raw_account_states AS account_row WHERE ($1::BYTEA IS NULL OR account_row.pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR account_row.slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR account_row.slot <= $3::TEXT::NUMERIC) ORDER BY account_row.slot ASC, account_row.pubkey ASC, account_row.state_hash ASC LIMIT $4 OFFSET $5) AS page ON TRUE"; const INSPECT_ACCOUNT_STATES_DESC_SQL: &str = "WITH filtered_count AS (SELECT COUNT(*)::TEXT AS filtered_count_text FROM ksp_raw_account_states WHERE ($1::BYTEA IS NULL OR pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR slot <= $3::TEXT::NUMERIC)), counts AS (SELECT filtered_count_text, CASE WHEN $1::BYTEA IS NULL AND $2::TEXT IS NULL AND $3::TEXT IS NULL THEN filtered_count_text ELSE (SELECT COUNT(*)::TEXT FROM ksp_raw_account_states) END AS total_count_text FROM filtered_count) SELECT counts.total_count_text, counts.filtered_count_text, page.pubkey IS NOT NULL AS page_present, page.pubkey, page.slot_text, page.state_hash, page.lamports_text, page.owner, page.executable, page.rent_epoch_text, page.data_length_bytes FROM counts LEFT JOIN LATERAL (SELECT account_row.pubkey, account_row.slot::TEXT AS slot_text, account_row.state_hash, account_row.lamports::TEXT AS lamports_text, account_row.owner, account_row.executable, account_row.rent_epoch::TEXT AS rent_epoch_text, OCTET_LENGTH(account_row.data)::BIGINT AS data_length_bytes FROM ksp_raw_account_states AS account_row WHERE ($1::BYTEA IS NULL OR account_row.pubkey = $1::BYTEA) AND ($2::TEXT IS NULL OR account_row.slot >= $2::TEXT::NUMERIC) AND ($3::TEXT IS NULL OR account_row.slot <= $3::TEXT::NUMERIC) ORDER BY account_row.slot DESC, account_row.pubkey DESC, account_row.state_hash DESC LIMIT $4 OFFSET $5) AS page ON TRUE"; const LIST_ACCOUNT_STATES_ASC_SQL: &str = "SELECT pubkey, slot::text AS slot_text, state_hash FROM ksp_raw_account_states WHERE ($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, pubkey, state_hash) > ($3::TEXT::NUMERIC, $4::BYTEA, $5::BYTEA)) ORDER BY slot ASC, pubkey ASC, state_hash ASC LIMIT $6"; @@ -38,6 +40,31 @@ struct RawAccountInspectionDbRow { total_count_text: std::string::String, } +struct RawAccountObservationInspectionDbRow { + account_pubkey: std::option::Option>, + account_slot_text: std::option::Option, + account_state_hash: std::option::Option>, + acquisition_method: std::option::Option, + capture_session_id: std::option::Option, + commitment: std::option::Option, + endpoint_id: std::option::Option, + filter_id: std::option::Option, + filtered_count_text: std::string::String, + is_startup: std::option::Option, + observation_key: std::option::Option>, + observed_at_unix_millis: std::option::Option, + origin: std::option::Option, + page_present: bool, + protocol: std::option::Option, + provider: std::option::Option, + received_at_unix_millis: std::option::Option, + source_payload_hash: std::option::Option>, + source_payload_size_bytes: std::option::Option, + total_count_text: std::string::String, + transaction_signature: std::option::Option>, + write_version_text: std::option::Option, +} + struct RawAccountObservationDbRow { account_pubkey: std::vec::Vec, account_slot_text: std::string::String, @@ -317,6 +344,121 @@ pub(crate) async fn inspect_raw_account_states( }; } +/// Inspects one safe random-access RAW account-observation window with exact counts. +pub(crate) async fn inspect_raw_account_observations( + pool: &deadpool_postgres::Pool, + network: &ksp_store_api::RawNetworkId, + query: &ksp_store_api::RawAccountObservationInspectionQuery, +) -> std::result::Result, crate::PostgresBackendError> { + if query.network() != network { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::WrongNetwork, + "raw_account_observation_inspection_network", + )); + } + let sql = match query.direction() { + ksp_store_api::RawSortDirection::Ascending => INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL, + ksp_store_api::RawSortDirection::Descending => INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL, + _ => { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::QueryInvalid, + "raw_account_observation_inspection_direction", + )); + }, + }; + let (sql_limit, sql_offset) = match raw_account_inspection_sql_window(query.page()) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let account_pubkey = query.account().map(|reference| return reference.pubkey().to_bytes().to_vec()); + let account_slot_text = query.account().map(|reference| return reference.slot().to_string()); + let account_state_hash = query.account().map(|reference| return reference.state_hash().as_bytes().to_vec()); + 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 rows_result = client.query(sql, &[&account_pubkey, &account_slot_text, &account_state_hash, &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_account_observation_inspection_query", + )); + }, + }; + if rows.is_empty() { + return std::result::Result::Err(data_invalid("raw_account_observation_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_account_observation_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_account_observation_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_account_observation_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_account_observation_inspection_counts")), + } + if physical.page_present { + let observation = match decode_raw_account_observation_inspection_row(network, physical) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if let std::option::Option::Some(reference) = query.account() { + if observation.account() != reference { + return std::result::Result::Err(data_invalid("raw_account_observation_inspection_reference")); + } + } + items.push(ksp_store_api::RawAccountObservationSummary::new( + observation.observation_key(), + observation.account().clone(), + observation.provenance().clone(), + observation.is_startup(), + observation.transaction_signature(), + observation.write_version(), + )); + } else if row_count != 1 || !raw_account_observation_inspection_empty_page_is_clean(&physical) { + return std::result::Result::Err(data_invalid("raw_account_observation_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_account_observation_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_account_observation_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_account_observation_inspection_page")), + }; + if item_count > query.page().limit().get() { + return std::result::Result::Err(data_invalid("raw_account_observation_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_account_observation_inspection_page")), + }; +} + /// Reads one persisted RAW account observation by producer-owned idempotence key. pub(crate) async fn get_raw_account_observation( pool: &deadpool_postgres::Pool, @@ -495,6 +637,141 @@ fn decode_raw_account_list_row( return std::result::Result::Ok((slot, reference)); } +fn raw_account_observation_inspection_empty_page_is_clean(row: &RawAccountObservationInspectionDbRow) -> bool { + return row.account_pubkey.is_none() + && row.account_slot_text.is_none() + && row.account_state_hash.is_none() + && row.acquisition_method.is_none() + && row.capture_session_id.is_none() + && row.commitment.is_none() + && row.endpoint_id.is_none() + && row.filter_id.is_none() + && row.is_startup.is_none() + && row.observation_key.is_none() + && row.observed_at_unix_millis.is_none() + && row.origin.is_none() + && row.protocol.is_none() + && row.provider.is_none() + && row.received_at_unix_millis.is_none() + && row.source_payload_hash.is_none() + && row.source_payload_size_bytes.is_none() + && row.transaction_signature.is_none() + && row.write_version_text.is_none(); +} + +fn raw_account_observation_inspection_db_row( + row: &tokio_postgres::Row, +) -> std::result::Result { + macro_rules! get_optional { + ($name:literal, $ty:ty) => { + match row.try_get::<_, std::option::Option<$ty>>($name) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_account_observation_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_account_observation_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_account_observation_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_account_observation_inspection_decode")), + }; + return std::result::Result::Ok(RawAccountObservationInspectionDbRow { + account_pubkey: get_optional!("account_pubkey", std::vec::Vec), + account_slot_text: get_optional!("account_slot_text", std::string::String), + account_state_hash: get_optional!("account_state_hash", std::vec::Vec), + acquisition_method: get_optional!("acquisition_method", std::string::String), + capture_session_id: get_optional!("capture_session_id", std::string::String), + commitment: get_optional!("commitment", std::string::String), + endpoint_id: get_optional!("endpoint_id", std::string::String), + filter_id: get_optional!("filter_id", std::string::String), + filtered_count_text, + is_startup: get_optional!("is_startup", bool), + observation_key: get_optional!("observation_key", std::vec::Vec), + observed_at_unix_millis: get_optional!("observed_at_unix_millis", i64), + origin: get_optional!("origin", std::string::String), + page_present, + protocol: get_optional!("protocol", std::string::String), + provider: get_optional!("provider", std::string::String), + received_at_unix_millis: get_optional!("received_at_unix_millis", i64), + source_payload_hash: get_optional!("source_payload_hash", std::vec::Vec), + source_payload_size_bytes: get_optional!("source_payload_size_bytes", i64), + total_count_text, + transaction_signature: get_optional!("transaction_signature", std::vec::Vec), + write_version_text: get_optional!("write_version_text", std::string::String), + }); +} + +fn decode_raw_account_observation_inspection_row( + network: &ksp_store_api::RawNetworkId, + row: RawAccountObservationInspectionDbRow, +) -> std::result::Result { + let account_pubkey = match inspection_required(row.account_pubkey, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let account_slot_text = match inspection_required(row.account_slot_text, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let account_state_hash = match inspection_required(row.account_state_hash, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let acquisition_method = match inspection_required(row.acquisition_method, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let observation_key = match inspection_required(row.observation_key, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let origin = match inspection_required(row.origin, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let protocol = match inspection_required(row.protocol, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let provider = match inspection_required(row.provider, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let received_at_unix_millis = match inspection_required(row.received_at_unix_millis, "raw_account_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let physical = RawAccountObservationDbRow { + account_pubkey, + account_slot_text, + account_state_hash, + acquisition_method, + capture_session_id: row.capture_session_id, + commitment: row.commitment, + endpoint_id: row.endpoint_id, + filter_id: row.filter_id, + is_startup: row.is_startup, + observation_key, + observed_at_unix_millis: row.observed_at_unix_millis, + origin, + protocol, + provider, + received_at_unix_millis, + source_payload_hash: row.source_payload_hash, + source_payload_size_bytes: row.source_payload_size_bytes, + transaction_signature: row.transaction_signature, + write_version_text: row.write_version_text, + }; + return decode_raw_account_observation_row(network, physical); +} + fn raw_account_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result { let account_pubkey: std::vec::Vec = match row.try_get("account_pubkey") { std::result::Result::Ok(value) => value, diff --git a/crates/ksp-store-postgres-lib/src/raw_transaction.rs b/crates/ksp-store-postgres-lib/src/raw_transaction.rs index a5653f8..1e7c52e 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: 5 +// version: 6 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_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"; 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"; @@ -48,6 +50,26 @@ struct RawInspectionDbRow { total_count_text: std::string::String, } +struct RawObservationInspectionDbRow { + acquisition_method: std::option::Option, + capture_session_id: std::option::Option, + commitment: std::option::Option, + endpoint_id: std::option::Option, + filter_id: std::option::Option, + filtered_count_text: std::string::String, + observation_key: std::option::Option>, + observed_at_unix_millis: std::option::Option, + origin: std::option::Option, + page_present: bool, + protocol: std::option::Option, + provider: std::option::Option, + received_at_unix_millis: std::option::Option, + source_payload_hash: std::option::Option>, + source_payload_size_bytes: std::option::Option, + total_count_text: std::string::String, + transaction_signature: std::option::Option>, +} + struct RawObservationDbRow { acquisition_method: std::string::String, capture_session_id: std::option::Option, @@ -307,6 +329,116 @@ pub(crate) async fn inspect_raw_transactions( }; } +/// Inspects one safe random-access RAW transaction-observation window with exact counts. +pub(crate) async fn inspect_raw_transaction_observations( + pool: &deadpool_postgres::Pool, + network: &ksp_store_api::RawNetworkId, + query: &ksp_store_api::RawTransactionObservationInspectionQuery, +) -> std::result::Result, crate::PostgresBackendError> { + if query.network() != network { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::WrongNetwork, + "raw_transaction_observation_inspection_network", + )); + } + let sql = match query.direction() { + ksp_store_api::RawSortDirection::Ascending => INSPECT_OBSERVATIONS_ASC_SQL, + ksp_store_api::RawSortDirection::Descending => INSPECT_OBSERVATIONS_DESC_SQL, + _ => { + return std::result::Result::Err(crate::PostgresBackendError::new( + crate::PostgresBackendErrorKind::QueryInvalid, + "raw_transaction_observation_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 signature = query.transaction().map(|reference| return reference.signature().as_bytes().to_vec()); + 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 rows_result = client.query(sql, &[&signature, &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_observation_inspection_query", + )); + }, + }; + if rows.is_empty() { + return std::result::Result::Err(data_invalid("raw_transaction_observation_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_observation_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_observation_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_observation_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_observation_inspection_counts")), + } + if physical.page_present { + let observation = match decode_raw_observation_inspection_row(network, physical) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + if let std::option::Option::Some(reference) = query.transaction() { + if observation.transaction() != reference { + return std::result::Result::Err(data_invalid("raw_transaction_observation_inspection_reference")); + } + } + items.push(ksp_store_api::RawTransactionObservationSummary::new( + observation.observation_key(), + observation.transaction().clone(), + observation.provenance().clone(), + )); + } else if row_count != 1 || !raw_observation_inspection_empty_page_is_clean(&physical) { + return std::result::Result::Err(data_invalid("raw_transaction_observation_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_observation_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_observation_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_observation_inspection_page")), + }; + if item_count > query.page().limit().get() { + return std::result::Result::Err(data_invalid("raw_transaction_observation_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_observation_inspection_page")), + }; +} + /// Reads one RAW transaction observation from the physical PostgreSQL backend. pub(crate) async fn get_raw_transaction_observation( pool: &deadpool_postgres::Pool, @@ -1401,6 +1533,116 @@ fn raw_transaction_db_row(row: &tokio_postgres::Row) -> std::result::Result bool { + return row.acquisition_method.is_none() + && row.capture_session_id.is_none() + && row.commitment.is_none() + && row.endpoint_id.is_none() + && row.filter_id.is_none() + && row.observation_key.is_none() + && row.observed_at_unix_millis.is_none() + && row.origin.is_none() + && row.protocol.is_none() + && row.provider.is_none() + && row.received_at_unix_millis.is_none() + && row.source_payload_hash.is_none() + && row.source_payload_size_bytes.is_none() + && row.transaction_signature.is_none(); +} + +fn raw_observation_inspection_db_row(row: &tokio_postgres::Row) -> std::result::Result { + macro_rules! get_optional { + ($name:literal, $ty:ty) => { + match row.try_get::<_, std::option::Option<$ty>>($name) { + std::result::Result::Ok(value) => value, + std::result::Result::Err(_) => return std::result::Result::Err(data_invalid("raw_transaction_observation_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_observation_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_observation_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_observation_inspection_decode")), + }; + return std::result::Result::Ok(RawObservationInspectionDbRow { + acquisition_method: get_optional!("acquisition_method", std::string::String), + capture_session_id: get_optional!("capture_session_id", std::string::String), + commitment: get_optional!("commitment", std::string::String), + endpoint_id: get_optional!("endpoint_id", std::string::String), + filter_id: get_optional!("filter_id", std::string::String), + filtered_count_text, + observation_key: get_optional!("observation_key", std::vec::Vec), + observed_at_unix_millis: get_optional!("observed_at_unix_millis", i64), + origin: get_optional!("origin", std::string::String), + page_present, + protocol: get_optional!("protocol", std::string::String), + provider: get_optional!("provider", std::string::String), + received_at_unix_millis: get_optional!("received_at_unix_millis", i64), + source_payload_hash: get_optional!("source_payload_hash", std::vec::Vec), + source_payload_size_bytes: get_optional!("source_payload_size_bytes", i64), + total_count_text, + transaction_signature: get_optional!("transaction_signature", std::vec::Vec), + }); +} + +fn decode_raw_observation_inspection_row( + network: &ksp_store_api::RawNetworkId, + row: RawObservationInspectionDbRow, +) -> std::result::Result { + let observation_key = match inspection_required(row.observation_key, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let received_at_unix_millis = match inspection_required(row.received_at_unix_millis, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let transaction_signature = match inspection_required(row.transaction_signature, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let acquisition_method = match inspection_required(row.acquisition_method, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let origin = match inspection_required(row.origin, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let protocol = match inspection_required(row.protocol, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let provider = match inspection_required(row.provider, "raw_transaction_observation_inspection_shape") { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let physical = RawObservationDbRow { + acquisition_method, + capture_session_id: row.capture_session_id, + commitment: row.commitment, + endpoint_id: row.endpoint_id, + filter_id: row.filter_id, + observation_key, + observed_at_unix_millis: row.observed_at_unix_millis, + origin, + protocol, + provider, + received_at_unix_millis, + source_payload_hash: row.source_payload_hash, + source_payload_size_bytes: row.source_payload_size_bytes, + transaction_signature, + }; + return decode_raw_observation_row(network, physical); +} + fn raw_observation_db_row(row: &tokio_postgres::Row) -> std::result::Result { macro_rules! required { ($name:literal, $ty:ty) => { diff --git a/crates/ksp-store-postgres-lib/src/runtime.rs b/crates/ksp-store-postgres-lib/src/runtime.rs index eb91952..a7a8469 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: 17 +// version: 18 const APPLICATION_NAME: &str = "ksp-store"; const MAX_CONNECTION_URI_BYTES: usize = 4_096; @@ -356,6 +356,14 @@ impl PostgresBackend { return crate::inspect_raw_account_states(&self.pool, &self.network, query).await; } + /// Inspects one safe random-access RAW account-observation window with exact logical counts. + pub async fn inspect_raw_account_observations( + &self, + query: &ksp_store_api::RawAccountObservationInspectionQuery, + ) -> std::result::Result, crate::PostgresBackendError> { + return crate::inspect_raw_account_observations(&self.pool, &self.network, query).await; + } + /// Reads one persisted RAW account observation by producer-owned idempotence key. pub async fn get_raw_account_observation( &self, @@ -405,6 +413,14 @@ impl PostgresBackend { return crate::inspect_raw_transactions(&self.pool, &self.network, query).await; } + /// Inspects one safe random-access RAW transaction-observation window with exact logical counts. + pub async fn inspect_raw_transaction_observations( + &self, + query: &ksp_store_api::RawTransactionObservationInspectionQuery, + ) -> std::result::Result, crate::PostgresBackendError> { + return crate::inspect_raw_transaction_observations(&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, @@ -485,6 +501,18 @@ impl std::fmt::Debug for PostgresBackend { } } +impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend { + fn inspect_raw_account_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawAccountObservationInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + return std::boxed::Box::pin(async move { + let result = PostgresBackend::inspect_raw_account_observations(self, query).await; + return result.map_err(map_capability_error); + }); + } +} + impl ksp_store_api::RawAccountObservationRead for PostgresBackend { fn get_raw_account_observation<'a>( &'a self, @@ -604,6 +632,18 @@ impl ksp_store_api::RawTransactionWrite for PostgresBackend { } } +impl ksp_store_api::RawTransactionObservationInspectionRead for PostgresBackend { + fn inspect_raw_transaction_observations<'a>( + &'a self, + query: &'a ksp_store_api::RawTransactionObservationInspectionQuery, + ) -> ksp_store_api::StoreApiFuture<'a, ksp_store_api::Result>> { + return std::boxed::Box::pin(async move { + let result = PostgresBackend::inspect_raw_transaction_observations(self, query).await; + return result.map_err(map_capability_error); + }); + } +} + impl ksp_store_api::RawTransactionObservationRead for PostgresBackend { fn get_raw_transaction_observation<'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 23a2d72..37bbc13 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: 28 +// version: 29 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -187,6 +187,7 @@ fn pre_005_raw_account_acquisition_is_atomic_idempotent_and_keeps_trait_impls_ou "impl ksp_store_api::RawAccountStateInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountStateRead for PostgresBackend", "impl ksp_store_api::RawAccountStateWrite for PostgresBackend", + "impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", ] { @@ -222,6 +223,7 @@ fn pre_006_raw_account_additional_observation_is_reference_guarded_cancellation_ "impl ksp_store_api::RawAccountStateInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountStateRead for PostgresBackend", "impl ksp_store_api::RawAccountStateWrite for PostgresBackend", + "impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", ] { @@ -408,6 +410,7 @@ fn pre_007_raw_account_pagination_is_keyset_cursor_bound_and_policy_free() { "impl ksp_store_api::RawAccountStateInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountStateRead for PostgresBackend", "impl ksp_store_api::RawAccountStateWrite for PostgresBackend", + "impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", ] { @@ -449,15 +452,17 @@ fn pre_007_raw_retention_is_atomic_compare_and_transition_without_fake_compactio } #[test] -fn v0_3_8_pre_005_backend_trait_implementations_cover_exact_twelve_raw_capabilities_in_runtime_bridge() { +fn v0_3_8_pre_009_backend_trait_implementations_cover_exact_fourteen_raw_capabilities_in_runtime_bridge() { let runtime = include_str!("../src/runtime.rs"); for implementation in [ + "impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", "impl ksp_store_api::RawAccountStateInspectionRead 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::RawTransactionObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationWrite for PostgresBackend", "impl ksp_store_api::RawTransactionRead for PostgresBackend", @@ -467,6 +472,6 @@ fn v0_3_8_pre_005_backend_trait_implementations_cover_exact_twelve_raw_capabilit ] { assert_eq!(runtime.matches(implementation).count(), 1, "unexpected PostgreSQL RAW capability inventory: {implementation}"); } - assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 12); + assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 14); return; } diff --git a/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs b/crates/ksp-store-postgres-lib/tests/hardening_completeness.rs index e54ee4c..16582a5 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: 21 +// version: 22 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -302,15 +302,17 @@ fn pre_009_live_raw_transaction_proof_is_opt_in_isolated_and_secret_safe() { } #[test] -fn v0_3_8_pre_005_raw_capability_implementation_inventory_is_exactly_twelve() { +fn v0_3_8_pre_009_raw_capability_implementation_inventory_is_exactly_fourteen() { let runtime = include_str!("../src/runtime.rs"); let capability_impls = [ + "impl ksp_store_api::RawAccountObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationRead for PostgresBackend", "impl ksp_store_api::RawAccountObservationWrite for PostgresBackend", "impl ksp_store_api::RawAccountStateInspectionRead 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::RawTransactionObservationInspectionRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationRead for PostgresBackend", "impl ksp_store_api::RawTransactionObservationWrite for PostgresBackend", "impl ksp_store_api::RawTransactionRead for PostgresBackend", @@ -321,7 +323,7 @@ fn v0_3_8_pre_005_raw_capability_implementation_inventory_is_exactly_twelve() { 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(), 12); + assert_eq!(runtime.matches("impl ksp_store_api::Raw").count(), 14); let migration = include_str!("../src/migration.rs"); assert!(migration.contains("raw_account_state")); assert!(migration.contains("crate::V002_RESOURCES")); @@ -394,13 +396,77 @@ fn v0_3_8_pre_004_transaction_inspection_sql_is_single_statement_payload_free_co } 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 v0_3_8_pre_009_observation_inspection_sql_is_counted_random_access_safe_and_family_local() { + let transaction = include_str!("../src/raw_transaction.rs"); + let account = include_str!("../src/raw_account.rs"); + let transaction_ascending = transaction.lines().find(|line| return line.starts_with("const INSPECT_OBSERVATIONS_ASC_SQL")); + let transaction_ascending = match transaction_ascending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing ascending transaction-observation inspection SQL"), + }; + let transaction_descending = transaction.lines().find(|line| return line.starts_with("const INSPECT_OBSERVATIONS_DESC_SQL")); + let transaction_descending = match transaction_descending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing descending transaction-observation inspection SQL"), + }; + for statement in [transaction_ascending, transaction_descending] { + for required in [ + "COUNT(*)::TEXT AS filtered_count_text", + "COUNT(*)::TEXT FROM ksp_raw_transaction_observations", + "LEFT JOIN LATERAL", + "transaction_signature = $1", + "LIMIT $2 OFFSET $3", + "received_at_unix_millis", + "observation_key", + ] { + assert!(statement.contains(required), "missing transaction-observation inspection SQL contract: {required}"); + } + for forbidden in ["ksp_raw_transactions AS", "payload", "archive_payload", "SELECT *"] { + assert!(!statement.contains(forbidden), "transaction-observation inspection leaked unrelated/raw material: {forbidden}"); + } + } + assert!(transaction_ascending.contains("ORDER BY received_at_unix_millis ASC, observation_key ASC")); + assert!(transaction_descending.contains("ORDER BY received_at_unix_millis DESC, observation_key DESC")); + let account_ascending = account.lines().find(|line| return line.starts_with("const INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL")); + let account_ascending = match account_ascending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing ascending account-observation inspection SQL"), + }; + let account_descending = account.lines().find(|line| return line.starts_with("const INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL")); + let account_descending = match account_descending { + std::option::Option::Some(value) => value, + std::option::Option::None => panic!("missing descending account-observation inspection SQL"), + }; + for statement in [account_ascending, account_descending] { + for required in [ + "COUNT(*)::TEXT AS filtered_count_text", + "COUNT(*)::TEXT FROM ksp_raw_account_observations", + "LEFT JOIN LATERAL", + "account_pubkey = $1", + "account_slot = $2::TEXT::NUMERIC", + "account_state_hash = $3", + "LIMIT $4 OFFSET $5", + "received_at_unix_millis", + "observation_key", + ] { + assert!(statement.contains(required), "missing account-observation inspection SQL contract: {required}"); + } + for forbidden in ["ksp_raw_account_states AS", "account_row.data", "SELECT *"] { + assert!(!statement.contains(forbidden), "account-observation inspection leaked unrelated/raw material: {forbidden}"); + } + } + assert!(account_ascending.contains("ORDER BY received_at_unix_millis ASC, observation_key ASC")); + assert!(account_descending.contains("ORDER BY received_at_unix_millis DESC, observation_key DESC")); + return; +} + #[test] fn pre_010_raw_account_private_sql_is_non_destructive_keyset_and_family_local() { let source = include_str!("../src/raw_account.rs"); @@ -493,6 +559,5 @@ fn v0_3_8_pre_005_raw_account_keyset_sql_remains_offset_free_while_inspection_is assert!(statement.contains("account_row.pubkey")); assert!(statement.contains("account_row.state_hash")); } - assert_eq!(source.matches(" OFFSET ").count(), 2); return; } diff --git a/crates/ksp-store-postgres-lib/tests/public_api.rs b/crates/ksp-store-postgres-lib/tests/public_api.rs index 0fec7e0..1cafa77 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: 16 +// version: 17 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -120,8 +120,10 @@ fn pre_005_raw_write_bridge_uses_only_backend_independent_models_and_outcomes() fn v0_3_8_pre_005_raw_list_and_inspection_bridges_use_only_backend_independent_models() { let _account_list = ksp_store_postgres_lib::PostgresBackend::list_raw_account_states; let _account_inspection = ksp_store_postgres_lib::PostgresBackend::inspect_raw_account_states; + let _account_observation_inspection = ksp_store_postgres_lib::PostgresBackend::inspect_raw_account_observations; let _transaction_list = ksp_store_postgres_lib::PostgresBackend::list_raw_transactions; let _transaction_inspection = ksp_store_postgres_lib::PostgresBackend::inspect_raw_transactions; + let _transaction_observation_inspection = ksp_store_postgres_lib::PostgresBackend::inspect_raw_transaction_observations; return; } @@ -133,12 +135,14 @@ fn pre_007_raw_retention_write_bridge_uses_backend_independent_transition_and_ou fn assert_raw_capabilities() where - T: ksp_store_api::RawAccountObservationRead + T: ksp_store_api::RawAccountObservationInspectionRead + + ksp_store_api::RawAccountObservationRead + ksp_store_api::RawAccountObservationWrite + ksp_store_api::RawAccountStateInspectionRead + ksp_store_api::RawAccountStateRead + ksp_store_api::RawAccountStateWrite + ksp_store_api::RawTransactionInspectionRead + + ksp_store_api::RawTransactionObservationInspectionRead + ksp_store_api::RawTransactionObservationRead + ksp_store_api::RawTransactionObservationWrite + ksp_store_api::RawTransactionRead @@ -151,7 +155,7 @@ where } #[test] -fn v0_3_8_pre_005_postgres_backend_implements_both_inspection_capabilities_as_12_of_12() { +fn v0_3_8_pre_009_postgres_backend_implements_entity_and_observation_inspection_capabilities_as_14_of_14() { assert_raw_capabilities::(); return; } diff --git a/deltas/0.3.8/pre.009.md b/deltas/0.3.8/pre.009.md new file mode 100644 index 0000000..c7054ad --- /dev/null +++ b/deltas/0.3.8/pre.009.md @@ -0,0 +1,114 @@ + + + +# Delta `0.3.8-pre.009` — inspection Store des observations RAW + +## Base requise + +Base directe attendue : `0.3.8-pre.008-fix.001` (`workspace.package.version = 0.3.8-pre.8.fix.1`). + +## Objet + +Ajouter les listings backend-neutres et random-access des observations RAW Transaction et Account, sans modifier le schéma PostgreSQL, les reads par clé existants ni l'UI Store Desk. + +## Store API + +Deux queries sont ajoutées : + +```text +RawTransactionObservationInspectionQuery +RawAccountObservationInspectionQuery +``` + +Elles portent uniquement : + +```text +network obligatoire +référence d'entité exacte optionnelle +direction canonique +RawInspectionPageRequest { offset, limit } +``` + +Une référence fournie doit appartenir au même network que la query. + +Deux summaries sûres sont ajoutées : + +```text +RawTransactionObservationSummary +RawAccountObservationSummary +``` + +La summary Transaction contient observation key, référence transaction et provenance d'acquisition sûre. La summary Account contient observation key, référence account, provenance sûre et metadata Yellowstone optionnelles déjà présentes dans le modèle (`is_startup`, transaction signature, `write_version`). Aucun payload transaction ni byte Account n'est exposé. + +Deux nouvelles capabilities complètent l'API : + +```text +RawTransactionObservationInspectionRead +RawAccountObservationInspectionRead +``` + +L'inventaire RAW passe de 12 à 14 capabilities. Les reads `RawTransactionObservationRead` / `RawAccountObservationRead` par observation key restent inchangés. + +## Store façade + +`ksp-store-lib::Store` implémente et dispatch les deux capabilities après le même contrôle de cohérence network pré-I/O que les autres opérations Store. Aucun type PostgreSQL, SQL ou handle physique ne traverse la façade. + +Les exports crate-root et les canaris d'inventaire sont réconciliés à 105 exports Store façade et 14 capabilities RAW symétriques. + +## PostgreSQL + +`ksp-store-postgres-lib` ajoute quatre statements privés : + +```text +INSPECT_OBSERVATIONS_ASC_SQL +INSPECT_OBSERVATIONS_DESC_SQL +INSPECT_ACCOUNT_OBSERVATIONS_ASC_SQL +INSPECT_ACCOUNT_OBSERVATIONS_DESC_SQL +``` + +Chaque famille utilise une seule instruction par requête avec : + +```text +filtered COUNT exact +total COUNT exact +LEFT JOIN LATERAL +LIMIT/OFFSET checked avant I/O +ordre déterministe received_at_unix_millis + observation_key +``` + +Transaction filtre optionnellement par signature exacte de la référence Transaction. Account filtre optionnellement par l'identité complète `(pubkey, slot, state_hash)`. + +Les projections sont strictement metadata-only : aucun payload/archive Transaction et aucun `account data` ne sont sélectionnés. Les statements keyset/cursor existants restent OFFSET-free et les reads par clé restent inchangés. + +Aucune migration ou index n'est ajouté dans cette tranche. + +## Règles d'Items + +Les nouveaux items `pub`/`pub(crate)` sont réexportés jusqu'aux `lib.rs` concernés. Les accès crate-wide utilisent les chemins `crate::Item`, y compris pour les helpers PostgreSQL visibles entre modules. + +## Hors scope + +Aucune UI Observation n'est ajoutée à Store Desk dans `pre.009`; l'intégration visuelle et le hardening UX restent réservés à `pre.010`. Aucun Config, capability Tauri, migration, index, action de rétention ou changement des DataTables Transactions/Accounts n'est introduit. + +## Version + +```text +delivery = 0.3.8-pre.009 +workspace.package.version = 0.3.8-pre.9 +commit = v0.3.8-pre.009 +tag = aucun +``` + +## Gate requis + +```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-api +cargo test -p ksp-store-postgres-lib +cargo test -p ksp-store-lib +cargo check -p ksp-store-lib --no-default-features +``` 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 d51abf5..99b28e0 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 @@ -608,7 +608,9 @@ Le détail reçoit uniquement l'identité de row `(pubkey, slot, state_hash)`, d ### pre.009 — extension Store observations -Ajouter les listings/summaries d'observations transaction/account backend-neutral. S'ils sont tabulaires, employer la même primitive d'inspection offset/count ; read-by-key existant reste intact. +Étendre le contrat Store backend-neutral avec `RawTransactionObservationInspectionRead` et `RawAccountObservationInspectionRead`, sans modifier les reads par clé existants. Les deux familles utilisent `RawInspectionPageRequest` / `RawInspectionPage` avec counts exacts et random access offset/limit. La requête porte le network obligatoire, une référence d'entité exacte optionnelle et la direction canonique ; la cohérence network/référence est validée avant dispatch. + +Les summaries d'observation restent strictement sans bytes RAW : clé d'observation, référence durable et `RawAcquisitionProvenance` sûre ; Account conserve seulement les metadata Yellowstone optionnelles déjà présentes dans le modèle (`is_startup`, transaction signature associée, `write_version`). PostgreSQL implémente quatre statements privés, un ASC et un DESC par famille, ordonnés déterministiquement par `(received_at_unix_millis, observation_key)` et utilisant counts + `LEFT JOIN LATERAL` + `LIMIT/OFFSET`. Les queries keyset/read-by-key historiques restent inchangées. L'inventaire symétrique Store/PostgreSQL passe de 12 à 14 capabilities RAW. Aucune migration, index, UI ou DTO DataTables n'est ajouté dans cette tranche ; l'intégration visuelle reste réservée à `pre.010`. ### pre.010 — observations UI + hardening UX diff --git a/docs/validation/025-V0_3_8_STORE_DESK.md b/docs/validation/025-V0_3_8_STORE_DESK.md index 7b06610..5312080 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 @@ -686,3 +686,54 @@ Le gate opérateur de `0.3.8-pre.8` confirme que les audits KSP, `cargo check -- - [X] la version workspace devient `0.3.8-pre.8.fix.1` parce qu'un fichier de test Rust est modifié. Le gate Cargo du fix reste à rejouer après application du delta. +## 34. Clôture `pre.008-fix.001` et assemblage `pre.009` — inspection Store observations + +Le gate opérateur de `0.3.8-pre.8.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), 709 file(s)) +cargo check --workspace: PASS +cargo clippy --workspace --all-targets: PASS +cargo test -p ksp-app-store-desk: PASS (28 unitaires + 5 dependency + 13 desktop contract + 12 security + 1 public API) +cargo test -p ksp-store-lib: PASS +cargo check -p ksp-store-lib --no-default-features: PASS +``` + +`pre.009` étend uniquement les contrats et implémentations Store d'inspection des observations RAW : + +- [X] `workspace.package.version = 0.3.8-pre.9` ; +- [X] `RawTransactionObservationInspectionQuery` et `RawAccountObservationInspectionQuery` réutilisent `RawInspectionPageRequest`, direction canonique et network obligatoire ; +- [X] le filtre entité est optionnel mais, lorsqu'il est présent, doit porter exactement le même network que la query ; +- [X] `RawTransactionObservationSummary` expose uniquement observation key, référence transaction et provenance sûre ; +- [X] `RawAccountObservationSummary` expose uniquement observation key, référence account, provenance sûre et metadata Yellowstone optionnelles déjà contractuelles ; +- [X] aucun payload transaction ni byte Account n'appartient aux summaries ; +- [X] `RawTransactionObservationInspectionRead` et `RawAccountObservationInspectionRead` portent l'inventaire API à 14 capabilities RAW ; +- [X] `ksp-store-lib::Store` dispatch les deux nouvelles capabilities après validation network pré-I/O ; +- [X] `ksp-store-postgres-lib` implémente les deux capabilities sans exposer de SQL ou type physique ; +- [X] PostgreSQL utilise deux statements privés par famille, ASC/DESC, avec count exact, `LEFT JOIN LATERAL`, `LIMIT/OFFSET` et ordre déterministe `(received_at_unix_millis, observation_key)` ; +- [X] les statements Observation Transaction ne sélectionnent jamais les payload/archive bytes ; +- [X] les statements Observation Account ne sélectionnent jamais les bytes `data` ; +- [X] les queries keyset/cursor historiques et les reads par observation key restent inchangés ; +- [X] les anciens canaris OFFSET restent scopés par statements nommés afin de ne pas confondre inspection random-access et navigation keyset ; +- [X] les inventaires Store/PostgreSQL sont symétriques à 14 capabilities ; +- [X] les exports Store API et Store façade sont exacts, respectivement 74 et 105 entrées dans leurs canaris d'inventaire ; +- [X] les items visibles restent réexportés jusqu'aux `lib.rs` et consommés via `crate::Item` ; +- [X] aucune migration, index, Config, Store Desk UI ou capability Tauri n'est modifié dans cette tranche. + +Contrôles statiques disponibles dans l'environnement d'assemblage : + +```text +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +API crate-root export inventory: 74 / 74 exact +Store crate-root re-export inventory: 105 / 105 exact +RAW capability inventory: API 14 / Store 14 / PostgreSQL 14 +Store operation-network validation inventory: 18 +``` + +Le gate Cargo de `pre.009` reste à rejouer après application du delta. +