v0.2.4-pre.003
This commit is contained in:
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":{"commitment":[1,2,3,4],"totalStake":1000000000},"id":1}
|
||||
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":410000000,"id":1}
|
||||
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":null,"id":1}
|
||||
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":1787072400,"id":1}
|
||||
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":250000,"id":1}
|
||||
@@ -0,0 +1 @@
|
||||
{"jsonrpc":"2.0","result":123456,"id":1}
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/lib.rs
|
||||
// version: 17
|
||||
// version: 18
|
||||
#![warn(missing_docs)]
|
||||
#![deny(unreachable_pub)]
|
||||
#![forbid(unsafe_code)]
|
||||
@@ -12,8 +12,8 @@
|
||||
//! are available. The four typed Solana HTTP foundation canaries plus all 22 typed `0.2.2` Accounts, Tokens and Cluster wrappers execute real JSON-RPC
|
||||
//! requests through the shared transport path. `0.2.3` exposes its shared Transaction wire/config primitives and all eleven Transaction wrappers through
|
||||
//! `pre.007`: eight reads, two write submissions with centralized no-resend protection, and retry-safe `simulateTransaction`, including complete
|
||||
//! modern/legacy `getTransaction` coverage. `0.2.4-pre.002` adds the shared Blocks/Economics wire, config and result primitives without yet advancing
|
||||
//! any `V0_2_4` method to a typed wrapper.
|
||||
//! modern/legacy `getTransaction` coverage. `0.2.4-pre.002` adds the shared Blocks/Economics wire, config and result primitives. `0.2.4-pre.003`
|
||||
//! activates the first five typed Blocks reads: `getBlockCommitment`, `getBlockHeight`, `getBlockTime`, `getFirstAvailableBlock` and `minimumLedgerSlot`.
|
||||
|
||||
mod client;
|
||||
mod constants;
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/src/rpc_blocks.rs
|
||||
// version: 2
|
||||
// version: 3
|
||||
|
||||
/// Transaction detail level accepted by modern `getBlock` requests.
|
||||
#[derive(Clone, Copy, Debug, Default, Eq, Hash, PartialEq)]
|
||||
@@ -236,7 +236,6 @@ impl SolanaBlockCommitment {
|
||||
}
|
||||
|
||||
/// Decodes a block-commitment result from its Solana JSON wire shape.
|
||||
#[cfg(test)]
|
||||
pub(crate) fn decode_wire(method: &str, value: serde_json::Value) -> ksp_core_lib::Result<Self> {
|
||||
let decoded = crate::decode_wire_json::<WireBlockCommitment>(method, value);
|
||||
return match decoded {
|
||||
@@ -587,6 +586,94 @@ impl SolanaPerformanceSample {
|
||||
}
|
||||
}
|
||||
|
||||
impl crate::HttpTransportPool {
|
||||
/// Executes typed `getBlockCommitment` through the common KSP HTTP transport path.
|
||||
pub async fn get_block_commitment(&self, role: &crate::HttpRoleName, slot: u64) -> ksp_core_lib::Result<crate::SolanaBlockCommitment> {
|
||||
let value = self.execute_blocks_rpc("getBlockCommitment", role, std::vec![serde_json::json!(slot)]).await;
|
||||
return match value {
|
||||
std::result::Result::Ok(value) => crate::SolanaBlockCommitment::decode_wire("getBlockCommitment", value),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(error),
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes typed `getBlockHeight` through the common KSP HTTP transport path.
|
||||
pub async fn get_block_height(&self, role: &crate::HttpRoleName, config: std::option::Option<&crate::SolanaContextConfig>) -> ksp_core_lib::Result<u64> {
|
||||
let mut params = std::vec::Vec::new();
|
||||
push_blocks_context_config(&mut params, config);
|
||||
return self.get_blocks_u64("getBlockHeight", role, params).await;
|
||||
}
|
||||
|
||||
/// Executes typed `getBlockTime` and preserves a runtime `null` as `None`.
|
||||
pub async fn get_block_time(&self, role: &crate::HttpRoleName, slot: u64) -> ksp_core_lib::Result<std::option::Option<i64>> {
|
||||
let value = self.execute_blocks_rpc("getBlockTime", role, std::vec![serde_json::json!(slot)]).await;
|
||||
return match value {
|
||||
std::result::Result::Ok(value) => crate::decode_wire_json::<std::option::Option<i64>>("getBlockTime", value),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(error),
|
||||
};
|
||||
}
|
||||
|
||||
/// Executes typed `getFirstAvailableBlock` through the common KSP HTTP transport path.
|
||||
pub async fn get_first_available_block(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result<u64> {
|
||||
return self.get_blocks_u64("getFirstAvailableBlock", role, std::vec::Vec::new()).await;
|
||||
}
|
||||
|
||||
/// Executes typed `minimumLedgerSlot` through the common KSP HTTP transport path.
|
||||
pub async fn minimum_ledger_slot(&self, role: &crate::HttpRoleName) -> ksp_core_lib::Result<u64> {
|
||||
return self.get_blocks_u64("minimumLedgerSlot", role, std::vec::Vec::new()).await;
|
||||
}
|
||||
|
||||
async fn get_blocks_u64(
|
||||
&self,
|
||||
method_name: &'static str,
|
||||
role: &crate::HttpRoleName,
|
||||
params: std::vec::Vec<serde_json::Value>,
|
||||
) -> ksp_core_lib::Result<u64> {
|
||||
let value = self.execute_blocks_rpc(method_name, role, params).await;
|
||||
return match value {
|
||||
std::result::Result::Ok(value) => crate::decode_wire_json::<u64>(method_name, value),
|
||||
std::result::Result::Err(error) => std::result::Result::Err(error),
|
||||
};
|
||||
}
|
||||
|
||||
async fn execute_blocks_rpc(
|
||||
&self,
|
||||
method_name: &'static str,
|
||||
role: &crate::HttpRoleName,
|
||||
params: std::vec::Vec<serde_json::Value>,
|
||||
) -> ksp_core_lib::Result<serde_json::Value> {
|
||||
let method = blocks_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 push_blocks_context_config(params: &mut std::vec::Vec<serde_json::Value>, config: std::option::Option<&crate::SolanaContextConfig>) {
|
||||
if let std::option::Option::Some(config) = config
|
||||
&& (config.commitment().is_some() || config.min_context_slot().is_some())
|
||||
{
|
||||
params.push((*config).to_json_value());
|
||||
}
|
||||
return;
|
||||
}
|
||||
|
||||
fn blocks_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::Blocks && descriptor.coverage_release() == crate::HttpRpcCoverageRelease::V0_2_4 =>
|
||||
{
|
||||
std::result::Result::Ok(descriptor)
|
||||
},
|
||||
_ => std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_INVALID_RESPONSE, "typed Blocks descriptor is missing from the audited 0.2.4 registry")
|
||||
.with_context("rpc_method", method),
|
||||
),
|
||||
};
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn decode_block_transaction_version(
|
||||
method: &str,
|
||||
@@ -647,7 +734,6 @@ fn decode_block_rewards(
|
||||
return std::result::Result::Ok(crate::SolanaWireField::Value(rewards));
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[derive(serde::Deserialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
struct WireBlockCommitment {
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/tests/public_api.rs
|
||||
// version: 16
|
||||
// version: 17
|
||||
|
||||
//! Integration tests for the public `ksp-onchain-transport-lib` consumer contract.
|
||||
|
||||
@@ -386,3 +386,12 @@ fn public_v0_2_4_pre_002_blocks_economics_types_are_available_from_crate_root()
|
||||
let _reward: std::option::Option<ksp_onchain_transport_lib::SolanaInflationReward> = std::option::Option::None;
|
||||
let _supply: std::option::Option<ksp_onchain_transport_lib::SolanaSupply> = std::option::Option::None;
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn public_v0_2_4_pre_003_simple_block_wrappers_are_available_from_crate_root() {
|
||||
let _get_block_commitment = ksp_onchain_transport_lib::HttpTransportPool::get_block_commitment;
|
||||
let _get_block_height = ksp_onchain_transport_lib::HttpTransportPool::get_block_height;
|
||||
let _get_block_time = ksp_onchain_transport_lib::HttpTransportPool::get_block_time;
|
||||
let _get_first_available_block = ksp_onchain_transport_lib::HttpTransportPool::get_first_available_block;
|
||||
let _minimum_ledger_slot = ksp_onchain_transport_lib::HttpTransportPool::minimum_ledger_slot;
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/tests/release_completeness.rs
|
||||
// version: 14
|
||||
// version: 15
|
||||
|
||||
//! Release-level completeness canaries for the staged HTTP wrapper sequence.
|
||||
|
||||
@@ -432,7 +432,7 @@ fn release_pre_008_ksp_transport_007_retro_audit_covers_all_typed_current_method
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn release_v0_2_4_descriptor_set_is_exact_and_remains_read_retry_safe_before_wrappers() {
|
||||
fn release_v0_2_4_descriptor_set_is_exact_and_remains_read_retry_safe_during_staged_wrappers() {
|
||||
let mut expected_blocks = std::vec![
|
||||
"getBlock",
|
||||
"getBlockCommitment",
|
||||
@@ -471,3 +471,28 @@ fn release_v0_2_4_descriptor_set_is_exact_and_remains_read_retry_safe_before_wra
|
||||
let get_block = ksp_onchain_transport_lib::find_http_rpc_method("getBlock").expect("getBlock descriptor must exist");
|
||||
assert!(get_block.request_form_status().has_deprecated_legacy());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn release_v0_2_4_pre_003_simple_blocks_subset_is_exact_and_retry_safe() {
|
||||
let mut expected = std::vec!["getBlockCommitment", "getBlockHeight", "getBlockTime", "getFirstAvailableBlock", "minimumLedgerSlot"];
|
||||
let mut actual = std::vec::Vec::new();
|
||||
for descriptor in ksp_onchain_transport_lib::current_http_rpc_methods() {
|
||||
if descriptor.coverage_release() != ksp_onchain_transport_lib::HttpRpcCoverageRelease::V0_2_4
|
||||
|| descriptor.category() != ksp_onchain_transport_lib::HttpRpcCategory::Blocks
|
||||
{
|
||||
continue;
|
||||
}
|
||||
match descriptor.method() {
|
||||
"getBlockCommitment" | "getBlockHeight" | "getBlockTime" | "getFirstAvailableBlock" | "minimumLedgerSlot" => {
|
||||
assert_eq!(descriptor.operation_kind(), ksp_onchain_transport_lib::RpcOperationKind::Read);
|
||||
assert_eq!(descriptor.transport_retry_class(), ksp_onchain_transport_lib::TransportRetryClass::RetrySafe);
|
||||
actual.push(descriptor.method());
|
||||
},
|
||||
_ => {},
|
||||
}
|
||||
}
|
||||
actual.sort_unstable();
|
||||
expected.sort_unstable();
|
||||
assert_eq!(actual, expected);
|
||||
assert_eq!(actual.len(), 5);
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-onchain-transport-lib/unit_tests/rpc_blocks.rs
|
||||
// version: 1
|
||||
// version: 2
|
||||
|
||||
#[test]
|
||||
fn transaction_details_and_get_block_config_preserve_all_modern_options() {
|
||||
@@ -148,3 +148,166 @@ fn block_reward_rejects_invalid_pubkey_without_echoing_value() {
|
||||
assert_eq!(error.code(), crate::ERROR_CODE_INVALID_RESPONSE);
|
||||
assert!(!error.to_string().contains("not-a-pubkey"));
|
||||
}
|
||||
|
||||
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<std::string::String>) {
|
||||
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::<usize>().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_block_commitment_serializes_slot_and_preserves_distribution() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_block_commitment.success.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let result = pool.get_block_commitment(&crate::HttpRoleName::new("default"), 430_000_123).await.expect("block commitment fixture must succeed");
|
||||
assert_eq!(result.commitment(), std::option::Option::Some(&[1, 2, 3, 4][..]));
|
||||
assert_eq!(result.total_stake(), 1_000_000_000);
|
||||
let request = handle.join().expect("fixture server must join");
|
||||
let body = request_body(request.as_str());
|
||||
assert_eq!(body["method"], serde_json::json!("getBlockCommitment"));
|
||||
assert_eq!(body["params"], serde_json::json!([430000123]));
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn typed_get_block_height_serializes_context_config() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_block_height.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 height = pool
|
||||
.get_block_height(&crate::HttpRoleName::new("default"), std::option::Option::Some(&config))
|
||||
.await
|
||||
.expect("block height fixture must succeed");
|
||||
assert_eq!(height, 410_000_000);
|
||||
let request = handle.join().expect("fixture server must join");
|
||||
assert_eq!(request_body(request.as_str())["params"], serde_json::json!([{"commitment":"finalized","minContextSlot":429000000}]));
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn typed_get_block_height_omits_explicitly_empty_config() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_block_height.success.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let config = crate::SolanaContextConfig::default();
|
||||
let height = pool
|
||||
.get_block_height(&crate::HttpRoleName::new("default"), std::option::Option::Some(&config))
|
||||
.await
|
||||
.expect("block height fixture must succeed");
|
||||
assert_eq!(height, 410_000_000);
|
||||
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_block_time_preserves_timestamp_and_null() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_block_time.success.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let timestamp = pool.get_block_time(&crate::HttpRoleName::new("default"), 430_000_123).await.expect("block time fixture must succeed");
|
||||
assert_eq!(timestamp, std::option::Option::Some(1_787_072_400));
|
||||
let request = handle.join().expect("fixture server must join");
|
||||
assert_eq!(request_body(request.as_str())["params"], serde_json::json!([430000123]));
|
||||
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_block_time.null.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let timestamp = pool.get_block_time(&crate::HttpRoleName::new("default"), 430_000_124).await.expect("null block time fixture must succeed");
|
||||
assert_eq!(timestamp, std::option::Option::None);
|
||||
handle.join().expect("fixture server must join");
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn typed_get_first_available_block_has_no_params() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/get_first_available_block.success.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let slot = pool.get_first_available_block(&crate::HttpRoleName::new("default")).await.expect("first available block fixture must succeed");
|
||||
assert_eq!(slot, 250_000);
|
||||
let request = handle.join().expect("fixture server must join");
|
||||
let body = request_body(request.as_str());
|
||||
assert_eq!(body["method"], serde_json::json!("getFirstAvailableBlock"));
|
||||
assert_eq!(body["params"], serde_json::json!([]));
|
||||
}
|
||||
|
||||
#[tokio::test(flavor = "current_thread")]
|
||||
async fn typed_minimum_ledger_slot_has_no_params() {
|
||||
let (url, handle) = serve_once(include_str!("../fixtures/http/minimum_ledger_slot.success.json"));
|
||||
let pool = pool_for_url(url.as_str());
|
||||
let slot = pool.minimum_ledger_slot(&crate::HttpRoleName::new("default")).await.expect("minimum ledger slot fixture must succeed");
|
||||
assert_eq!(slot, 123_456);
|
||||
let request = handle.join().expect("fixture server must join");
|
||||
let body = request_body(request.as_str());
|
||||
assert_eq!(body["method"], serde_json::json!("minimumLedgerSlot"));
|
||||
assert_eq!(body["params"], serde_json::json!([]));
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user