From b32de48c0bb2284cf40a59273fa3ae4d2d185ab3 Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Mon, 24 Aug 2026 13:17:24 +0200 Subject: [PATCH] v0.2.9-pre.006 --- Cargo.toml | 4 +- .../src/{client.rs => http_client.rs} | 6 +- .../src/{executor.rs => http_executor.rs} | 6 +- .../src/{pool.rs => http_pool.rs} | 6 +- .../src/{resilience.rs => http_resilience.rs} | 6 +- .../src/{settings.rs => http_settings.rs} | 6 +- crates/ksp-onchain-transport-lib/src/lib.rs | 95 ++--- .../tests/release_completeness.rs | 41 ++- .../unit_tests/http_client.rs | 62 ++++ .../unit_tests/http_executor.rs | 141 ++++++++ .../unit_tests/http_pool.rs | 336 ++++++++++++++++++ .../unit_tests/http_resilience.rs | 186 ++++++++++ .../unit_tests/http_settings.rs | 207 +++++++++++ deltas/0.2.9/pre.006.md | 167 +++++++++ .../plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md | 74 ++-- .../validation/012-V0_2_9_YELLOWSTONE_GRPC.md | 71 +++- 16 files changed, 1299 insertions(+), 115 deletions(-) rename crates/ksp-onchain-transport-lib/src/{client.rs => http_client.rs} (99%) rename crates/ksp-onchain-transport-lib/src/{executor.rs => http_executor.rs} (98%) rename crates/ksp-onchain-transport-lib/src/{pool.rs => http_pool.rs} (99%) rename crates/ksp-onchain-transport-lib/src/{resilience.rs => http_resilience.rs} (99%) rename crates/ksp-onchain-transport-lib/src/{settings.rs => http_settings.rs} (99%) create mode 100644 crates/ksp-onchain-transport-lib/unit_tests/http_client.rs create mode 100644 crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs create mode 100644 crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs create mode 100644 crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs create mode 100644 crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs create mode 100644 deltas/0.2.9/pre.006.md diff --git a/Cargo.toml b/Cargo.toml index 557c3b0..7544bb0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 245 +# version: 246 [workspace] resolver = "3" members = ["crates/ksp-app-config-desk", "crates/ksp-app-wallet-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-logging-lib", "crates/ksp-onchain-transport-lib", "crates/ksp-wallet-lib"] [workspace.package] -version = "0.2.9-pre.5.fix.1" +version = "0.2.9-pre.6" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/src/client.rs b/crates/ksp-onchain-transport-lib/src/http_client.rs similarity index 99% rename from crates/ksp-onchain-transport-lib/src/client.rs rename to crates/ksp-onchain-transport-lib/src/http_client.rs index ab6bc8a..7f90fe4 100644 --- a/crates/ksp-onchain-transport-lib/src/client.rs +++ b/crates/ksp-onchain-transport-lib/src/http_client.rs @@ -1,5 +1,5 @@ -// file: crates/ksp-onchain-transport-lib/src/client.rs -// version: 7 +// file: crates/ksp-onchain-transport-lib/src/http_client.rs +// version: 8 /// Passive runtime availability reported for one logical HTTP endpoint or role. #[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)] @@ -510,5 +510,5 @@ fn build_reqwest_client(settings: &crate::HttpEndpointSettings) -> std::result:: } #[cfg(test)] -#[path = "../unit_tests/client.rs"] +#[path = "../unit_tests/http_client.rs"] mod tests; diff --git a/crates/ksp-onchain-transport-lib/src/executor.rs b/crates/ksp-onchain-transport-lib/src/http_executor.rs similarity index 98% rename from crates/ksp-onchain-transport-lib/src/executor.rs rename to crates/ksp-onchain-transport-lib/src/http_executor.rs index dda00fb..f4cd338 100644 --- a/crates/ksp-onchain-transport-lib/src/executor.rs +++ b/crates/ksp-onchain-transport-lib/src/http_executor.rs @@ -1,5 +1,5 @@ -// file: crates/ksp-onchain-transport-lib/src/executor.rs -// version: 3 +// file: crates/ksp-onchain-transport-lib/src/http_executor.rs +// version: 4 const HTTP_BAD_GATEWAY: u16 = 502; const HTTP_GATEWAY_TIMEOUT: u16 = 504; @@ -228,5 +228,5 @@ fn http_status_error(method: &crate::HttpRpcMethodDescriptor, status: u16) -> ks } #[cfg(test)] -#[path = "../unit_tests/executor.rs"] +#[path = "../unit_tests/http_executor.rs"] mod tests; diff --git a/crates/ksp-onchain-transport-lib/src/pool.rs b/crates/ksp-onchain-transport-lib/src/http_pool.rs similarity index 99% rename from crates/ksp-onchain-transport-lib/src/pool.rs rename to crates/ksp-onchain-transport-lib/src/http_pool.rs index 1e15447..343ecb7 100644 --- a/crates/ksp-onchain-transport-lib/src/pool.rs +++ b/crates/ksp-onchain-transport-lib/src/http_pool.rs @@ -1,5 +1,5 @@ -// file: crates/ksp-onchain-transport-lib/src/pool.rs -// version: 7 +// file: crates/ksp-onchain-transport-lib/src/http_pool.rs +// version: 8 /// Safe snapshot of the logical HTTP endpoint pool. #[derive(Clone, Debug, Eq, PartialEq)] @@ -616,5 +616,5 @@ fn request_timeout(role: &crate::HttpRoleName, request_kind: &crate::HttpRequest } #[cfg(test)] -#[path = "../unit_tests/pool.rs"] +#[path = "../unit_tests/http_pool.rs"] mod tests; diff --git a/crates/ksp-onchain-transport-lib/src/resilience.rs b/crates/ksp-onchain-transport-lib/src/http_resilience.rs similarity index 99% rename from crates/ksp-onchain-transport-lib/src/resilience.rs rename to crates/ksp-onchain-transport-lib/src/http_resilience.rs index 175fb90..b361f7e 100644 --- a/crates/ksp-onchain-transport-lib/src/resilience.rs +++ b/crates/ksp-onchain-transport-lib/src/http_resilience.rs @@ -1,5 +1,5 @@ -// file: crates/ksp-onchain-transport-lib/src/resilience.rs -// version: 3 +// file: crates/ksp-onchain-transport-lib/src/http_resilience.rs +// version: 4 const DEFAULT_RATE_LIMIT_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(1); const MAX_PROVIDER_RETRY_AFTER: std::time::Duration = std::time::Duration::from_secs(60); @@ -401,5 +401,5 @@ fn retry_backoff(settings: &crate::HttpRetrySettings, retry_number: u32) -> std: } #[cfg(test)] -#[path = "../unit_tests/resilience.rs"] +#[path = "../unit_tests/http_resilience.rs"] mod tests; diff --git a/crates/ksp-onchain-transport-lib/src/settings.rs b/crates/ksp-onchain-transport-lib/src/http_settings.rs similarity index 99% rename from crates/ksp-onchain-transport-lib/src/settings.rs rename to crates/ksp-onchain-transport-lib/src/http_settings.rs index d74343e..cb84a67 100644 --- a/crates/ksp-onchain-transport-lib/src/settings.rs +++ b/crates/ksp-onchain-transport-lib/src/http_settings.rs @@ -1,5 +1,5 @@ -// file: crates/ksp-onchain-transport-lib/src/settings.rs -// version: 6 +// file: crates/ksp-onchain-transport-lib/src/http_settings.rs +// version: 7 /// Runtime HTTP endpoint URL owned by Transport. /// @@ -582,5 +582,5 @@ fn invalid_settings(message: &str, field: &str) -> ksp_core_lib::Result<()> { } #[cfg(test)] -#[path = "../unit_tests/settings.rs"] +#[path = "../unit_tests/http_settings.rs"] mod tests; diff --git a/crates/ksp-onchain-transport-lib/src/lib.rs b/crates/ksp-onchain-transport-lib/src/lib.rs index 639baf9..77f17d4 100644 --- a/crates/ksp-onchain-transport-lib/src/lib.rs +++ b/crates/ksp-onchain-transport-lib/src/lib.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/lib.rs -// version: 38 +// version: 39 #![warn(missing_docs)] #![deny(unreachable_pub)] @@ -36,19 +36,21 @@ //! `0.2.9-pre.003` adds bounded TLS/WebPKI connection establishment, generic redacted ASCII request metadata and the seven standard Yellowstone unary RPCs //! through KSP-owned DTOs. //! `0.2.9-pre.004` materializes the provider-neutral standard `SubscribeRequest` foundation: all seven named filter maps, global filter-name bounds/uniqueness, -//! commitment, ordered account-data slices, ping and `from_slot`. Family-specific account/slot/transaction/block filter fields remain staged for `pre.005–007`. +//! commitment, ordered account-data slices, ping and `from_slot`. Family-specific account/slot filters land in `pre.005`; transaction/block filters remain staged for `pre.007–008`. +//! `0.2.9-pre.006` normalizes the five unambiguously HTTP-owned private implementation modules with an `http_` prefix while preserving shared `rpc_*`, JSON-RPC, error and constants modules. -mod client; mod constants; mod error; -mod executor; mod grpc_channel; mod grpc_settings; mod grpc_subscribe; mod grpc_unary; +mod http_client; +mod http_executor; +mod http_pool; +mod http_resilience; +mod http_settings; mod json_rpc; -mod pool; -mod resilience; mod rpc_accounts; mod rpc_blocks; mod rpc_canary; @@ -58,7 +60,6 @@ mod rpc_economics; mod rpc_method; mod rpc_tokens; mod rpc_transactions; -mod settings; mod ws_accounts; mod ws_blocks; mod ws_cluster; @@ -71,13 +72,13 @@ mod ws_subscription; mod ws_transactions; /// Passive runtime availability reported for one logical HTTP endpoint. -pub use self::client::HttpEndpointAvailability; +pub use self::http_client::HttpEndpointAvailability; /// Shareable logical HTTP endpoint client owned by KSP Transport. -pub use self::client::HttpEndpointClient; +pub use self::http_client::HttpEndpointClient; /// Safe routing snapshot for one configured endpoint role. -pub use self::client::HttpEndpointRoleSnapshot; +pub use self::http_client::HttpEndpointRoleSnapshot; /// Safe metadata snapshot for one logical HTTP endpoint. -pub use self::client::HttpEndpointSnapshot; +pub use self::http_client::HttpEndpointSnapshot; /// Error code used when no logical endpoint can satisfy a request. pub use self::error::ERROR_CODE_ENDPOINT_SELECTION_FAILED; /// Error code used when a Yellowstone gRPC channel cannot be prepared safely. @@ -134,28 +135,24 @@ pub use self::grpc_settings::YellowstoneGrpcReconnectSettings; pub use self::grpc_settings::YellowstoneGrpcSessionSettings; /// Complete runtime settings consumed by the KSP Yellowstone gRPC transport engine. pub use self::grpc_settings::YellowstoneGrpcTransportSettings; -/// One validated standard Yellowstone account predicate. -pub use self::grpc_subscribe::YellowstoneAccountFilterPredicate; /// Typed account payload carried by one standard Yellowstone account update. pub use self::grpc_subscribe::YellowstoneAccountInfo; -/// Lamport comparison used by standard Yellowstone account filters. -pub use self::grpc_subscribe::YellowstoneAccountLamportsFilter; +/// One validated standard Yellowstone account predicate. +pub use self::grpc_subscribe::YellowstoneAccountFilterPredicate; +/// Standard Yellowstone account-update projection owned by KSP. +pub use self::grpc_subscribe::YellowstoneAccountUpdate; /// Validated standard Yellowstone account memcmp predicate. pub use self::grpc_subscribe::YellowstoneAccountMemcmp; /// Encoding selected by one Yellowstone account memcmp predicate. pub use self::grpc_subscribe::YellowstoneAccountMemcmpEncoding; -/// Standard Yellowstone account-update projection owned by KSP. -pub use self::grpc_subscribe::YellowstoneAccountUpdate; +/// Lamport comparison used by standard Yellowstone account filters. +pub use self::grpc_subscribe::YellowstoneAccountLamportsFilter; /// One standard Yellowstone account-data slice. pub use self::grpc_subscribe::YellowstoneAccountsDataSlice; -/// Wire-preserving KSP representation of a standard Yellowstone Cuckoo filter. -pub use self::grpc_subscribe::YellowstoneCuckooFilter; /// Hash algorithm carried by a standard Yellowstone Cuckoo filter. pub use self::grpc_subscribe::YellowstoneCuckooHashAlgorithm; -/// Current standard Yellowstone slot status. -pub use self::grpc_subscribe::YellowstoneSlotStatus; -/// Standard Yellowstone slot-update projection owned by KSP. -pub use self::grpc_subscribe::YellowstoneSlotUpdate; +/// Wire-preserving KSP representation of a standard Yellowstone Cuckoo filter. +pub use self::grpc_subscribe::YellowstoneCuckooFilter; /// Complete account-family filter group for standard Yellowstone Subscribe. pub use self::grpc_subscribe::YellowstoneSubscribeAccountFilter; /// Block-family filter-group shell for standard Yellowstone Subscribe. @@ -172,12 +169,16 @@ pub use self::grpc_subscribe::YellowstoneSubscribePing; pub use self::grpc_subscribe::YellowstoneSubscribeRequest; /// Complete slot-family filter group for standard Yellowstone Subscribe. pub use self::grpc_subscribe::YellowstoneSubscribeSlotFilter; -/// Transaction-family filter-group shell shared by transactions and transaction-status maps. -pub use self::grpc_subscribe::YellowstoneSubscribeTransactionFilter; +/// Current standard Yellowstone slot status. +pub use self::grpc_subscribe::YellowstoneSlotStatus; +/// Standard Yellowstone slot-update projection owned by KSP. +pub use self::grpc_subscribe::YellowstoneSlotUpdate; /// Fixed-width transaction signature attached to Yellowstone account updates when available. pub use self::grpc_subscribe::YellowstoneTransactionSignature; /// Timestamp attached to standard Yellowstone update envelopes. pub use self::grpc_subscribe::YellowstoneUpdateTimestamp; +/// Transaction-family filter-group shell shared by transactions and transaction-status maps. +pub use self::grpc_subscribe::YellowstoneSubscribeTransactionFilter; /// Standard Solana Yellowstone unary facade over one KSP-owned physical gRPC channel. pub use self::grpc_unary::SolanaYellowstoneGrpcUnaryClient; /// Block height returned by the standard Yellowstone unary surface. @@ -209,21 +210,21 @@ pub use self::json_rpc::parse_json_rpc_response_text; /// Validates a decoded JSON value as one JSON-RPC HTTP response. pub use self::json_rpc::parse_json_rpc_response_value; /// Result of one logical endpoint selection. -pub use self::pool::HttpEndpointSelection; +pub use self::http_pool::HttpEndpointSelection; /// Runtime admission permit for one HTTP request. -pub use self::pool::HttpRequestPermit; +pub use self::http_pool::HttpRequestPermit; /// Shareable logical HTTP endpoint pool with priority routing, admission limits and bounded deadlines. -pub use self::pool::HttpTransportPool; +pub use self::http_pool::HttpTransportPool; /// Safe snapshot of the logical HTTP endpoint pool. -pub use self::pool::HttpTransportPoolSnapshot; +pub use self::http_pool::HttpTransportPoolSnapshot; /// Dispatch knowledge used to prevent ambiguous automatic resubmission. -pub use self::resilience::HttpDispatchState; +pub use self::http_resilience::HttpDispatchState; /// Transport-level cause considered by the bounded retry policy. -pub use self::resilience::HttpRetryCause; +pub use self::http_resilience::HttpRetryCause; /// Result of evaluating one bounded transport retry opportunity. -pub use self::resilience::HttpRetryDecision; +pub use self::http_resilience::HttpRetryDecision; /// Evaluates the centralized bounded HTTP retry policy for one audited RPC method. -pub use self::resilience::evaluate_transport_retry; +pub use self::http_resilience::evaluate_transport_retry; /// Typed transport-level Solana account without Program/SPL decoding. pub use self::rpc_accounts::SolanaAccount; /// Address and lamport balance returned by `getLargestAccounts`. @@ -397,25 +398,25 @@ pub use self::rpc_transactions::SolanaTransactionVersion; /// Three-state wire field used when Solana distinguishes omission from an explicit JSON `null`. pub use self::rpc_transactions::SolanaWireField; /// Open cluster or network descriptor used by HTTP endpoint settings. -pub use self::settings::HttpClusterName; +pub use self::http_settings::HttpClusterName; /// Runtime settings for one role declared by an HTTP endpoint. -pub use self::settings::HttpEndpointRoleSettings; +pub use self::http_settings::HttpEndpointRoleSettings; /// Runtime settings for one named Solana HTTP endpoint. -pub use self::settings::HttpEndpointSettings; +pub use self::http_settings::HttpEndpointSettings; /// Runtime HTTP endpoint URL with redacted diagnostics. -pub use self::settings::HttpEndpointUrl; +pub use self::http_settings::HttpEndpointUrl; /// Open provider descriptor used by HTTP endpoint settings. -pub use self::settings::HttpProviderName; +pub use self::http_settings::HttpProviderName; /// Open request-kind descriptor used by logical endpoint capabilities. -pub use self::settings::HttpRequestKind; +pub use self::http_settings::HttpRequestKind; /// Bounded retry settings owned by the HTTP transport runtime. -pub use self::settings::HttpRetrySettings; +pub use self::http_settings::HttpRetrySettings; /// Local limits attached to one logical HTTP endpoint role. -pub use self::settings::HttpRoleLimits; +pub use self::http_settings::HttpRoleLimits; /// Open logical endpoint role descriptor. -pub use self::settings::HttpRoleName; +pub use self::http_settings::HttpRoleName; /// Complete runtime settings consumed by the Solana HTTP transport foundation. -pub use self::settings::HttpTransportSettings; +pub use self::http_settings::HttpTransportSettings; /// Configuration accepted by the standard Solana `accountSubscribe` WebSocket method. pub use self::ws_accounts::SolanaAccountSubscribeConfig; /// One `programNotification` payload preserving contextual and non-contextual upstream forms. @@ -504,17 +505,17 @@ pub use self::ws_transactions::SolanaSignatureSubscribeConfig; /// Owning tracing target for events emitted by the on-chain transport crate. pub(crate) use self::constants::TRACING_TARGET; /// Crate-internal `HttpConcurrencyPermit` state shared across the owning crate. -pub(crate) use self::resilience::HttpConcurrencyPermit; +pub(crate) use self::http_resilience::HttpConcurrencyPermit; /// Crate-internal `HttpRoleRuntime` state shared across the owning crate. -pub(crate) use self::resilience::HttpRoleRuntime; +pub(crate) use self::http_resilience::HttpRoleRuntime; /// Crate-internal `RoleAdmissionAttempt` variants used by the owning crate. -pub(crate) use self::resilience::RoleAdmissionAttempt; +pub(crate) use self::http_resilience::RoleAdmissionAttempt; /// Decodes one private serde wire type into the shared Transport error domain for typed RPC adapters. pub(crate) use self::rpc_common::decode_wire_json; /// Parses a base58 public key without echoing its wire value into diagnostics for typed RPC adapters. pub(crate) use self::rpc_common::parse_wire_pubkey; /// Validates endpoint settings. -pub(crate) use self::settings::validate_endpoint_settings; +pub(crate) use self::http_settings::validate_endpoint_settings; /// Crate-internal command surface shared by the physical session and typed subscription handle. pub(crate) use self::ws_session::WsSessionCommand; /// Crate-internal notification dispatch result. diff --git a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs index e1f0cc2..7ed02c0 100644 --- a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs +++ b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/release_completeness.rs -// version: 36 +// version: 37 //! Release-level completeness canaries for staged HTTP and WebSocket Transport coverage. @@ -1071,19 +1071,7 @@ fn release_v0_2_9_pre_003_adds_tls_metadata_and_exactly_seven_standard_unary_met fn release_v0_2_9_pre_004_materializes_only_standard_subscribe_common_contract() { let source = include_str!("../src/grpc_subscribe.rs"); let crate_root = include_str!("../src/lib.rs"); - for wire_field in [ - "accounts:", - "slots:", - "transactions:", - "transactions_status:", - "blocks:", - "blocks_meta:", - "entry:", - "commitment:", - "accounts_data_slice:", - "ping:", - "from_slot:", - ] { + for wire_field in ["accounts:", "slots:", "transactions:", "transactions_status:", "blocks:", "blocks_meta:", "entry:", "commitment:", "accounts_data_slice:", "ping:", "from_slot:"] { assert!(source.contains(wire_field), "missing standard Yellowstone SubscribeRequest field: {wire_field}"); } assert!(source.contains("MAX_GRPC_SUBSCRIBE_FILTER_GROUP_COUNT")); @@ -1155,3 +1143,28 @@ fn release_v0_2_9_pre_005_completes_standard_accounts_and_slots_without_advancin let _account = std::any::type_name::(); let _slot = std::any::type_name::(); } + +#[test] +fn release_v0_2_9_pre_006_namespaces_unambiguously_http_owned_private_modules() { + let crate_root = include_str!("../src/lib.rs"); + let http_client = include_str!("../src/http_client.rs"); + let http_executor = include_str!("../src/http_executor.rs"); + let http_pool = include_str!("../src/http_pool.rs"); + let http_resilience = include_str!("../src/http_resilience.rs"); + let http_settings = include_str!("../src/http_settings.rs"); + for declaration in ["mod http_client;", "mod http_executor;", "mod http_pool;", "mod http_resilience;", "mod http_settings;"] { + assert!(crate_root.contains(declaration), "missing HTTP-owned private module declaration: {declaration}"); + } + for historical in ["mod client;", "mod executor;", "mod pool;", "mod resilience;", "mod settings;"] { + assert!(!crate_root.contains(historical), "historical ambiguous private module declaration remains: {historical}"); + } + assert!(http_client.contains("HttpEndpointClient")); + assert!(http_executor.contains("execute_standard_rpc")); + assert!(http_pool.contains("HttpTransportPool")); + assert!(http_resilience.contains("HttpRetryDecision")); + assert!(http_settings.contains("HttpTransportSettings")); + // Shared/protocol-oriented modules deliberately keep their existing names. + for shared in ["mod constants;", "mod error;", "mod json_rpc;", "mod rpc_common;", "mod rpc_accounts;", "mod rpc_blocks;", "mod rpc_transactions;"] { + assert!(crate_root.contains(shared), "shared/protocol module was incorrectly HTTP-prefixed: {shared}"); + } +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/http_client.rs b/crates/ksp-onchain-transport-lib/unit_tests/http_client.rs new file mode 100644 index 0000000..d164f6a --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/http_client.rs @@ -0,0 +1,62 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/http_client.rs +// version: 4 + +fn endpoint(enabled: bool, url_text: &str) -> crate::HttpEndpointSettings { + let url = crate::HttpEndpointUrl::parse(url_text).expect("test endpoint URL must parse"); + let role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 10, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ); + return crate::HttpEndpointSettings::new( + "endpoint", + enabled, + crate::HttpProviderName::new("provider"), + crate::HttpClusterName::new("devnet"), + url, + std::time::Duration::from_secs(1), + std::time::Duration::from_secs(2), + std::option::Option::Some(4), + std::vec![role], + ); +} + +#[test] +fn endpoint_client_snapshot_never_contains_url_or_secret_material() { + let client = crate::HttpEndpointClient::new(endpoint(true, "https://provider.invalid/rpc?api-key=SECRET-CANARY")).expect("client must build"); + let snapshot = client.snapshot(); + let rendered = format!("{snapshot:?} {client:?}"); + assert_eq!(snapshot.availability(), crate::HttpEndpointAvailability::Available); + assert!(!rendered.contains("SECRET-CANARY")); + assert!(!rendered.contains("provider.invalid")); + assert!(!rendered.contains("https://")); +} + +#[test] +fn disabled_endpoint_client_is_visible_but_not_selectable() { + let client = crate::HttpEndpointClient::new(endpoint(false, "https://api.devnet.solana.com")).expect("disabled client must still build"); + assert_eq!(client.snapshot().availability(), crate::HttpEndpointAvailability::Disabled); + assert!(!client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))); +} + +#[test] +fn endpoint_client_matches_exact_and_wildcard_capabilities() { + let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build"); + assert!(client.supports(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance"))); + assert!(!client.supports(&crate::HttpRoleName::new("write"), &crate::HttpRequestKind::new("get_balance"))); +} + +#[test] +fn endpoint_role_snapshot_exposes_safe_resilience_state() { + let client = crate::HttpEndpointClient::new(endpoint(true, "https://api.devnet.solana.com")).expect("client must build"); + let snapshot = client.snapshot(); + let role = &snapshot.roles()[0]; + assert_eq!(role.availability(), crate::HttpEndpointAvailability::Available); + assert_eq!(role.in_flight_requests(), std::option::Option::None); + assert_eq!(role.cooldown_remaining(), std::option::Option::None); + assert_eq!(role.success_count(), 0); + assert_eq!(role.failure_count(), 0); + assert_eq!(role.rate_limit_count(), 0); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs b/crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs new file mode 100644 index 0000000..4694b55 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs @@ -0,0 +1,141 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs +// version: 3 + +fn pool_for_url(url: &str, request_timeout: std::time::Duration, max_retries: u32) -> crate::HttpTransportPool { + let role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 10, + crate::HttpRoleLimits::new( + std::option::Option::None, + std::option::Option::None, + std::option::Option::None, + std::option::Option::Some(std::time::Duration::from_millis(1)), + ), + ); + let endpoint = crate::HttpEndpointSettings::new( + "fixture", + true, + crate::HttpProviderName::new("fixture"), + crate::HttpClusterName::new("local"), + crate::HttpEndpointUrl::parse(url).expect("fixture URL must parse"), + std::time::Duration::from_millis(100), + request_timeout, + std::option::Option::Some(1), + std::vec![role], + ); + let settings = crate::HttpTransportSettings::new( + std::vec![endpoint], + crate::HttpRetrySettings::new(max_retries, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)), + ); + return crate::HttpTransportPool::new(settings).expect("fixture pool must build"); +} + +fn health_method() -> &'static crate::HttpRpcMethodDescriptor { + return crate::find_http_rpc_method("getHealth").expect("getHealth descriptor must exist"); +} + +fn serve_rate_limit_then_success() -> (std::string::String, std::thread::JoinHandle) { + let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind"); + let address = listener.local_addr().expect("fixture listener address must resolve"); + let handle = std::thread::spawn(move || { + let mut count = 0_usize; + while count < 2 { + let (mut stream, _) = listener.accept().expect("fixture server must accept request"); + let _ = read_request(&mut stream); + let response = if count == 0 { + "HTTP/1.1 429 Too Many Requests\r\nRetry-After: 0\r\nContent-Length: 0\r\nConnection: close\r\n\r\n".to_owned() + } else { + let body = include_str!("../fixtures/http/get_health.success.json"); + format!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", body.len(), body) + }; + std::io::Write::write_all(&mut stream, response.as_bytes()).expect("fixture response must write"); + count = count.saturating_add(1); + } + return count; + }); + return (format!("http://{address}"), handle); +} + +fn serve_timeout() -> (std::string::String, std::thread::JoinHandle<()>) { + let listener = std::net::TcpListener::bind("127.0.0.1:0").expect("fixture listener must bind"); + let address = listener.local_addr().expect("fixture listener address must resolve"); + let handle = std::thread::spawn(move || { + let (mut stream, _) = listener.accept().expect("fixture server must accept request"); + let _ = read_request(&mut stream); + std::thread::sleep(std::time::Duration::from_millis(100)); + return; + }); + return (format!("http://{address}"), handle); +} + +fn read_request(stream: &mut std::net::TcpStream) -> std::string::String { + let mut bytes = std::vec::Vec::new(); + let mut buffer = [0_u8; 1024]; + loop { + let count = std::io::Read::read(stream, &mut buffer).expect("fixture request must read"); + if count == 0 { + break; + } + bytes.extend_from_slice(&buffer[..count]); + if request_complete(bytes.as_slice()) { + break; + } + } + return std::string::String::from_utf8(bytes).expect("fixture request must be UTF-8"); +} + +fn request_complete(bytes: &[u8]) -> bool { + let text = match std::str::from_utf8(bytes) { + std::result::Result::Ok(text) => text, + std::result::Result::Err(_) => return false, + }; + let header_end = match text.find("\r\n\r\n") { + std::option::Option::Some(value) => value, + std::option::Option::None => return false, + }; + let mut content_length = 0_usize; + for line in text[..header_end].lines() { + let (name, value) = match line.split_once(':') { + std::option::Option::Some(parts) => parts, + std::option::Option::None => continue, + }; + if name.eq_ignore_ascii_case("content-length") { + content_length = value.trim().parse::().expect("content length must parse"); + } + } + return bytes.len() >= header_end.saturating_add(4).saturating_add(content_length); +} + +#[tokio::test(flavor = "current_thread")] +async fn executor_applies_retry_after_and_retries_http_429_for_retry_safe_method() { + let (url, handle) = serve_rate_limit_then_success(); + let pool = pool_for_url(url.as_str(), std::time::Duration::from_millis(500), 1); + let result = pool + .execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new()) + .await + .expect("retry-safe request must recover from one 429"); + assert_eq!(result, serde_json::json!("ok")); + assert_eq!(handle.join().expect("fixture server must join"), 2); + let snapshot = pool.snapshot(); + assert_eq!(snapshot.endpoints()[0].roles()[0].rate_limit_count(), 1); + assert_eq!(snapshot.endpoints()[0].roles()[0].success_count(), 1); +} + +#[tokio::test(flavor = "current_thread")] +async fn executor_maps_reqwest_timeout_to_ksp_timeout_error_without_endpoint_secret_leak() { + const SECRET_CANARY: &str = "SECRET-REQWEST-URL-CANARY"; + let (url, handle) = serve_timeout(); + let sensitive_url = format!("{url}/rpc?api-key={SECRET_CANARY}"); + let pool = pool_for_url(sensitive_url.as_str(), std::time::Duration::from_millis(20), 0); + let error = pool + .execute_standard_rpc(&crate::HttpRoleName::new("default"), health_method(), std::vec::Vec::new()) + .await + .expect_err("timed out request must fail"); + assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); + assert!(!format!("{error:?}").contains(SECRET_CANARY)); + let source = std::error::Error::source(&error).expect("transport timeout should preserve a sanitized reqwest source"); + assert!(!format!("{source:?}").contains(SECRET_CANARY)); + handle.join().expect("fixture server must join"); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs b/crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs new file mode 100644 index 0000000..24e87df --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs @@ -0,0 +1,336 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs +// version: 5 + +fn role(name: &str, priority: u32, request_kinds: std::vec::Vec) -> crate::HttpEndpointRoleSettings { + return crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new(name), + true, + request_kinds, + priority, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ); +} + +fn endpoint(name: &str, enabled: bool, priority: u32, request_kinds: std::vec::Vec) -> crate::HttpEndpointSettings { + return crate::HttpEndpointSettings::new( + name, + enabled, + crate::HttpProviderName::new("provider"), + crate::HttpClusterName::new("devnet"), + crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"), + std::time::Duration::from_secs(1), + std::time::Duration::from_secs(2), + std::option::Option::Some(4), + std::vec![role("default", priority, request_kinds)], + ); +} + +fn non_zero(value: u32) -> std::num::NonZeroU32 { + return std::num::NonZeroU32::new(value).expect("test limit must be non-zero"); +} + +fn limited_endpoint( + name: &str, + priority: u32, + requests_per_second: std::option::Option, + burst_capacity: std::option::Option, + max_concurrent_requests: std::option::Option, + cooldown: std::option::Option, +) -> crate::HttpEndpointSettings { + let limits = crate::HttpRoleLimits::new( + requests_per_second.map(|value| return non_zero(value)), + burst_capacity.map(|value| return non_zero(value)), + max_concurrent_requests.map(|value| return non_zero(value)), + cooldown, + ); + let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], priority, limits); + return crate::HttpEndpointSettings::new( + name, + true, + crate::HttpProviderName::new("provider"), + crate::HttpClusterName::new("devnet"), + crate::HttpEndpointUrl::parse(format!("https://{name}.invalid/rpc?token=SECRET-CANARY")).expect("test URL must parse"), + std::time::Duration::from_secs(1), + std::time::Duration::from_secs(2), + std::option::Option::Some(4), + std::vec![role], + ); +} + +fn settings(endpoints: std::vec::Vec) -> crate::HttpTransportSettings { + return crate::HttpTransportSettings::new( + endpoints, + crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(10), std::time::Duration::from_millis(50)), + ); +} + +#[test] +fn pool_prefers_lowest_priority_tier() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + endpoint("secondary", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), + endpoint("primary", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), + ])) + .expect("pool must build"); + let selection = pool + .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) + .expect("selection must succeed"); + assert_eq!(selection.endpoint_name(), "primary"); + assert_eq!(selection.priority(), 10); +} + +#[test] +fn pool_round_robins_fairly_inside_best_priority_tier() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + endpoint("one", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), + endpoint("two", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), + endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), + ])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let first = pool.select_for_request_kind(&role, &kind).expect("first selection must succeed"); + let second = pool.select_for_request_kind(&role, &kind).expect("second selection must succeed"); + let third = pool.select_for_request_kind(&role, &kind).expect("third selection must succeed"); + assert_eq!(first.endpoint_name(), "one"); + assert_eq!(second.endpoint_name(), "two"); + assert_eq!(third.endpoint_name(), "one"); +} + +#[test] +fn disabled_best_priority_endpoint_falls_back_to_next_tier() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + endpoint("disabled-primary", false, 1, std::vec![crate::HttpRequestKind::wildcard()]), + endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), + ])) + .expect("pool must build"); + let selection = pool + .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) + .expect("fallback must be selected"); + assert_eq!(selection.endpoint_name(), "fallback"); +} + +#[test] +fn pool_filters_role_and_capability_before_priority() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + endpoint("wrong-capability", true, 1, std::vec![crate::HttpRequestKind::new("send_transaction")]), + endpoint("matching", true, 50, std::vec![crate::HttpRequestKind::new("get_balance")]), + ])) + .expect("pool must build"); + let selection = pool + .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) + .expect("matching capability must be selected"); + assert_eq!(selection.endpoint_name(), "matching"); +} + +#[test] +fn pool_returns_structured_error_when_no_endpoint_matches() { + let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("read-only", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)])) + .expect("pool must build"); + let error = pool + .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("send_transaction")) + .expect_err("unsupported request kind must fail selection"); + assert_eq!(error.code(), crate::ERROR_CODE_ENDPOINT_SELECTION_FAILED); + assert!(!format!("{error:?}").contains("SECRET-CANARY")); +} + +#[test] +fn standard_method_selection_uses_registry_request_kind() { + let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("balance", true, 10, std::vec![crate::HttpRequestKind::new("get_balance")],)])) + .expect("pool must build"); + let method = crate::find_http_rpc_method("getBalance").expect("audited method must exist"); + let selection = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect("standard method must route"); + assert_eq!(selection.request_kind().as_str(), "get_balance"); +} + +#[test] +fn pool_snapshot_is_safe_and_preserves_disabled_endpoints() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + endpoint("enabled", true, 10, std::vec![crate::HttpRequestKind::wildcard()]), + endpoint("disabled", false, 10, std::vec![crate::HttpRequestKind::wildcard()]), + ])) + .expect("pool must build"); + let snapshot = pool.snapshot(); + let rendered = format!("{snapshot:?} {pool:?}"); + assert_eq!(snapshot.endpoint_count(), 2); + assert_eq!(snapshot.available_endpoint_count(), 1); + assert!(!rendered.contains("SECRET-CANARY")); + assert!(!rendered.contains(".invalid/rpc")); +} + +#[test] +fn disabled_role_is_excluded_before_priority_selection() { + let base = endpoint("disabled-role", true, 1, std::vec![crate::HttpRequestKind::wildcard()]); + let disabled_role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + false, + std::vec![crate::HttpRequestKind::wildcard()], + 1, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ); + let enabled_non_matching_role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("maintenance"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 1, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ); + let disabled_role_endpoint = crate::HttpEndpointSettings::new( + base.name(), + true, + base.provider().clone(), + base.cluster().clone(), + base.url().clone(), + base.connect_timeout(), + base.request_timeout(), + base.max_idle_connections_per_host(), + std::vec![disabled_role, enabled_non_matching_role], + ); + let pool = crate::HttpTransportPool::new(settings(std::vec![ + disabled_role_endpoint, + endpoint("fallback", true, 20, std::vec![crate::HttpRequestKind::wildcard()]), + ])) + .expect("pool must build"); + let selection = pool + .select_for_request_kind(&crate::HttpRoleName::new("default"), &crate::HttpRequestKind::new("get_balance")) + .expect("enabled fallback role must be selected"); + assert_eq!(selection.endpoint_name(), "fallback"); +} + +#[test] +fn removed_standard_method_is_rejected_before_endpoint_routing() { + let pool = crate::HttpTransportPool::new(settings(std::vec![endpoint("wildcard", true, 10, std::vec![crate::HttpRequestKind::wildcard()],)])) + .expect("pool must build"); + let method = crate::find_http_rpc_method("confirmTransaction").expect("historical method must exist"); + let error = pool.select_for_method(&crate::HttpRoleName::new("default"), method).expect_err("removed standard method must be rejected before routing"); + assert_eq!(error.code(), crate::ERROR_CODE_METHOD_REMOVED); +} + +#[tokio::test] +async fn runtime_concurrency_saturation_falls_back_to_lower_priority_tier() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + limited_endpoint("primary", 1, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None), + limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None), + ])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must acquire primary"); + assert_eq!(first.selection().endpoint_name(), "primary"); + let second = pool.acquire_for_request_kind(&role, &kind).await.expect("second request must fall back while primary is saturated"); + assert_eq!(second.selection().endpoint_name(), "fallback"); +} + +#[tokio::test] +async fn runtime_token_bucket_exhaustion_falls_back_without_busy_waiting() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + limited_endpoint("primary", 1, std::option::Option::Some(1), std::option::Option::Some(1), std::option::Option::None, std::option::Option::None), + limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first request must consume primary token"); + assert_eq!(first.selection().endpoint_name(), "primary"); + drop(first); + let second = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be used while primary token bucket refills"); + assert_eq!(second.selection().endpoint_name(), "fallback"); +} + +#[tokio::test] +async fn provider_cooldown_excludes_rate_limited_role_and_uses_fallback() { + let pool = crate::HttpTransportPool::new(settings(std::vec![ + limited_endpoint( + "primary", + 1, + std::option::Option::None, + std::option::Option::None, + std::option::Option::None, + std::option::Option::Some(std::time::Duration::from_millis(50)), + ), + limited_endpoint("fallback", 20, std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let primary = pool.acquire_for_request_kind(&role, &kind).await.expect("primary must be acquired"); + assert_eq!(primary.selection().endpoint_name(), "primary"); + let pause = primary.record_rate_limited(std::option::Option::None); + assert_eq!(pause, std::time::Duration::from_millis(50)); + drop(primary); + let fallback = pool.acquire_for_request_kind(&role, &kind).await.expect("fallback must be selected during primary cooldown"); + assert_eq!(fallback.selection().endpoint_name(), "fallback"); +} + +#[tokio::test] +async fn admission_waits_for_released_concurrency_without_holding_a_sync_mutex_across_await() { + let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( + "primary", + 1, + std::option::Option::None, + std::option::Option::None, + std::option::Option::Some(1), + std::option::Option::None, + )])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let first = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired"); + let release_task = tokio::spawn(async move { + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + drop(first); + }); + let second = pool + .acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(100)) + .await + .expect("second permit must wake after concurrency release"); + assert_eq!(second.selection().endpoint_name(), "primary"); + release_task.await.expect("release task must complete"); +} + +#[tokio::test] +async fn admission_timeout_is_bounded_when_concurrency_never_becomes_available() { + let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( + "primary", + 1, + std::option::Option::None, + std::option::Option::None, + std::option::Option::Some(1), + std::option::Option::None, + )])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let _held = pool.acquire_for_request_kind(&role, &kind).await.expect("first permit must be acquired"); + let error = pool + .acquire_for_request_kind_with_timeout(&role, &kind, std::time::Duration::from_millis(20)) + .await + .expect_err("second permit must time out while concurrency remains saturated"); + assert_eq!(error.code(), crate::ERROR_CODE_TIMEOUT); +} + +#[tokio::test] +async fn passive_health_snapshot_moves_from_degraded_back_to_available_after_success() { + let pool = crate::HttpTransportPool::new(settings(std::vec![limited_endpoint( + "primary", + 1, + std::option::Option::None, + std::option::Option::None, + std::option::Option::None, + std::option::Option::None, + )])) + .expect("pool must build"); + let role = crate::HttpRoleName::new("default"); + let kind = crate::HttpRequestKind::new("get_balance"); + let first = pool.acquire_for_request_kind(&role, &kind).await.expect("request permit must be acquired"); + first.record_failure(); + drop(first); + let degraded = pool.snapshot(); + assert_eq!(degraded.endpoints()[0].availability(), crate::HttpEndpointAvailability::Degraded); + assert_eq!(degraded.endpoints()[0].roles()[0].failure_count(), 1); + let second = pool.acquire_for_request_kind(&role, &kind).await.expect("degraded endpoint remains eligible for passive recovery"); + second.record_success(); + drop(second); + let recovered = pool.snapshot(); + assert_eq!(recovered.endpoints()[0].availability(), crate::HttpEndpointAvailability::Available); + assert_eq!(recovered.endpoints()[0].roles()[0].success_count(), 1); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs b/crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs new file mode 100644 index 0000000..f81a915 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs @@ -0,0 +1,186 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs +// version: 3 + +fn non_zero(value: u32) -> std::num::NonZeroU32 { + return std::num::NonZeroU32::new(value).expect("test limit must be non-zero"); +} + +fn retry_settings() -> crate::HttpRetrySettings { + return crate::HttpRetrySettings::new(4, std::time::Duration::from_millis(100), std::time::Duration::from_millis(500)); +} + +fn method(name: &str) -> &'static crate::HttpRpcMethodDescriptor { + return crate::find_http_rpc_method(name).expect("audited test method must exist"); +} + +fn role_limits( + requests_per_second: std::option::Option, + burst_capacity: std::option::Option, + max_concurrent_requests: std::option::Option, + cooldown: std::option::Option, +) -> crate::HttpEndpointRoleSettings { + return crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 10, + crate::HttpRoleLimits::new( + requests_per_second.map(|value| return non_zero(value)), + burst_capacity.map(|value| return non_zero(value)), + max_concurrent_requests.map(|value| return non_zero(value)), + cooldown, + ), + ); +} + +#[test] +fn retry_backoff_is_exponential_and_bounded() { + let settings = retry_settings(); + assert_eq!(super::retry_backoff(&settings, 1), std::time::Duration::from_millis(100)); + assert_eq!(super::retry_backoff(&settings, 2), std::time::Duration::from_millis(200)); + assert_eq!(super::retry_backoff(&settings, 3), std::time::Duration::from_millis(400)); + assert_eq!(super::retry_backoff(&settings, 4), std::time::Duration::from_millis(500)); + assert_eq!(super::retry_backoff(&settings, 32), std::time::Duration::from_millis(500)); +} + +#[test] +fn retry_safe_timeout_is_retried_until_budget_is_exhausted() { + let settings = retry_settings(); + let first = crate::evaluate_transport_retry( + method("getBalance"), + &settings, + crate::HttpRetryCause::Timeout, + crate::HttpDispatchState::DispatchedAmbiguous, + 0, + std::option::Option::None, + ); + assert_eq!(first, crate::HttpRetryDecision::RetryAfter(std::time::Duration::from_millis(100))); + let exhausted = crate::evaluate_transport_retry( + method("getBalance"), + &settings, + crate::HttpRetryCause::Timeout, + crate::HttpDispatchState::DispatchedAmbiguous, + settings.max_retries(), + std::option::Option::None, + ); + assert_eq!(exhausted, crate::HttpRetryDecision::Stop); +} + +#[test] +fn write_submission_never_retries_after_ambiguous_dispatch() { + let decision = crate::evaluate_transport_retry( + method("sendTransaction"), + &retry_settings(), + crate::HttpRetryCause::Connection, + crate::HttpDispatchState::DispatchedAmbiguous, + 0, + std::option::Option::None, + ); + assert_eq!(decision, crate::HttpRetryDecision::Stop); +} + +#[test] +fn write_submission_can_retry_when_transport_proves_no_dispatch() { + let decision = crate::evaluate_transport_retry( + method("sendTransaction"), + &retry_settings(), + crate::HttpRetryCause::Connection, + crate::HttpDispatchState::NotDispatched, + 0, + std::option::Option::None, + ); + assert!(decision.should_retry()); +} + +#[test] +fn rpc_application_and_invalid_response_are_not_transport_retries() { + for cause in [crate::HttpRetryCause::RpcApplication, crate::HttpRetryCause::InvalidResponse, crate::HttpRetryCause::Request] { + let decision = crate::evaluate_transport_retry( + method("getBalance"), + &retry_settings(), + cause, + crate::HttpDispatchState::NotDispatched, + 0, + std::option::Option::None, + ); + assert_eq!(decision, crate::HttpRetryDecision::Stop); + } +} + +#[test] +fn provider_retry_after_can_extend_backoff_but_is_defensively_bounded() { + let settings = retry_settings(); + let extended = crate::evaluate_transport_retry( + method("getBalance"), + &settings, + crate::HttpRetryCause::RateLimited, + crate::HttpDispatchState::DispatchedAmbiguous, + 0, + std::option::Option::Some(std::time::Duration::from_secs(3)), + ); + assert_eq!(extended.delay(), std::option::Option::Some(std::time::Duration::from_secs(3))); + let bounded = crate::evaluate_transport_retry( + method("getBalance"), + &settings, + crate::HttpRetryCause::RateLimited, + crate::HttpDispatchState::DispatchedAmbiguous, + 0, + std::option::Option::Some(std::time::Duration::from_secs(600)), + ); + assert_eq!(bounded.delay(), std::option::Option::Some(std::time::Duration::from_secs(60))); +} + +#[test] +fn token_bucket_consumes_burst_then_refills_from_elapsed_time() { + let start = std::time::Instant::now(); + let mut bucket = super::HttpTokenBucketState::new(2, 2, start); + assert!(bucket.try_consume_at(start).is_none()); + assert!(bucket.try_consume_at(start).is_none()); + assert!(bucket.try_consume_at(start).is_some()); + let later = start.checked_add(std::time::Duration::from_millis(500)).expect("test instant must advance"); + assert!(bucket.try_consume_at(later).is_none()); +} + +#[test] +fn absent_burst_capacity_defaults_to_one_second_of_rps_capacity() { + let role = role_limits(std::option::Option::Some(2), std::option::Option::None, std::option::Option::None, std::option::Option::None); + let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new()))); + let now = std::time::Instant::now(); + let first = runtime.try_acquire(now); + let second = runtime.try_acquire(now); + let third = runtime.try_acquire(now); + assert!(matches!(first, crate::RoleAdmissionAttempt::Ready(_))); + assert!(matches!(second, crate::RoleAdmissionAttempt::Ready(_))); + assert!(matches!(third, crate::RoleAdmissionAttempt::BlockedUntil(_))); +} + +#[test] +fn concurrency_semaphore_releases_capacity_when_permit_is_dropped() { + let role = role_limits(std::option::Option::None, std::option::Option::None, std::option::Option::Some(1), std::option::Option::None); + let runtime = std::sync::Arc::new(crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new()))); + let now = std::time::Instant::now(); + let first = runtime.try_acquire(now); + let held = match first { + crate::RoleAdmissionAttempt::Ready(permit) => permit, + _ => panic!("first concurrency permit must be available"), + }; + assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::ConcurrencySaturated)); + drop(held); + assert!(matches!(runtime.try_acquire(now), crate::RoleAdmissionAttempt::Ready(_))); +} + +#[test] +fn rate_limit_cooldown_marks_role_and_caps_provider_delay() { + let role = role_limits( + std::option::Option::None, + std::option::Option::None, + std::option::Option::None, + std::option::Option::Some(std::time::Duration::from_millis(10)), + ); + let runtime = crate::HttpRoleRuntime::new(&role, std::sync::Arc::new(tokio::sync::Notify::new())); + let pause = runtime.record_rate_limited(std::option::Option::Some(std::time::Duration::from_secs(600))); + assert_eq!(pause, std::time::Duration::from_secs(60)); + assert_eq!(runtime.rate_limit_count(), 1); + assert_eq!(runtime.failure_count(), 1); + assert_eq!(runtime.availability(std::time::Instant::now()), crate::HttpEndpointAvailability::RateLimited); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs b/crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs new file mode 100644 index 0000000..91e2f29 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs @@ -0,0 +1,207 @@ +// file: crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs +// version: 3 + +fn non_zero(value: u32) -> std::num::NonZeroU32 { + return std::num::NonZeroU32::new(value).expect("test non-zero value must remain non-zero"); +} + +fn valid_settings(url_text: &str) -> crate::HttpTransportSettings { + let url = crate::HttpEndpointUrl::parse(url_text).expect("test URL must be valid"); + let limits = crate::HttpRoleLimits::new( + std::option::Option::Some(non_zero(10)), + std::option::Option::Some(non_zero(20)), + std::option::Option::Some(non_zero(4)), + std::option::Option::Some(std::time::Duration::from_millis(500)), + ); + let role = crate::HttpEndpointRoleSettings::new(crate::HttpRoleName::new("default"), true, std::vec![crate::HttpRequestKind::wildcard()], 100, limits); + let endpoint = crate::HttpEndpointSettings::new( + "devnet_public", + true, + crate::HttpProviderName::new("solana-public"), + crate::HttpClusterName::new("devnet"), + url, + std::time::Duration::from_secs(5), + std::time::Duration::from_secs(15), + std::option::Option::Some(8), + std::vec![role], + ); + return crate::HttpTransportSettings::new( + std::vec![endpoint], + crate::HttpRetrySettings::new(2, std::time::Duration::from_millis(100), std::time::Duration::from_secs(2)), + ); +} + +#[test] +fn endpoint_url_accepts_http_and_https() { + assert!(crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").is_ok()); + assert!(crate::HttpEndpointUrl::parse("http://127.0.0.1:8899").is_ok()); +} + +#[test] +fn endpoint_url_rejects_non_http_schemes() { + let result = crate::HttpEndpointUrl::parse("ws://api.devnet.solana.com"); + let error = result.expect_err("WebSocket URL must not be accepted by HTTP settings"); + assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); +} + +#[test] +fn endpoint_url_debug_redacts_secret_material() { + let url = crate::HttpEndpointUrl::parse("https://provider.invalid/rpc?api-key=SECRET-CANARY").expect("test URL must parse"); + let rendered = format!("{url:?}"); + assert!(rendered.contains("")); + assert!(!rendered.contains("SECRET-CANARY")); + assert!(!rendered.contains("provider.invalid")); +} + +#[test] +fn valid_transport_settings_pass_validation() { + let settings = valid_settings("https://api.devnet.solana.com"); + assert!(settings.validate().is_ok()); +} + +#[test] +fn transport_settings_debug_does_not_leak_endpoint_url() { + let settings = valid_settings("https://provider.invalid/rpc?api-key=SECRET-CANARY"); + let rendered = format!("{settings:?}"); + assert!(!rendered.contains("SECRET-CANARY")); + assert!(!rendered.contains("provider.invalid")); + assert!(rendered.contains("HttpEndpointUrl()")); +} + +#[test] +fn transport_settings_require_one_enabled_endpoint() { + let url = crate::HttpEndpointUrl::parse("https://api.devnet.solana.com").expect("test URL must parse"); + let role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 100, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::None, std::option::Option::None, std::option::Option::None), + ); + let endpoint = crate::HttpEndpointSettings::new( + "disabled", + false, + crate::HttpProviderName::new("provider"), + crate::HttpClusterName::new("devnet"), + url, + std::time::Duration::from_secs(1), + std::time::Duration::from_secs(1), + std::option::Option::None, + std::vec![role], + ); + let settings = crate::HttpTransportSettings::new( + std::vec![endpoint], + crate::HttpRetrySettings::new(1, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)), + ); + let error = settings.validate().expect_err("all-disabled settings must fail"); + assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); +} + +#[test] +fn transport_settings_reject_duplicate_endpoint_names() { + let first = valid_settings("https://one.invalid"); + let second = valid_settings("https://two.invalid"); + let settings = crate::HttpTransportSettings::new(std::vec![first.endpoints()[0].clone(), second.endpoints()[0].clone()], first.retry().clone()); + let error = settings.validate().expect_err("duplicate endpoint names must fail"); + assert_eq!(error.code(), crate::ERROR_CODE_INVALID_SETTINGS); +} + +#[test] +fn transport_settings_reject_duplicate_roles() { + let base = valid_settings("https://api.devnet.solana.com"); + let endpoint = &base.endpoints()[0]; + let duplicated_endpoint = crate::HttpEndpointSettings::new( + endpoint.name(), + true, + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.url().clone(), + endpoint.connect_timeout(), + endpoint.request_timeout(), + endpoint.max_idle_connections_per_host(), + std::vec![endpoint.roles()[0].clone(), endpoint.roles()[0].clone()], + ); + let settings = crate::HttpTransportSettings::new(std::vec![duplicated_endpoint], base.retry().clone()); + assert!(settings.validate().is_err()); +} + +#[test] +fn transport_settings_reject_wildcard_mixed_with_specific_kind() { + let base = valid_settings("https://api.devnet.solana.com"); + let endpoint = &base.endpoints()[0]; + let role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard(), crate::HttpRequestKind::new("get_balance")], + 100, + endpoint.roles()[0].limits().clone(), + ); + let modified_endpoint = crate::HttpEndpointSettings::new( + endpoint.name(), + true, + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.url().clone(), + endpoint.connect_timeout(), + endpoint.request_timeout(), + endpoint.max_idle_connections_per_host(), + std::vec![role], + ); + let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); + assert!(settings.validate().is_err()); +} + +#[test] +fn transport_settings_reject_burst_without_rps() { + let base = valid_settings("https://api.devnet.solana.com"); + let endpoint = &base.endpoints()[0]; + let role = crate::HttpEndpointRoleSettings::new( + crate::HttpRoleName::new("default"), + true, + std::vec![crate::HttpRequestKind::wildcard()], + 100, + crate::HttpRoleLimits::new(std::option::Option::None, std::option::Option::Some(non_zero(2)), std::option::Option::None, std::option::Option::None), + ); + let modified_endpoint = crate::HttpEndpointSettings::new( + endpoint.name(), + true, + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.url().clone(), + endpoint.connect_timeout(), + endpoint.request_timeout(), + endpoint.max_idle_connections_per_host(), + std::vec![role], + ); + let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); + assert!(settings.validate().is_err()); +} + +#[test] +fn transport_settings_reject_reversed_retry_backoff() { + let base = valid_settings("https://api.devnet.solana.com"); + let settings = crate::HttpTransportSettings::new( + base.endpoints().to_vec(), + crate::HttpRetrySettings::new(2, std::time::Duration::from_secs(2), std::time::Duration::from_secs(1)), + ); + assert!(settings.validate().is_err()); +} + +#[test] +fn transport_settings_reject_zero_request_timeout() { + let base = valid_settings("https://api.devnet.solana.com"); + let endpoint = &base.endpoints()[0]; + let modified_endpoint = crate::HttpEndpointSettings::new( + endpoint.name(), + true, + endpoint.provider().clone(), + endpoint.cluster().clone(), + endpoint.url().clone(), + endpoint.connect_timeout(), + std::time::Duration::ZERO, + endpoint.max_idle_connections_per_host(), + endpoint.roles().to_vec(), + ); + let settings = crate::HttpTransportSettings::new(std::vec![modified_endpoint], base.retry().clone()); + assert!(settings.validate().is_err()); +} diff --git a/deltas/0.2.9/pre.006.md b/deltas/0.2.9/pre.006.md new file mode 100644 index 0000000..1f32892 --- /dev/null +++ b/deltas/0.2.9/pre.006.md @@ -0,0 +1,167 @@ + + + +# Delta `0.2.9-pre.006` — namespace privé HTTP explicite + +## 1. Base requise + +```text +0.2.9-pre.5.fix.1 +``` + +Le gate opérateur de `pre.005-fix.001` est intégralement vert : fmt, audit Rust, check, Clippy, Transport 364 unit + 45 public API + 38 release-completeness + 4 doctests, dependency canary Core 3/3 et workspace complet. + +## 2. Objectif + +Réduire l'ambiguïté croissante dans `ksp-onchain-transport-lib` maintenant que HTTP, WebSocket et Yellowstone gRPC coexistent dans la même crate. + +Cinq modules privés sont exclusivement propriétaires de la pile HTTP et reçoivent un préfixe explicite : + +```text +client.rs -> http_client.rs +executor.rs -> http_executor.rs +pool.rs -> http_pool.rs +resilience.rs -> http_resilience.rs +settings.rs -> http_settings.rs +``` + +Les unit tests miroirs suivent exactement les mêmes renommages. + +## 3. Frontière du renommage + +Le renommage n'est pas appliqué mécaniquement à tous les anciens modules. + +`rpc_accounts`, `rpc_blocks`, `rpc_transactions` et `rpc_common` portent déjà des DTOs/types Solana réutilisés par WebSocket et/ou gRPC. Le préfixe `http_` y serait donc architecturalement faux. Les autres `rpc_*` restent dans la même famille cohérente. + +`json_rpc` décrit une couche de protocole/enveloppe et conserve son nom. `constants` et `error` sont transverses à plusieurs transports et restent également inchangés. + +Aucun type ou nom public n'est renommé : les surfaces HTTP publiques utilisent déjà des noms `Http*` et restent réexportées depuis le crate root. + +## 4. Forecast recalibré + +L'ancien `pre.006` fonctionnel est décalé afin de ne pas mélanger ce refactor de fichiers avec l'ajout Transactions/transaction_status. + +```text +pre.006 namespace privé HTTP explicite +pre.007 Transactions + transaction_status +pre.008 Blocks + block_meta + entry +pre.009 bidi/backpressure/half-close/shutdown +pre.010 reconnect/replay/gaps/duplicates +pre.011 Config V3 + protocol/provider + profils PublicNode +pre.012 PublicNode live + compliance + docs/prompt 0.2.10 +rel.001 +``` + +Le contenu fonctionnel des tranches décalées ne change pas. + +## 5. Fichiers ajoutés par renommage + +```text +crates/ksp-onchain-transport-lib/src/http_client.rs +crates/ksp-onchain-transport-lib/src/http_executor.rs +crates/ksp-onchain-transport-lib/src/http_pool.rs +crates/ksp-onchain-transport-lib/src/http_resilience.rs +crates/ksp-onchain-transport-lib/src/http_settings.rs +crates/ksp-onchain-transport-lib/unit_tests/http_client.rs +crates/ksp-onchain-transport-lib/unit_tests/http_executor.rs +crates/ksp-onchain-transport-lib/unit_tests/http_pool.rs +crates/ksp-onchain-transport-lib/unit_tests/http_resilience.rs +crates/ksp-onchain-transport-lib/unit_tests/http_settings.rs +``` + +## 6. Fichiers supprimés + +Un ZIP overlay ne peut pas supprimer ces chemins. L'opérateur doit donc les retirer explicitement après extraction : + +```text +crates/ksp-onchain-transport-lib/src/client.rs +crates/ksp-onchain-transport-lib/src/executor.rs +crates/ksp-onchain-transport-lib/src/pool.rs +crates/ksp-onchain-transport-lib/src/resilience.rs +crates/ksp-onchain-transport-lib/src/settings.rs +crates/ksp-onchain-transport-lib/unit_tests/client.rs +crates/ksp-onchain-transport-lib/unit_tests/executor.rs +crates/ksp-onchain-transport-lib/unit_tests/pool.rs +crates/ksp-onchain-transport-lib/unit_tests/resilience.rs +crates/ksp-onchain-transport-lib/unit_tests/settings.rs +``` + +Commande opérateur : + +```bash +rm \ + crates/ksp-onchain-transport-lib/src/{client,executor,pool,resilience,settings}.rs \ + crates/ksp-onchain-transport-lib/unit_tests/{client,executor,pool,resilience,settings}.rs +``` + +## 7. Fichiers modifiés + +```text +Cargo.toml +crates/ksp-onchain-transport-lib/src/lib.rs +crates/ksp-onchain-transport-lib/tests/release_completeness.rs +docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md +docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md +deltas/0.2.9/pre.006.md +``` + +`workspace.package.version` devient : + +```text +0.2.9-pre.6 +``` + +## 8. Canaries + +Le nouveau canari release-completeness vérifie : + +- les cinq déclarations privées `mod http_*` ; +- l'absence des cinq anciennes déclarations ambiguës ; +- la présence des types/fonctions HTTP structurants dans les nouveaux fichiers ; +- le maintien volontaire des modules partagés/protocolaires `rpc_*`, `json_rpc`, `constants`, `error` ; +- l'absence de changement de surface publique requise par ce refactor. + +## 9. Documentation Markdown + +Les tableaux touchés dans `016` et `012` sont reformattés selon le comportement JetBrains RustRover : largeur de chaque colonne basée sur son contenu le plus large, puis exactement un espace de padding de part et d'autre du contenu avant les pipes. + +## 10. Validations exécutées lors de la préparation + +```text +python3 scripts/audit_rust_workspace_rules.py +General Rust rule audit: clean +Rust export completeness audit: 0 candidate(s) +KSP workspace Rust rule audit: clean +``` + +Contrôles statiques supplémentaires : + +```text +5 nouveaux modules HTTP source présents +5 nouveaux unit tests miroirs présents +5 anciens modules source absents du worktree final +5 anciens unit tests absents du worktree final +aucun renommage des modules rpc_* partagés +aucun changement de dépendance Cargo +aucune modification de public_api.rs +``` + +## 11. Validations non exécutées dans l'environnement de préparation + +Cargo/Rust ne sont pas disponibles dans l'environnement de préparation. L'opérateur doit exécuter après extraction **et suppression des anciens chemins** : + +```bash +cargo fmt --all +python3 scripts/audit_rust_workspace_rules.py +cargo check --workspace +cargo clippy --workspace --all-targets +cargo test -p ksp-onchain-transport-lib +cargo test -p ksp-core-lib --test workspace_dependencies +cargo test --workspace +``` + +Aucun `cargo tree` n'est requis : aucune dépendance ni feature ne change. + +## 12. Verdict + +`pre.006` est une candidate structurelle sans changement fonctionnel ni API publique. Sa fermeture exige l'absence effective des dix anciens chemins et un gate Cargo intégralement vert. diff --git a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md index e3978a8..c884d7e 100644 --- a/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md +++ b/docs/plans/016-V0_2_9_YELLOWSTONE_GRPC_PLAN.md @@ -1,9 +1,9 @@ - + # Plan `0.2.9` — moteur Yellowstone gRPC + standard Solana + PublicNode -> **Statut : `0.2.9-pre.004-fix.001` est fermée sur gate opérateur sans warning : fmt/audit/check/Clippy/workspace PASS, Transport 359 unit + 44 public API + 37 completeness + 4 doctests. `0.2.9-pre.005` est fonctionnellement verte sur le gate opérateur (364 unit + 45 public API + 38 completeness + workspace PASS) mais nécessite `pre.005-fix.001` pour supprimer quatre warnings `dead_code` test-only et deux diagnostics Clippy de forme dans `grpc_subscribe.rs`. Le fix ne change ni le contrat Accounts/Slots ni le wire. Transactions/Blocks, PublicNode et Config V3 restent hors tranche. `0.2.9` reste bornée à un moteur client Yellowstone partagé, une façade Solana Yellowstone standard et une première intégration concrète PublicNode. Seuls OrbitFlare puis Helius LaserStream gRPC sont actuellement planifiés comme releases provider suivantes ; les autres providers restent en TODO/IDEAS sans numéro réservé. Chaque prerelease vise 15–20 minutes de travail effectif et la release complète doit rester clôturable dans une seule session de chat.** +> **Statut : `0.2.9-pre.005-fix.001` est fermée sur gate opérateur intégralement vert et sans warning : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 38 completeness + 4 doctests. `0.2.9-pre.006` est une tranche structurelle dédiée au nommage des modules privés : les cinq modules sans ambiguïté HTTP (`client`, `executor`, `pool`, `resilience`, `settings`) deviennent `http_*`, avec leurs unit tests miroirs. Les modules `rpc_*`, `json_rpc`, `constants` et `error` ne sont pas préfixés artificiellement car ils décrivent une couche/protocole partagé ou portent déjà des DTOs utilisés par WS/gRPC. Cette insertion décale l’ancien forecast fonctionnel `pre.006–011` vers `pre.007–012` sans changer son contenu. Aucun contrat public ni comportement runtime ne change.** ## 1. Objet, base et état d'ouverture @@ -442,7 +442,7 @@ accounts_data_slice length <= 64 MiB offset + length aucun overflow u64 ``` -Les bornes account/owner/include/exclude/required, memcmp et Cuckoo restent volontairement dans les tranches de famille `pre.005–007`. +Les bornes Accounts sont matérialisées en `pre.005`; les bornes include/exclude/required et Cuckoo propres aux Transactions/Blocks restent dans `pre.007–008`. ## 8. Matrice `SubscribeUpdate` @@ -709,7 +709,7 @@ max logical filter groups/names reconnect attempts/backoff ``` -La forme JSON exacte et les bornes sont matérialisées en `pre.010`, mais **la décision V3 + `grpc_endpoints` séparés + metadata publique/secrète séparée est fermée par `pre.001`**. +La forme JSON exacte et les bornes sont matérialisées en `pre.011`, mais **la décision V3 + `grpc_endpoints` séparés + metadata publique/secrète séparée est fermée par `pre.001`**. ## 14. Audit fournisseurs gRPC gratuits et durables @@ -822,7 +822,7 @@ OrbitFlare et Helius sont validés dans leurs releases dédiées. Les autres pro Le smoke 1/2 peut vivre dans `ksp-onchain-transport-lib/tests` car il construit ses settings programmatiquement et ne teste que Transport. -Le smoke 3 ne doit pas être ajouté à `ksp-config-lib` par facilité. Si aucune surface d'intégration dédiée n'existe encore, il peut rester une procédure opérateur/documentée ou être placé sur une surface de composition déjà légitime ; le plan doit revalider l'owner au moment de `pre.010/pre.011`. +Le smoke 3 ne doit pas être ajouté à `ksp-config-lib` par facilité. Si aucune surface d'intégration dédiée n'existe encore, il peut rester une procédure opérateur/documentée ou être placé sur une surface de composition déjà légitime ; le plan doit revalider l'owner au moment de `pre.011/pre.012`. Aucun secret provider n'est versionné. @@ -922,31 +922,34 @@ pre.003 DONE — moteur TLS/metadata + façade N2 unary + fixture locale + 7 un pre.004 DONE — standard Solana : Subscribe foundation + maps/commitment/ping/from_slot/data slices/bounds budget : 15–20 min ; gate final fix.001 : fmt/audit/check/Clippy/workspace PASS sans warning + Transport 359/44/37/4 -pre.005 CANDIDATE — standard Solana : Accounts + Slots filters/updates - budget : 15–20 min ; preuve : request wire complet + Cuckoo/memcmp/predicates + Account/Slot updates exacts/malformed/adversarial +pre.005 DONE — standard Solana : Accounts + Slots filters/updates + budget : 15–20 min ; gate final fix.001 : fmt/audit/check/Clippy/workspace PASS sans warning + Transport 364/45/38/4 -pre.006 standard Solana : Transactions + transaction_status +pre.006 CANDIDATE — structure Transport : namespace privé HTTP explicite + budget : 15–20 min ; preuve : 5 modules + 5 unit tests renommés `http_*`, API publique inchangée, canari d'absence des anciens modules + +pre.007 standard Solana : Transactions + transaction_status budget : 15–20 min ; preuve : include/exclude/required/Cuckoo/token expansion + tx/meta -pre.007 standard Solana : Blocks + block_meta + entry +pre.008 standard Solana : Blocks + block_meta + entry budget : 15–20 min ; preuve : counts/arrays/optional/oneof/payload bounds -pre.008 moteur partagé : bidi mutation + Ping/Pong + half-close + backpressure + shutdown +pre.009 moteur partagé : bidi mutation + Ping/Pong + half-close + backpressure + shutdown budget : 15–20 min ; preuve : actor/session local + bounded queues + cleanup déterministe -pre.009 moteur partagé : reconnect/resubscribe + from_slot/ReplayInfo + gaps/duplicates +pre.010 moteur partagé : reconnect/resubscribe + from_slot/ReplayInfo + gaps/duplicates budget : 15–20 min ; preuve : reconnect local déterministe + aucune promesse lossless -pre.010 Config V3 + séparation protocol/provider + profils PublicNode Mainnet/Testnet +pre.011 Config V3 + séparation protocol/provider + profils PublicNode Mainnet/Testnet budget : 15–20 min ; preuve : V1/V2 backward + schema/mapping/redaction + Config -> Transport -pre.011 intégration PublicNode + smokes live + compliance + docs finales + prompt 0.2.10 OrbitFlare +pre.012 intégration PublicNode + smokes live + compliance + docs finales + prompt 0.2.10 OrbitFlare budget : 15–20 min ; preuve : Mainnet/Testnet opt-in + HTTP 52/14 + WS 18/18 + Helius + cargo graphs + workspace final rel.001 publication stable stricte ``` -Prévision : **11 prereleases**, soit environ **165–220 minutes de travail effectif nominal hors temps d'attente des commandes**, compatible avec une session complète. Si une tranche réelle excède son budget ou si `pre.011` ne peut pas raisonnablement fermer la release dans la session, on scinde avant de poursuivre au lieu de prolonger artificiellement `0.2.9`. +Prévision : **12 prereleases**, soit environ **180–240 minutes de travail effectif nominal hors temps d'attente des commandes**, compatible avec une session complète. Si une tranche réelle excède son budget ou si `pre.012` ne peut pas raisonnablement fermer la release dans la session, on scinde avant de poursuivre au lieu de prolonger artificiellement `0.2.9`. ### Critères de split @@ -1118,7 +1121,7 @@ Slot update : filters/created_at + slot + parent? + 7 SlotStatus + dead_error? Les bornes provider-neutral KSP de cette tranche couvrent notamment les sélecteurs Accounts, predicates, memcmp, Cuckoo, taille account-data, noms de filtres d'update, timestamp nanos, signature fixe 64 octets et `dead_error`. Les Debug KSP n'exposent ni pubkeys sélectionnées, ni payload memcmp/Cuckoo/account-data, ni texte `dead_error`. -Les conversions protobuf et décodeurs update restent sous `#[cfg(test)]` jusqu'à `pre.008`, car aucun stream runtime ne les consomme encore. Cette décision évite du faux `dead_code` sans créer une seconde implémentation : les mêmes helpers seront remis en runtime au moment de l'ouverture bidi. +Les conversions protobuf et décodeurs update restent sous `#[cfg(test)]` jusqu'à `pre.009`, car aucun stream runtime ne les consomme encore. Cette décision évite du faux `dead_code` sans créer une seconde implémentation : les mêmes helpers seront remis en runtime au moment de l'ouverture bidi. OUT de `pre.005` : Transactions/transaction_status, Blocks/block_meta/entry, stream bidi, provider PublicNode, Config V3 et `SubscribeDeshred`. @@ -1176,13 +1179,40 @@ Pour toute intégration provider future, la règle reste : N1 n'est jamais dupli Le premier gate opérateur de `pre.005` confirme la surface fonctionnelle : 364/364 unit, 45/45 public API, 38/38 release-completeness et workspace complet PASS. Les seuls écarts sont quatre constantes utilisées uniquement par les décodeurs `#[cfg(test)]`, une convention `to_wire(&self)` sur un type `Copy`, et une closure `is_some_and` soumise à `-D clippy::implicit-return`. -| Correction | Traitement | Impact runtime | -|-------------------------------------------|-------------------------------------------------|-------------------------| -| quatre constantes de bounds update | `#[cfg(test)]` | aucun avant `pre.008` | -| `YellowstoneSubscribeSlotFilter::to_wire` | receiver `self` | aucun, helper test-only | -| closure `dead_error.is_some_and` | `return` explicite | aucun | -| tableaux `016` et `012` | réalignement systématique des colonnes Markdown | documentaire uniquement | +| Correction | Traitement | Impact runtime | +|-------------------------------------------|----------------------------------------------------------------------|-------------------------| +| quatre constantes de bounds update | `#[cfg(test)]` | aucun avant `pre.009` | +| `YellowstoneSubscribeSlotFilter::to_wire` | receiver `self` | aucun, helper test-only | +| closure `dead_error.is_some_and` | `return` explicite | aucun | +| tableaux `016` et `012` | reformatage JetBrains RustRover (largeur max + un espace de padding) | documentaire uniquement | -Le réalignement documentaire conserve le contenu des cellules et ne modifie que les espaces de padding et les séparateurs de tableaux. +Le reformatage documentaire reproduit le comportement RustRover : largeur calculée sur le contenu le plus large de chaque colonne et exactement un espace de padding de chaque côté du contenu avant les pipes. **Verdict `pre.005-fix.001` : correctif minimal prêt ; `pre.005` reste ouverte jusqu'à réexécution sans warning de check/Clippy/workspace.** + +### 19.6 Gate final `pre.005-fix.001` + +Le second gate opérateur confirme la fermeture complète de `pre.005` : fmt, audit Rust, check et Clippy passent sans warning ; Transport passe 364 unit + 45 public API + 38 release-completeness + 4 doctests ; le dependency canary Core passe 3/3 et `cargo test --workspace` est vert. + +**Verdict : `pre.005` fermée.** + +## 20. `pre.006` — namespace privé HTTP explicite + +Le développement simultané de HTTP, WebSocket et Yellowstone gRPC rend les anciens noms privés `client`, `executor`, `pool`, `resilience` et `settings` trop ambigus. Ces cinq modules sont exclusivement propriétaires de la pile HTTP et deviennent donc : + +```text +client.rs -> http_client.rs +executor.rs -> http_executor.rs +pool.rs -> http_pool.rs +resilience.rs -> http_resilience.rs +settings.rs -> http_settings.rs +``` + +Les unit tests miroirs reçoivent les mêmes noms. Le changement reste privé à la crate : les types publics sont déjà explicitement nommés `Http*` et leurs chemins au crate root ne changent pas. + +Le mini-audit interdit un renommage aveugle des autres modules : `rpc_accounts`, `rpc_blocks`, `rpc_transactions` et `rpc_common` portent des DTOs/types déjà réutilisés par WebSocket et/ou gRPC ; `json_rpc` décrit le protocole d'enveloppe ; `constants` et `error` sont transverses. Ils conservent donc leur nom actuel. + +Cette tranche est volontairement séparée de Transactions afin de respecter le budget 15–20 minutes et d'éviter de combiner un refactor de fichiers avec une nouvelle surface protobuf. L'ancien forecast fonctionnel `pre.006–011` est décalé vers `pre.007–012` sans changement de contenu. + +**Gate candidat :** audit Rust clean, public API inchangée, cinq anciens modules absents après suppression opérateur, cinq nouveaux modules `http_*` compilés, release-completeness canary dédié, workspace complet vert. + diff --git a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md index 74bacf7..57eac00 100644 --- a/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md +++ b/docs/validation/012-V0_2_9_YELLOWSTONE_GRPC.md @@ -1,9 +1,9 @@ - + # Validation `0.2.9` — moteur Yellowstone + standard Solana + PublicNode -> **Statut : `pre.004-fix.001` est fermé sur gate opérateur intégralement vert et sans les 10 warnings `dead_code` précédents : fmt/audit/check/Clippy/workspace PASS, Transport 359 unit + 44 public API + 37 completeness + 4 doctests. `0.2.9-pre.005` est fonctionnellement verte (364 unit + 45 public API + 38 completeness + workspace PASS) mais nécessite `pre.005-fix.001` : quatre warnings `dead_code` provenant de constantes utilisées seulement par les décodeurs test-only, un warning `wrong_self_convention` et un `implicit_return` Clippy. Aucun stream bidi n'est encore ouvert ; Transactions/Blocks restent `pre.006–007`, puis lifecycle en `pre.008`. PublicNode et Config V3 restent hors tranche. OrbitFlare et Helius sont les seules releases provider suivantes planifiées ; les autres providers restent en TODO/IDEAS sans numéro réservé.** +> **Statut : `0.2.9-pre.005-fix.001` est fermé sur gate opérateur intégralement vert et sans warning : fmt/audit/check/Clippy/workspace PASS, Transport 364 unit + 45 public API + 38 completeness + 4 doctests. `pre.006` est une tranche structurelle de namespace privé : `client/executor/pool/resilience/settings` deviennent `http_*` avec leurs tests miroirs, sans changement d’API publique. Les `rpc_*`, `json_rpc`, `constants` et `error` restent volontairement non préfixés lorsqu’ils sont partagés ou déjà groupés. Le forecast fonctionnel est décalé d’un numéro : Transactions `pre.007`, Blocks `pre.008`, bidi `pre.009`, reconnect `pre.010`, Config V3 `pre.011`, PublicNode final `pre.012`.** ## 1. Autorités du gate @@ -57,7 +57,7 @@ crate proto : 12.6.0 | RPC | Forme | Classification | Scope | Preuve cible | État | |-----------------------|-------|------------------------------------------------------|-------|-----------------------|------------------------------------------------------| -| `Subscribe` | bidi | standard | IN | fixture locale + live | PARTIAL pre.005 Accounts/Slots / stream TODO pre.008 | +| `Subscribe` | bidi | standard | IN | fixture locale + live | PARTIAL pre.005 Accounts/Slots / stream TODO pre.009 | | `SubscribeDeshred` | bidi | Triton extension/pré-exécution malgré présence proto | OUT | canari d'absence/API | OUT | | `SubscribeReplayInfo` | unary | standard | IN | fixture unary | DONE pre.003 | | `Ping` | unary | standard | IN | fixture unary | DONE pre.003 | @@ -247,7 +247,7 @@ prost/prost-types 0.14.x (0.14.4 latest observé) | no raw Tonic client escape hatch | public API canary | IMPLEMENTED / Cargo pending | | moteur Yellowstone partagé sans duplication provider | source/API canary | N1 FOUNDATION IMPLEMENTED | | façade Solana Yellowstone standard distincte du moteur | public API canary | PARTIAL : 7 unary N2 / Subscribe TODO | -| PublicNode représenté comme provider/capabilities, pas comme nouveau protocole | public API/config canary | TODO `pre.010/011` | +| PublicNode représenté comme provider/capabilities, pas comme nouveau protocole | public API/config canary | TODO `pre.011/012` | | provider peut réutiliser, restreindre ou étendre N2 sans dupliquer N1 | capability/completeness review | architecture conservée | | aucune équivalence provider/standard présumée sans preuve | provider matrix/tests | architecture conservée | | façade provider spécialisée seulement si delta réel | completeness review | TODO provider integration | @@ -405,14 +405,15 @@ Les autres providers restent en TODO/IDEAS sans release dédiée. Aucune façade pre.001 DONE audit/sizing/architecture 15–20 min nominal pre.002 DONE moteur: deps/settings/errors/channel 15–20 min ; gate final fix.002 PASS pre.003 DONE TLS/metadata + fixture + 7 unary standard 15–20 min ; gate final fix.001 PASS -pre.004 CANDIDATE standard: Subscribe common/from_slot/bounds 15–20 min -pre.005 TODO standard: accounts + slots 15–20 min -pre.006 TODO standard: transactions + transaction_status 15–20 min -pre.007 TODO standard: blocks + block_meta + entry 15–20 min -pre.008 TODO moteur: bidi/backpressure/half-close/shutdown 15–20 min -pre.009 TODO moteur: reconnect/replay/gap/duplicate 15–20 min -pre.010 TODO Config V3 + protocol/provider + profils PublicNode 15–20 min -pre.011 TODO PublicNode live + compliance + docs/prompt 0.2.10 15–20 min +pre.004 DONE standard: Subscribe common/from_slot/bounds 15–20 min ; gate final fix.001 PASS +pre.005 DONE standard: accounts + slots 15–20 min ; gate final fix.001 PASS +pre.006 CANDIDATE structure: namespace privé HTTP `http_*` 15–20 min +pre.007 TODO standard: transactions + transaction_status 15–20 min +pre.008 TODO standard: blocks + block_meta + entry 15–20 min +pre.009 TODO moteur: bidi/backpressure/half-close/shutdown 15–20 min +pre.010 TODO moteur: reconnect/replay/gap/duplicate 15–20 min +pre.011 TODO Config V3 + protocol/provider + profils PublicNode 15–20 min +pre.012 TODO PublicNode live + compliance + docs/prompt 0.2.10 15–20 min rel.001 TODO stable ``` @@ -649,7 +650,7 @@ Le proto publié `yellowstone-grpc-proto 12.6.0` a été recontrôlé avant impl | Core dependency canary | 3/3 PASS | inchangé | | `cargo test --workspace` | PASS + mêmes warnings à la compilation | à réexécuter sans warning | -Les dix symboles concernés (`to_wire` et `commitment_to_wire`) servent uniquement aux tests de projection protobuf de `pre.004`. Le stream bidi n'étant pas ouvert avant `pre.008`, ils ne font pas encore partie du runtime. `#[cfg(test)]` évite donc un faux code mort de production sans créer une seconde implémentation ni modifier le contrat public. Lors de l'intégration runtime du stream, ces helpers seront naturellement retirés du `cfg(test)` au moment où ils auront un consommateur de production. +Les dix symboles concernés (`to_wire` et `commitment_to_wire`) servent uniquement aux tests de projection protobuf de `pre.004`. Le stream bidi n'étant pas ouvert avant `pre.009`, ils ne font pas encore partie du runtime. `#[cfg(test)]` évite donc un faux code mort de production sans créer une seconde implémentation ni modifier le contrat public. Lors de l'intégration runtime du stream, ces helpers seront naturellement retirés du `cfg(test)` au moment où ils auront un consommateur de production. **Verdict `pre.004-fix.001` : fermé.** @@ -693,7 +694,7 @@ Les dix symboles concernés (`to_wire` et `commitment_to_wire`) servent uniqueme | audit Rust workspace local | PASS / clean | | fmt/check/Clippy/tests | opérateur TODO | -Les conversions request protobuf et les décodeurs Account/Slot sont test-only jusqu'à `pre.008`, faute de consommateur runtime avant l'ouverture du stream. Le contrat public reste entièrement KSP-owned. +Les conversions request protobuf et les décodeurs Account/Slot sont test-only jusqu'à `pre.009`, faute de consommateur runtime avant l'ouverture du stream. Le contrat public reste entièrement KSP-owned. **Verdict `pre.005` : candidate source prête ; fermeture après gate Cargo opérateur.** @@ -714,8 +715,48 @@ Les conversions request protobuf et les décodeurs Account/Slot sont test-only j | Core dependency canary | 3/3 PASS | inchangé | | `cargo test --workspace` | PASS + 4 warnings à la compilation | à réexécuter sans warning | -Les quatre constantes concernées sont exclusivement consommées par les décodeurs Account/Slot eux-mêmes sous `#[cfg(test)]` jusqu'à l'ouverture du stream en `pre.008`. Elles passent donc sous le même `cfg(test)` plutôt que d'introduire un `allow(dead_code)`. Le receiver du helper Slot devient `self` car le type est `Copy`; la closure `is_some_and` reçoit un `return` explicite conformément à la politique Clippy KSP. +Les quatre constantes concernées sont exclusivement consommées par les décodeurs Account/Slot eux-mêmes sous `#[cfg(test)]` jusqu'à l'ouverture du stream en `pre.009`. Elles passent donc sous le même `cfg(test)` plutôt que d'introduire un `allow(dead_code)`. Le receiver du helper Slot devient `self` car le type est `Copy`; la closure `is_some_and` reçoit un `return` explicite conformément à la politique Clippy KSP. Le même fix réaligne tous les tableaux Markdown de `016` et `012` sans modifier leur contenu sémantique. **Verdict `pre.005-fix.001` : correctif source/documentaire prêt ; fermeture de `pre.005` après gate opérateur sans warning.** + +### 20.2 Gate final opérateur `pre.005-fix.001` + +| Gate | Résultat final | +|------------------------------------------|----------------| +| `cargo fmt --all` | PASS | +| audit Rust workspace | PASS / clean | +| `cargo check --workspace` | PASS | +| `cargo clippy --workspace --all-targets` | PASS | +| Transport unit | 364/364 PASS | +| Transport `public_api` | 45/45 PASS | +| Transport `release_completeness` | 38/38 PASS | +| Transport doctests | 4/4 PASS | +| Core dependency canary | 3/3 PASS | +| `cargo test --workspace` | PASS | + +**Verdict : `pre.005` fermée.** + +## 21. Gate `pre.006` — namespace privé HTTP + +| Surface / invariant | État candidate | +|----------------------------------------------------------------|----------------| +| `client.rs -> http_client.rs` | SOURCE PASS | +| `executor.rs -> http_executor.rs` | SOURCE PASS | +| `pool.rs -> http_pool.rs` | SOURCE PASS | +| `resilience.rs -> http_resilience.rs` | SOURCE PASS | +| `settings.rs -> http_settings.rs` | SOURCE PASS | +| unit tests miroirs `http_*` | SOURCE PASS | +| types/fonctions publics HTTP | INCHANGÉS | +| `rpc_common/accounts/blocks/transactions` restent partagés | SOURCE PASS | +| `json_rpc`, `constants`, `error` non artificiellement préfixés | SOURCE PASS | +| anciens fichiers supprimés après application overlay | opérateur TODO | +| release-completeness canary namespace | SOURCE PASS | +| audit Rust workspace local | PASS / clean | +| fmt/check/Clippy/tests/workspace | opérateur TODO | + +Le renommage est limité aux cinq modules dont l'ownership HTTP est sans ambiguïté. Il ne modifie aucune signature publique et n'introduit aucune dépendance. L'overlay ZIP ajoute les nouveaux chemins ; les cinq anciens fichiers source et leurs cinq unit tests miroirs doivent être supprimés explicitement par l'opérateur après extraction. + +**Verdict `pre.006` : candidate structurelle prête ; fermeture après suppressions opérateur et gate Cargo complet.** +