From 83cb861e54ee26ba34e43a4a1151cc75cd0bd98d Mon Sep 17 00:00:00 2001 From: SinuS Von SifriduS Date: Tue, 18 Aug 2026 08:17:10 +0200 Subject: [PATCH] v0.2.2-pre.005 --- Cargo.toml | 4 +- .../get_cluster_nodes.invalid_pubkey.json | 1 + .../http/get_cluster_nodes.success.json | 1 + .../fixtures/http/get_epoch_info.success.json | 1 + .../http/get_epoch_schedule.success.json | 1 + .../http/get_highest_snapshot_slot.error.json | 1 + .../get_highest_snapshot_slot.success.json | 1 + .../http/get_identity.invalid_pubkey.json | 1 + .../fixtures/http/get_identity.success.json | 1 + .../http/get_max_retransmit_slot.success.json | 1 + .../get_max_shred_insert_slot.success.json | 1 + .../src/rpc_cluster.rs | 139 ++++++++- .../tests/public_api.rs | 13 +- .../tests/release_completeness.rs | 29 +- .../unit_tests/rpc_cluster.rs | 212 +++++++++++++- deltas/0.2.2/pre.005.md | 269 ++++++++++++++++++ 16 files changed, 662 insertions(+), 14 deletions(-) create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json create mode 100644 crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json create mode 100644 deltas/0.2.2/pre.005.md diff --git a/Cargo.toml b/Cargo.toml index 4c3db8c..6321070 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,12 +1,12 @@ # file: Cargo.toml -# version: 114 +# version: 115 [workspace] resolver = "3" members = ["crates/ksp-app-config-desk", "crates/ksp-config-lib", "crates/ksp-core-lib", "crates/ksp-logging-lib", "crates/ksp-onchain-transport-lib"] [workspace.package] -version = "0.2.2-pre.4" +version = "0.2.2-pre.5" edition = "2024" license = "MIT" repository = "https://git.sasedev.com/Sasedev/khadhroony-solana-project" diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json new file mode 100644 index 0000000..6246742 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":[{"pubkey":"not-a-pubkey","rpc":"127.0.0.1:8899"}],"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json new file mode 100644 index 0000000..d97297b --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":[{"pubkey":"11111111111111111111111111111111","featureSet":3073396398,"gossip":"127.0.0.1:8001","pubsub":null,"rpc":"127.0.0.1:8899","serveRepair":"127.0.0.1:8004","shredVersion":50093,"tpu":"127.0.0.1:8003","tpuForwards":"127.0.0.1:8004","tpuForwardsQuic":"127.0.0.1:8006","tpuQuic":"127.0.0.1:8009","tpuVote":"127.0.0.1:8005","tvu":"127.0.0.1:8000","version":"4.2.1","clientId":"Agave"},{"pubkey":"Stake11111111111111111111111111111111111111"}],"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json new file mode 100644 index 0000000..a9504c4 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":{"absoluteSlot":430000001,"blockHeight":429900000,"epoch":995,"slotIndex":12345,"slotsInEpoch":432000,"transactionCount":null},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json new file mode 100644 index 0000000..4a5e333 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":{"firstNormalEpoch":0,"firstNormalSlot":0,"leaderScheduleSlotOffset":432000,"slotsPerEpoch":432000,"warmup":false},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json new file mode 100644 index 0000000..35fedfb --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","error":{"code":-32008,"message":"No snapshot"},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json new file mode 100644 index 0000000..904e0d5 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":{"full":429990000,"incremental":null},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json new file mode 100644 index 0000000..79da0f8 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":{"identity":"invalid-identity"},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json new file mode 100644 index 0000000..decc1ab --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":{"identity":"ComputeBudget111111111111111111111111111111"},"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json new file mode 100644 index 0000000..1a76768 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":430000010,"id":1} diff --git a/crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json b/crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json new file mode 100644 index 0000000..b911e61 --- /dev/null +++ b/crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json @@ -0,0 +1 @@ +{"jsonrpc":"2.0","result":430000011,"id":1} diff --git a/crates/ksp-onchain-transport-lib/src/rpc_cluster.rs b/crates/ksp-onchain-transport-lib/src/rpc_cluster.rs index 070b740..2bdf42b 100644 --- a/crates/ksp-onchain-transport-lib/src/rpc_cluster.rs +++ b/crates/ksp-onchain-transport-lib/src/rpc_cluster.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/src/rpc_cluster.rs -// version: 2 +// version: 3 /// Contact information returned for one cluster node. #[derive(Clone, Debug, Eq, PartialEq)] @@ -99,7 +99,6 @@ impl SolanaClusterNode { } /// Decodes one cluster-node contact record from the Solana JSON wire shape. - #[cfg(test)] pub(crate) fn decode_wire(method: &str, value: serde_json::Value) -> ksp_core_lib::Result { let decoded = crate::decode_wire_json::(method, value); let wire = match decoded { @@ -175,7 +174,6 @@ impl SolanaEpochInfo { } /// Decodes epoch information from the Solana JSON wire shape. - #[cfg(test)] pub(crate) fn decode_wire(method: &str, value: serde_json::Value) -> ksp_core_lib::Result { let decoded = crate::decode_wire_json::(method, value); return match decoded { @@ -230,7 +228,6 @@ impl SolanaEpochSchedule { } /// Decodes an epoch schedule from the Solana JSON wire shape. - #[cfg(test)] pub(crate) fn decode_wire(method: &str, value: serde_json::Value) -> ksp_core_lib::Result { let decoded = crate::decode_wire_json::(method, value); return match decoded { @@ -266,7 +263,6 @@ impl SolanaSnapshotSlotInfo { } /// Decodes snapshot-slot information from the Solana JSON wire shape. - #[cfg(test)] pub(crate) fn decode_wire(method: &str, value: serde_json::Value) -> ksp_core_lib::Result { let decoded = crate::decode_wire_json::(method, value); return match decoded { @@ -617,7 +613,135 @@ impl SolanaVoteAccountStatus { } } -#[cfg(test)] +#[derive(serde::Deserialize)] +struct WireIdentity { + identity: std::string::String, +} + +impl crate::HttpTransportPool { + /// Executes typed `getClusterNodes` through the common KSP HTTP transport path. + pub async fn get_cluster_nodes(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result> { + let value = self.execute_cluster_simple_rpc("getClusterNodes", role, std::vec::Vec::new()).await; + let value = match value { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let decoded = crate::decode_wire_json::>("getClusterNodes", value); + let values = match decoded { + std::result::Result::Ok(values) => values, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let mut nodes = std::vec::Vec::with_capacity(values.len()); + for value in values { + let node = crate::SolanaClusterNode::decode_wire("getClusterNodes", value); + match node { + std::result::Result::Ok(node) => nodes.push(node), + std::result::Result::Err(error) => return std::result::Result::Err(error), + } + } + return std::result::Result::Ok(nodes); + } + + /// Executes typed `getEpochInfo` through the common KSP HTTP transport path. + pub async fn get_epoch_info( + &self, + role: &crate::HttpRoleName, + config: std::option::Option<&crate::SolanaContextConfig>, + ) -> ksp_core_lib::Result { + let mut params = std::vec::Vec::new(); + if let std::option::Option::Some(config) = config + && (config.commitment().is_some() || config.min_context_slot().is_some()) + { + params.push(config.to_json_value()); + } + let value = self.execute_cluster_simple_rpc("getEpochInfo", role, params).await; + return match value { + std::result::Result::Ok(value) => crate::SolanaEpochInfo::decode_wire("getEpochInfo", value), + std::result::Result::Err(error) => std::result::Result::Err(error), + }; + } + + /// Executes typed `getEpochSchedule` through the common KSP HTTP transport path. + pub async fn get_epoch_schedule(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + let value = self.execute_cluster_simple_rpc("getEpochSchedule", role, std::vec::Vec::new()).await; + return match value { + std::result::Result::Ok(value) => crate::SolanaEpochSchedule::decode_wire("getEpochSchedule", value), + std::result::Result::Err(error) => std::result::Result::Err(error), + }; + } + + /// Executes typed `getHighestSnapshotSlot` through the common KSP HTTP transport path. + pub async fn get_highest_snapshot_slot(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + let value = self.execute_cluster_simple_rpc("getHighestSnapshotSlot", role, std::vec::Vec::new()).await; + return match value { + std::result::Result::Ok(value) => crate::SolanaSnapshotSlotInfo::decode_wire("getHighestSnapshotSlot", value), + std::result::Result::Err(error) => std::result::Result::Err(error), + }; + } + + /// Executes typed `getIdentity` through the common KSP HTTP transport path. + pub async fn get_identity(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + let value = self.execute_cluster_simple_rpc("getIdentity", role, std::vec::Vec::new()).await; + let value = match value { + std::result::Result::Ok(value) => value, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + let decoded = crate::decode_wire_json::("getIdentity", value); + let wire = match decoded { + std::result::Result::Ok(wire) => wire, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + return crate::parse_wire_pubkey("getIdentity", "identity", wire.identity.as_str()); + } + + /// Executes typed `getMaxRetransmitSlot` through the common KSP HTTP transport path. + pub async fn get_max_retransmit_slot(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + return self.get_cluster_simple_slot("getMaxRetransmitSlot", role).await; + } + + /// Executes typed `getMaxShredInsertSlot` through the common KSP HTTP transport path. + pub async fn get_max_shred_insert_slot(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + return self.get_cluster_simple_slot("getMaxShredInsertSlot", role).await; + } + + async fn get_cluster_simple_slot(&self, method_name: &'static str, role: &crate::HttpRoleName) -> ksp_core_lib::Result { + let value = self.execute_cluster_simple_rpc(method_name, role, std::vec::Vec::new()).await; + return match value { + std::result::Result::Ok(value) => crate::decode_wire_json::(method_name, value), + std::result::Result::Err(error) => std::result::Result::Err(error), + }; + } + + async fn execute_cluster_simple_rpc( + &self, + method_name: &'static str, + role: &crate::HttpRoleName, + params: std::vec::Vec, + ) -> ksp_core_lib::Result { + let method = cluster_descriptor(method_name); + let method = match method { + std::result::Result::Ok(method) => method, + std::result::Result::Err(error) => return std::result::Result::Err(error), + }; + return self.execute_standard_rpc(role, method, params).await; + } +} + +fn cluster_descriptor(method: &str) -> ksp_core_lib::Result<&'static crate::HttpRpcMethodDescriptor> { + let descriptor = crate::find_http_rpc_method(method); + return match descriptor { + std::option::Option::Some(descriptor) + if descriptor.category() == crate::HttpRpcCategory::Cluster && descriptor.coverage_release() == crate::HttpRpcCoverageRelease::V0_2_2 => + { + std::result::Result::Ok(descriptor) + }, + _ => std::result::Result::Err( + ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RESPONSE, "typed Cluster descriptor is missing from the audited 0.2.2 registry") + .with_context("rpc_method", method), + ), + }; +} + #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] struct WireClusterNode { @@ -652,7 +776,6 @@ struct WireClusterNode { client_id: std::option::Option, } -#[cfg(test)] #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] struct WireEpochInfo { @@ -664,7 +787,6 @@ struct WireEpochInfo { transaction_count: std::option::Option, } -#[cfg(test)] #[derive(serde::Deserialize)] #[serde(rename_all = "camelCase")] struct WireEpochSchedule { @@ -675,7 +797,6 @@ struct WireEpochSchedule { warmup: bool, } -#[cfg(test)] #[derive(serde::Deserialize)] struct WireSnapshotSlotInfo { full: u64, diff --git a/crates/ksp-onchain-transport-lib/tests/public_api.rs b/crates/ksp-onchain-transport-lib/tests/public_api.rs index 966cb6e..f8d49ff 100644 --- a/crates/ksp-onchain-transport-lib/tests/public_api.rs +++ b/crates/ksp-onchain-transport-lib/tests/public_api.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/tests/public_api.rs -// version: 7 +// version: 8 //! Integration tests for the public `ksp-onchain-transport-lib` consumer contract. @@ -221,3 +221,14 @@ fn public_pre_004_token_wrappers_are_available_from_crate_root() { let selector = ksp_onchain_transport_lib::SolanaTokenAccountSelector::Mint(mint); assert!(matches!(selector, ksp_onchain_transport_lib::SolanaTokenAccountSelector::Mint(_))); } + +#[test] +fn public_pre_005_simple_cluster_wrappers_are_available_from_crate_root() { + let _get_cluster_nodes = ksp_onchain_transport_lib::HttpTransportPool::get_cluster_nodes; + let _get_epoch_info = ksp_onchain_transport_lib::HttpTransportPool::get_epoch_info; + let _get_epoch_schedule = ksp_onchain_transport_lib::HttpTransportPool::get_epoch_schedule; + let _get_highest_snapshot_slot = ksp_onchain_transport_lib::HttpTransportPool::get_highest_snapshot_slot; + let _get_identity = ksp_onchain_transport_lib::HttpTransportPool::get_identity; + let _get_max_retransmit_slot = ksp_onchain_transport_lib::HttpTransportPool::get_max_retransmit_slot; + let _get_max_shred_insert_slot = ksp_onchain_transport_lib::HttpTransportPool::get_max_shred_insert_slot; +} diff --git a/crates/ksp-onchain-transport-lib/tests/release_completeness.rs b/crates/ksp-onchain-transport-lib/tests/release_completeness.rs index 1258b26..e48a8ed 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: 3 +// version: 4 //! Release-level completeness canaries for the `0.2.1` HTTP foundation contract. @@ -84,3 +84,30 @@ fn release_pre_004_tokens_subset_is_exact_and_retry_safe() { std::vec!["getTokenAccountBalance", "getTokenAccountsByDelegate", "getTokenAccountsByOwner", "getTokenLargestAccounts", "getTokenSupply",], ); } + +#[test] +fn release_pre_005_simple_cluster_subset_is_exact_and_retry_safe() { + let expected = std::vec![ + "getClusterNodes", + "getEpochInfo", + "getEpochSchedule", + "getHighestSnapshotSlot", + "getIdentity", + "getMaxRetransmitSlot", + "getMaxShredInsertSlot", + ]; + let deferred = std::vec!["getLeaderSchedule", "getSlot", "getSlotLeader", "getSlotLeaders", "getVoteAccounts"]; + for method_name in &expected { + let descriptor = ksp_onchain_transport_lib::find_http_rpc_method(method_name).expect("pre.005 cluster descriptor must exist"); + assert_eq!(descriptor.category(), ksp_onchain_transport_lib::HttpRpcCategory::Cluster); + assert_eq!(descriptor.coverage_release(), ksp_onchain_transport_lib::HttpRpcCoverageRelease::V0_2_2); + assert_eq!(descriptor.operation_kind(), ksp_onchain_transport_lib::RpcOperationKind::Read); + assert_eq!(descriptor.transport_retry_class(), ksp_onchain_transport_lib::TransportRetryClass::RetrySafe); + } + for method_name in &deferred { + let descriptor = ksp_onchain_transport_lib::find_http_rpc_method(method_name).expect("pre.006 cluster descriptor must remain registered"); + assert_eq!(descriptor.coverage_release(), ksp_onchain_transport_lib::HttpRpcCoverageRelease::V0_2_2); + } + assert_eq!(expected.len(), 7); + assert_eq!(deferred.len(), 5); +} diff --git a/crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs b/crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs index fd304ca..b8706ef 100644 --- a/crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs +++ b/crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs @@ -1,5 +1,5 @@ // file: crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs -// version: 2 +// version: 3 #[test] fn cluster_node_fixture_preserves_v4_client_id_and_optional_fields() { @@ -111,3 +111,213 @@ fn staged_vote_status_and_config_helpers_preserve_wire_shapes() { assert_eq!(status.current().len(), 1); assert!(status.delinquent().is_empty()); } + +fn pool_for_url(url: &str) -> 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::None), + ); + 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_secs(1), + std::time::Duration::from_secs(1), + std::option::Option::Some(1), + std::vec![role], + ); + let settings = crate::HttpTransportSettings::new( + std::vec![endpoint], + crate::HttpRetrySettings::new(0, std::time::Duration::from_millis(1), std::time::Duration::from_millis(1)), + ); + return crate::HttpTransportPool::new(settings).expect("fixture pool must build"); +} + +fn serve_once(body: &'static str) -> (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 one request"); + let request = read_request(&mut stream); + let response = 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"); + return request; + }); + 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); +} + +fn request_body(request: &str) -> serde_json::Value { + let body = request.split("\r\n\r\n").nth(1).expect("fixture request body must exist"); + return serde_json::from_str(body).expect("fixture request body must be JSON"); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_cluster_nodes_preserves_optional_v4_fields_and_omissions() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_cluster_nodes.success.json")); + let pool = pool_for_url(url.as_str()); + let nodes = pool.get_cluster_nodes(&crate::HttpRoleName::new("default")).await.expect("cluster nodes fixture must succeed"); + assert_eq!(nodes.len(), 2); + assert_eq!(nodes[0].client_id(), std::option::Option::Some("Agave")); + assert_eq!(nodes[0].serve_repair(), std::option::Option::Some("127.0.0.1:8004")); + assert_eq!(nodes[1].rpc(), std::option::Option::None); + assert_eq!(nodes[1].client_id(), std::option::Option::None); + let request = handle.join().expect("fixture server must join"); + let body = request_body(request.as_str()); + assert_eq!(body["method"], serde_json::json!("getClusterNodes")); + assert_eq!(body["params"], serde_json::json!([])); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_cluster_nodes_rejects_invalid_node_pubkey() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_cluster_nodes.invalid_pubkey.json")); + let pool = pool_for_url(url.as_str()); + let result = pool.get_cluster_nodes(&crate::HttpRoleName::new("default")).await; + let error = result.expect_err("invalid cluster node pubkey must reject typed response"); + assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE); + handle.join().expect("fixture server must join"); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_epoch_info_serializes_context_config_and_preserves_nullable_transaction_count() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_epoch_info.success.json")); + let pool = pool_for_url(url.as_str()); + let config = crate::SolanaContextConfig::new(std::option::Option::Some(crate::SolanaCommitment::Finalized), std::option::Option::Some(429_000_000)); + let info = pool + .get_epoch_info(&crate::HttpRoleName::new("default"), std::option::Option::Some(&config)) + .await + .expect("epoch info fixture must succeed"); + assert_eq!(info.absolute_slot(), 430_000_001); + assert_eq!(info.transaction_count(), std::option::Option::None); + let request = handle.join().expect("fixture server must join"); + let body = request_body(request.as_str()); + assert_eq!(body["params"], serde_json::json!([{"commitment":"finalized","minContextSlot":429000000}])); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_epoch_info_omits_explicitly_empty_config() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_epoch_info.success.json")); + let pool = pool_for_url(url.as_str()); + let config = crate::SolanaContextConfig::default(); + let info = pool + .get_epoch_info(&crate::HttpRoleName::new("default"), std::option::Option::Some(&config)) + .await + .expect("epoch info fixture must succeed"); + assert_eq!(info.epoch(), 995); + let request = handle.join().expect("fixture server must join"); + assert_eq!(request_body(request.as_str())["params"], serde_json::json!([])); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_epoch_schedule_preserves_fixed_wire_shape() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_epoch_schedule.success.json")); + let pool = pool_for_url(url.as_str()); + let schedule = pool.get_epoch_schedule(&crate::HttpRoleName::new("default")).await.expect("epoch schedule fixture must succeed"); + assert_eq!(schedule.slots_per_epoch(), 432_000); + assert!(!schedule.warmup()); + let request = handle.join().expect("fixture server must join"); + assert_eq!(request_body(request.as_str())["method"], serde_json::json!("getEpochSchedule")); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_highest_snapshot_slot_preserves_nullable_incremental_slot() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_highest_snapshot_slot.success.json")); + let pool = pool_for_url(url.as_str()); + let snapshot = pool.get_highest_snapshot_slot(&crate::HttpRoleName::new("default")).await.expect("snapshot fixture must succeed"); + assert_eq!(snapshot.full(), 429_990_000); + assert_eq!(snapshot.incremental(), std::option::Option::None); + let request = handle.join().expect("fixture server must join"); + assert_eq!(request_body(request.as_str())["params"], serde_json::json!([])); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_highest_snapshot_slot_preserves_no_snapshot_rpc_error() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_highest_snapshot_slot.error.json")); + let pool = pool_for_url(url.as_str()); + let result = pool.get_highest_snapshot_slot(&crate::HttpRoleName::new("default")).await; + let error = result.expect_err("no snapshot must remain an RPC application error"); + assert_eq!(error.code(), crate::ERROR_CODE_RPC_APPLICATION_ERROR); + handle.join().expect("fixture server must join"); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_identity_decodes_pubkey_object() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_identity.success.json")); + let pool = pool_for_url(url.as_str()); + let identity = pool.get_identity(&crate::HttpRoleName::new("default")).await.expect("identity fixture must succeed"); + assert_eq!(identity.to_string(), "ComputeBudget111111111111111111111111111111"); + let request = handle.join().expect("fixture server must join"); + assert_eq!(request_body(request.as_str())["method"], serde_json::json!("getIdentity")); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_identity_rejects_invalid_wire_pubkey() { + let (url, handle) = serve_once(include_str!("../fixtures/http/get_identity.invalid_pubkey.json")); + let pool = pool_for_url(url.as_str()); + let result = pool.get_identity(&crate::HttpRoleName::new("default")).await; + let error = result.expect_err("invalid identity pubkey must reject typed response"); + assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE); + handle.join().expect("fixture server must join"); +} + +#[tokio::test(flavor = "current_thread")] +async fn typed_get_max_cluster_slots_decode_u64_without_params() { + let (retransmit_url, retransmit_handle) = serve_once(include_str!("../fixtures/http/get_max_retransmit_slot.success.json")); + let retransmit_pool = pool_for_url(retransmit_url.as_str()); + let retransmit = retransmit_pool.get_max_retransmit_slot(&crate::HttpRoleName::new("default")).await.expect("max retransmit slot fixture must succeed"); + assert_eq!(retransmit, 430_000_010); + let retransmit_request = retransmit_handle.join().expect("fixture server must join"); + let retransmit_body = request_body(retransmit_request.as_str()); + assert_eq!(retransmit_body["method"], serde_json::json!("getMaxRetransmitSlot")); + assert_eq!(retransmit_body["params"], serde_json::json!([])); + let (shred_url, shred_handle) = serve_once(include_str!("../fixtures/http/get_max_shred_insert_slot.success.json")); + let shred_pool = pool_for_url(shred_url.as_str()); + let shred = shred_pool.get_max_shred_insert_slot(&crate::HttpRoleName::new("default")).await.expect("max shred insert slot fixture must succeed"); + assert_eq!(shred, 430_000_011); + let shred_request = shred_handle.join().expect("fixture server must join"); + let shred_body = request_body(shred_request.as_str()); + assert_eq!(shred_body["method"], serde_json::json!("getMaxShredInsertSlot")); + assert_eq!(shred_body["params"], serde_json::json!([])); +} diff --git a/deltas/0.2.2/pre.005.md b/deltas/0.2.2/pre.005.md new file mode 100644 index 0000000..12ee893 --- /dev/null +++ b/deltas/0.2.2/pre.005.md @@ -0,0 +1,269 @@ + + + +# Delta `0.2.2-pre.005` — sept wrappers HTTP Cluster simples typés + +## Base requise + +Livraison précédente validée localement par l'opérateur : + +```text +0.2.2-pre.004 +workspace.package.version = "0.2.2-pre.4" +``` + +La validation opérateur du 2026-08-18 a confirmé : + +```text +cargo fmt --all OK +cargo check --workspace OK +cargo clippy --workspace --all-targets OK +cargo test -p ksp-onchain-transport-lib OK +``` + +Résultats Transport de cette base : + +```text +106 unit tests +11 public API tests +4 release completeness tests +0 warning signalé par check/clippy +``` + +Le plan canonique reste `docs/plans/009-V0_2_2_HTTP_ACCOUNTS_TOKENS_CLUSTER_PLAN.md` version 3. + +## Contrôle documentaire avant nouvelle tranche + +Le contrôle de cohérence demandé entre prereleases confirme : + +- `ROADMAP.md` conserve correctement `0.2.2` en cours `[/]` ; +- le plan `009` attribue exactement à `pre.005` les sept méthodes Cluster simples implémentées ici ; +- `pre.006` conserve exactement cinq méthodes différées : `getLeaderSchedule`, `getSlot`, `getSlotLeader`, `getSlotLeaders`, `getVoteAccounts` ; +- les index `docs/000-README.md`, `docs/plans/000-README.md` et `docs/plans/002-FUNCTIONAL_RELEASE_SEQUENCE.md` restent cohérents et ne nécessitent aucune modification. + +## Objectif + +Implémenter les sept wrappers publics typés Cluster simples prévus par le plan : + +```text +getClusterNodes +getEpochInfo +getEpochSchedule +getHighestSnapshotSlot +getIdentity +getMaxRetransmitSlot +getMaxShredInsertSlot +``` + +Tous passent par la foundation HTTP commune : + +```text +wrapper typé + -> descriptor audité + -> execute_standard_rpc + -> pool/admission/retry/deadline + -> reqwest HTTP + -> JSON-RPC validation + -> décodage DTO KSP +``` + +Aucun wrapper Cluster ne contacte `reqwest` directement. + +## Version Cargo + +Conformément au cycle prerelease KSP : + +```text +0.2.2-pre.4 -> 0.2.2-pre.5 +``` + +Aucune dépendance ni feature Cargo n'est ajoutée ou retirée. + +## Réaudit ciblé du contrat courant + +Les pages RPC Solana courantes ont été revérifiées le 2026-08-18 : + +```text +https://solana.com/docs/rpc/http/getclusternodes +https://solana.com/docs/rpc/http/getepochinfo +https://solana.com/docs/rpc/http/getepochschedule +https://solana.com/docs/rpc/http/gethighestsnapshotslot +https://solana.com/docs/rpc/http/getidentity +https://solana.com/docs/rpc/http/getmaxretransmitslot +https://solana.com/docs/rpc/http/getmaxshredinsertslot +``` + +Le contrat reste cohérent avec le plan `009` : + +- `getClusterNodes` ne prend aucun paramètre et renvoie une liste de contacts de noeuds ; +- `getEpochInfo` accepte uniquement la config commune `commitment/minContextSlot` et conserve `transactionCount` nullable ; +- `getEpochSchedule` ne prend aucun paramètre et renvoie la structure fixe d'epoch schedule ; +- `getHighestSnapshotSlot` ne prend aucun paramètre, conserve `incremental` nullable et laisse l'absence de snapshot comme erreur JSON-RPC distante ; +- `getIdentity` ne prend aucun paramètre et renvoie l'objet `{ identity }`, converti en `Pubkey` KSP ; +- `getMaxRetransmitSlot` et `getMaxShredInsertSlot` ne prennent aucun paramètre et renvoient chacun un `u64`. + +Le champ Agave `v4.2.1` `clientId: Option` identifié par `pre.001-fix.001` reste conservé par `SolanaClusterNode`, sans devenir obligatoire. + +## Surface typée Cluster simple + +### `getClusterNodes` + +```text +role -> Vec +``` + +Chaque `pubkey` est validée et convertie en `ksp_core_lib::Pubkey`. Les endpoints réseau restent des chaînes optionnelles. Les champs optionnels absents, y compris `clientId`, restent `None`. + +### `getEpochInfo` + +```text +role + Option -> SolanaEpochInfo +``` + +Un config explicitement vide est omis de `params`. `transactionCount: null` reste `None`. + +### `getEpochSchedule` + +```text +role -> SolanaEpochSchedule +``` + +Aucune transformation métier n'est introduite. + +### `getHighestSnapshotSlot` + +```text +role -> SolanaSnapshotSlotInfo +``` + +`incremental` reste optionnel. Une erreur RPC « no snapshot » n'est pas convertie en valeur sentinelle. + +### `getIdentity` + +```text +role -> Pubkey +``` + +Le wrapper valide la chaîne `identity` avant de l'exposer comme `Pubkey`, sans recopier une valeur wire invalide dans le diagnostic. + +### `getMaxRetransmitSlot` / `getMaxShredInsertSlot` + +```text +role -> u64 +``` + +Les deux wrappers partagent un chemin privé simple de décodage `u64` et ne créent aucun DTO artificiel. + +## Activation ciblée des helpers préparés en `pre.002` + +Les helpers Cluster nécessaires à cette tranche deviennent production-live uniquement maintenant qu'ils ont des consommateurs runtime : + +- `SolanaClusterNode::decode_wire` + `WireClusterNode` ; +- `SolanaEpochInfo::decode_wire` + `WireEpochInfo` ; +- `SolanaEpochSchedule::decode_wire` + `WireEpochSchedule` ; +- `SolanaSnapshotSlotInfo::decode_wire` + `WireSnapshotSlotInfo`. + +Les helpers spécifiques à `getLeaderSchedule` et `getVoteAccounts` restent sous `#[cfg(test)]` jusqu'à `pre.006`. Aucun `#[allow(dead_code)]` n'est ajouté. + +## Fixtures HTTP déterministes ajoutées + +```text +crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json +crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json +crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json +crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json +``` + +Les tests utilisent uniquement un serveur HTTP loopback local. + +## Couverture de tests ajoutée + +Les tests couvrent notamment : + +- requête sans paramètres pour les six méthodes sans config ; +- sérialisation exacte `commitment + minContextSlot` pour `getEpochInfo` ; +- omission d'un `SolanaContextConfig` explicitement vide ; +- `transactionCount: null` ; +- champs ClusterNode optionnels présents et absents ; +- conservation de `clientId` Agave `v4.2.1` ; +- rejet d'une `pubkey` de noeud invalide ; +- `incremental: null` ; +- propagation d'une erreur RPC « no snapshot » ; +- validation de l'identity `Pubkey` ; +- décodage des deux slots max comme `u64`. + +Le test public compile explicitement les sept nouvelles méthodes depuis `HttpTransportPool`. + +Un canari release fige le sous-ensemble `pre.005` exact et vérifie `Read + RetrySafe` pour les sept descriptors, tout en laissant les cinq méthodes Cluster de `pre.006` enregistrées mais différées. + +Après application, la cible Transport attendue devient : + +```text +116 unit tests +12 public API tests +5 release completeness tests +``` + +## Fichiers ajoutés + +```text +crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_cluster_nodes.invalid_pubkey.json +crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_info.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_epoch_schedule.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_highest_snapshot_slot.error.json +crates/ksp-onchain-transport-lib/fixtures/http/get_identity.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_identity.invalid_pubkey.json +crates/ksp-onchain-transport-lib/fixtures/http/get_max_retransmit_slot.success.json +crates/ksp-onchain-transport-lib/fixtures/http/get_max_shred_insert_slot.success.json +deltas/0.2.2/pre.005.md +``` + +## Fichiers modifiés + +```text +Cargo.toml +crates/ksp-onchain-transport-lib/src/rpc_cluster.rs +crates/ksp-onchain-transport-lib/unit_tests/rpc_cluster.rs +crates/ksp-onchain-transport-lib/tests/public_api.rs +crates/ksp-onchain-transport-lib/tests/release_completeness.rs +``` + +## Fichiers supprimés + +Aucun. + +## Documentation durable + +Aucune modification de ROADMAP/plan/index n'est requise par cette tranche après le contrôle documentaire : les documents existants décrivent déjà exactement ce découpage. + +`CHANGELOG.md` reste réservé à la clôture stable. + +## Validation à exécuter sur le checkout opérateur + +```bash +cargo fmt --all +cargo check --workspace +cargo clippy --workspace --all-targets +cargo test -p ksp-onchain-transport-lib +``` + +La livraison n'affirme pas que ces commandes ont été exécutées dans l'environnement de génération. + +## Critères de validation de `pre.005` + +`pre.005` est validable lorsque : + +- les quatre commandes ci-dessus passent sans warning nouveau ; +- les 116 unit tests passent ; +- les 12 tests public API passent ; +- les 5 tests release completeness passent ; +- aucun helper Cluster simple production-live n'est `dead_code` ; +- les cinq wrappers complexes de `pre.006` restent hors de cette tranche.