v0.3.6-pre.004

This commit is contained in:
2026-09-01 12:57:25 +02:00
parent 20cf700f22
commit 478548b7f4
9 changed files with 425 additions and 46 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/src/http_executor.rs
// version: 4
// version: 7
const HTTP_BAD_GATEWAY: u16 = 502;
const HTTP_GATEWAY_TIMEOUT: u16 = 504;
@@ -8,6 +8,63 @@ const HTTP_REQUEST_TIMEOUT: u16 = 408;
const HTTP_SERVICE_UNAVAILABLE: u16 = 503;
const HTTP_TOO_MANY_REQUESTS: u16 = 429;
/// Typed value returned by an observed HTTP RPC path together with the safe identity of the endpoint that produced the successful response.
///
/// Endpoint URLs, headers and raw HTTP bodies are intentionally absent.
#[derive(Clone, PartialEq)]
pub struct HttpObservedValue<T> {
value: T,
endpoint_name: std::string::String,
provider: crate::HttpProviderName,
}
impl<T> HttpObservedValue<T> {
/// Returns the typed RPC value.
#[must_use]
pub const fn value(&self) -> &T {
return &self.value;
}
/// Returns the safe configured identity of the endpoint that produced the successful response.
#[must_use]
pub fn endpoint_name(&self) -> &str {
return self.endpoint_name.as_str();
}
/// Returns the safe provider descriptor attached to the successful endpoint.
#[must_use]
pub const fn provider(&self) -> &crate::HttpProviderName {
return &self.provider;
}
/// Consumes the observation and returns only the typed value.
#[must_use]
pub fn into_value(self) -> T {
return self.value;
}
/// Builds an observed value from a successful Transport attempt and its safe routing identity.
pub(crate) fn new(value: T, endpoint_name: std::string::String, provider: crate::HttpProviderName) -> Self {
return Self { value, endpoint_name, provider };
}
/// Consumes the observation into its typed value and safe routing identity for crate-internal typed decoding.
pub(crate) fn into_parts(self) -> (T, std::string::String, crate::HttpProviderName) {
return (self.value, self.endpoint_name, self.provider);
}
}
impl<T> std::fmt::Debug for HttpObservedValue<T> {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
return formatter
.debug_struct("HttpObservedValue")
.field("endpoint_name", &self.endpoint_name)
.field("provider", &self.provider)
.field("value", &"<available>")
.finish();
}
}
impl crate::HttpTransportPool {
/// Executes one audited standard Solana HTTP JSON-RPC method through KSP routing, admission and bounded retry policy.
///
@@ -19,6 +76,33 @@ impl crate::HttpTransportPool {
method: &crate::HttpRpcMethodDescriptor,
params: std::vec::Vec<serde_json::Value>,
) -> ksp_core_lib::Result<serde_json::Value> {
return self.execute_standard_rpc_with(role, method, params, |value, _permit| return value).await;
}
/// Executes one audited standard Solana HTTP JSON-RPC method and retains only safe routing identity for the successful attempt.
pub(crate) async fn execute_standard_rpc_observed(
&self,
role: &crate::HttpRoleName,
method: &crate::HttpRpcMethodDescriptor,
params: std::vec::Vec<serde_json::Value>,
) -> ksp_core_lib::Result<crate::HttpObservedValue<serde_json::Value>> {
return self
.execute_standard_rpc_with(role, method, params, |value, permit| {
return crate::HttpObservedValue::new(value, permit.selection().endpoint_name().to_owned(), permit.client().provider().clone());
})
.await;
}
async fn execute_standard_rpc_with<T, F>(
&self,
role: &crate::HttpRoleName,
method: &crate::HttpRpcMethodDescriptor,
params: std::vec::Vec<serde_json::Value>,
on_success: F,
) -> ksp_core_lib::Result<T>
where
F: std::ops::FnOnce(serde_json::Value, &crate::HttpRequestPermit) -> T,
{
let support = method.ensure_runtime_supported();
if let std::result::Result::Err(error) = support {
return std::result::Result::Err(error);
@@ -156,7 +240,12 @@ impl crate::HttpTransportPool {
http_status = status,
"completed Solana HTTP JSON-RPC request"
);
return parsed.into_result();
let value = parsed.into_result();
let value = match value {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return std::result::Result::Ok(on_success(value, &permit));
}
}
}
@@ -203,14 +292,11 @@ async fn wait_retry_delay(delay: std::time::Duration, deadline: std::time::Insta
return std::time::Instant::now() < deadline;
}
fn execution_timeout(method: &crate::HttpRpcMethodDescriptor, message: &str) -> ksp_core_lib::Result<serde_json::Value> {
fn execution_timeout<T>(method: &crate::HttpRpcMethodDescriptor, message: &str) -> ksp_core_lib::Result<T> {
return std::result::Result::Err(ksp_core_lib::Error::new(crate::ERROR_CODE_TIMEOUT, message).with_context("rpc_method", method.method()));
}
fn rate_limited_error(
method: &crate::HttpRpcMethodDescriptor,
provider_retry_after: std::option::Option<std::time::Duration>,
) -> ksp_core_lib::Result<serde_json::Value> {
fn rate_limited_error<T>(method: &crate::HttpRpcMethodDescriptor, provider_retry_after: std::option::Option<std::time::Duration>) -> ksp_core_lib::Result<T> {
let mut error = ksp_core_lib::Error::new(crate::ERROR_CODE_RATE_LIMITED, "Solana HTTP endpoint rate-limited the JSON-RPC request")
.with_context("rpc_method", method.method());
if let std::option::Option::Some(delay) = provider_retry_after {
@@ -219,7 +305,7 @@ fn rate_limited_error(
return std::result::Result::Err(error);
}
fn http_status_error(method: &crate::HttpRpcMethodDescriptor, status: u16) -> ksp_core_lib::Result<serde_json::Value> {
fn http_status_error<T>(method: &crate::HttpRpcMethodDescriptor, status: u16) -> ksp_core_lib::Result<T> {
return std::result::Result::Err(
ksp_core_lib::Error::new(crate::ERROR_CODE_HTTP_REQUEST_FAILED, "Solana HTTP endpoint returned an unsuccessful status")
.with_context("rpc_method", method.method())

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/src/lib.rs
// version: 44
// version: 45
#![warn(missing_docs)]
#![deny(unreachable_pub)]
@@ -271,6 +271,8 @@ pub use self::http_client::HttpEndpointClient;
pub use self::http_client::HttpEndpointRoleSnapshot;
/// Safe metadata snapshot for one logical HTTP endpoint.
pub use self::http_client::HttpEndpointSnapshot;
/// Typed RPC value paired with the safe identity of the HTTP endpoint that produced the successful response.
pub use self::http_executor::HttpObservedValue;
/// Result of one logical endpoint selection.
pub use self::http_pool::HttpEndpointSelection;
/// Runtime admission permit for one HTTP request.

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/src/rpc_transactions.rs
// version: 10
// version: 11
const MAX_RECENT_PRIORITIZATION_FEE_ACCOUNTS: usize = 128;
const MAX_SIGNATURES_FOR_ADDRESS_LIMIT: usize = 1_000;
@@ -1249,25 +1249,42 @@ impl crate::HttpTransportPool {
signature: &str,
config: std::option::Option<&crate::SolanaGetTransactionConfig>,
) -> ksp_core_lib::Result<std::option::Option<crate::SolanaConfirmedTransaction>> {
if let std::option::Option::Some(config) = config
&& config.commitment() == std::option::Option::Some(crate::SolanaCommitment::Processed)
{
return std::result::Result::Err(
ksp_core_lib::Error::new(
crate::ERROR_CODE_INVALID_RPC_PARAMETERS,
"getTransaction commitment must be confirmed or finalized when explicitly provided",
)
.with_context("rpc_method", "getTransaction")
.with_context("commitment", "processed"),
);
}
let mut params = std::vec![serde_json::Value::String(signature.to_owned())];
if let std::option::Option::Some(config) = config {
params.push((*config).to_json_value());
}
let params = get_transaction_params(signature, config);
let params = match params {
std::result::Result::Ok(params) => params,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
return self.execute_get_transaction(role, params).await;
}
/// Executes the current object-form `getTransaction` request and reports the safe identity of the endpoint that produced the successful response.
///
/// Routing, admission, timeout and retry behavior are identical to [`Self::get_transaction`]. The returned observation never contains an endpoint URL,
/// HTTP headers or a raw HTTP body.
pub async fn get_transaction_observed(
&self,
role: &crate::HttpRoleName,
signature: &str,
config: std::option::Option<&crate::SolanaGetTransactionConfig>,
) -> ksp_core_lib::Result<crate::HttpObservedValue<std::option::Option<crate::SolanaConfirmedTransaction>>> {
let params = get_transaction_params(signature, config);
let params = match params {
std::result::Result::Ok(params) => params,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let observed = self.execute_transaction_rpc_observed("getTransaction", role, params).await;
let observed = match observed {
std::result::Result::Ok(observed) => observed,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
let (value, endpoint_name, provider) = observed.into_parts();
let transaction = decode_get_transaction(value);
return match transaction {
std::result::Result::Ok(transaction) => std::result::Result::Ok(crate::HttpObservedValue::new(transaction, endpoint_name, provider)),
std::result::Result::Err(error) => std::result::Result::Err(error),
};
}
/// Executes the deprecated bare-encoding `getTransaction` request form retained by Solana RPC for backwards compatibility.
#[deprecated(note = "use HttpTransportPool::get_transaction with SolanaGetTransactionConfig; the bare encoding request form is deprecated")]
pub async fn get_transaction_legacy(
@@ -1297,14 +1314,7 @@ impl crate::HttpTransportPool {
std::result::Result::Ok(value) => value,
std::result::Result::Err(error) => return std::result::Result::Err(error),
};
if value.is_null() {
return std::result::Result::Ok(std::option::Option::None);
}
let transaction = crate::SolanaConfirmedTransaction::decode_wire("getTransaction", value);
return match transaction {
std::result::Result::Ok(transaction) => std::result::Result::Ok(std::option::Option::Some(transaction)),
std::result::Result::Err(error) => std::result::Result::Err(error),
};
return decode_get_transaction(value);
}
/// Executes typed `requestAirdrop` through the common KSP HTTP transport path.
@@ -1466,6 +1476,54 @@ impl crate::HttpTransportPool {
};
return self.execute_standard_rpc(role, method, params).await;
}
async fn execute_transaction_rpc_observed(
&self,
method_name: &'static str,
role: &crate::HttpRoleName,
params: std::vec::Vec<serde_json::Value>,
) -> ksp_core_lib::Result<crate::HttpObservedValue<serde_json::Value>> {
let method = transaction_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_observed(role, method, params).await;
}
}
fn get_transaction_params(
signature: &str,
config: std::option::Option<&crate::SolanaGetTransactionConfig>,
) -> ksp_core_lib::Result<std::vec::Vec<serde_json::Value>> {
if let std::option::Option::Some(config) = config
&& config.commitment() == std::option::Option::Some(crate::SolanaCommitment::Processed)
{
return std::result::Result::Err(
ksp_core_lib::Error::new(
crate::ERROR_CODE_INVALID_RPC_PARAMETERS,
"getTransaction commitment must be confirmed or finalized when explicitly provided",
)
.with_context("rpc_method", "getTransaction")
.with_context("commitment", "processed"),
);
}
let mut params = std::vec![serde_json::Value::String(signature.to_owned())];
if let std::option::Option::Some(config) = config {
params.push((*config).to_json_value());
}
return std::result::Result::Ok(params);
}
fn decode_get_transaction(value: serde_json::Value) -> ksp_core_lib::Result<std::option::Option<crate::SolanaConfirmedTransaction>> {
if value.is_null() {
return std::result::Result::Ok(std::option::Option::None);
}
let transaction = crate::SolanaConfirmedTransaction::decode_wire("getTransaction", value);
return match transaction {
std::result::Result::Ok(transaction) => std::result::Result::Ok(std::option::Option::Some(transaction)),
std::result::Result::Err(error) => std::result::Result::Err(error),
};
}
#[derive(serde::Deserialize)]

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/tests/public_api.rs
// version: 48
// version: 49
//! Integration tests for the public `ksp-onchain-transport-lib` consumer contract.
@@ -199,6 +199,15 @@ fn public_pre_002_shared_rpc_types_are_constructible_from_crate_root() {
assert_eq!(vote.vote_pubkey(), std::option::Option::Some(&pubkey));
}
#[test]
fn public_v0_3_6_pre_004_observed_get_transaction_surface_is_available_from_crate_root() {
let _get_transaction_observed = ksp_onchain_transport_lib::HttpTransportPool::get_transaction_observed;
let observed: std::option::Option<
ksp_onchain_transport_lib::HttpObservedValue<std::option::Option<ksp_onchain_transport_lib::SolanaConfirmedTransaction>>,
> = std::option::Option::None;
assert!(observed.is_none());
}
#[test]
fn public_pre_003_account_wrappers_are_available_from_crate_root() {
let _get_account_info = ksp_onchain_transport_lib::HttpTransportPool::get_account_info;

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/unit_tests/rpc_transactions.rs
// version: 10
// version: 11
#[test]
fn transaction_encoding_strings_match_current_and_legacy_wire_labels() {
@@ -966,6 +966,80 @@ async fn typed_get_transaction_preserves_raw_json_meta_version_and_transaction_i
);
}
#[tokio::test(flavor = "current_thread")]
async fn typed_get_transaction_observed_reports_actual_winner_after_retry_reroute() {
let (first_url, first_handle) = serve_transaction_status_and_count("429 Too Many Requests");
let (winner_url, winner_handle) = serve_transaction_once(include_str!("../fixtures/http/get_transaction.base64.json"));
let urls = [(first_url.as_str(), "first-endpoint", "first-provider"), (winner_url.as_str(), "winner-endpoint", "winner-provider")];
let mut endpoints = std::vec::Vec::with_capacity(urls.len());
for (url, endpoint_name, provider) in urls {
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),
);
endpoints.push(crate::HttpEndpointSettings::new(
endpoint_name,
true,
crate::HttpProviderName::new(provider),
crate::HttpClusterName::new("local"),
crate::HttpEndpointUrl::parse(url).expect("fixture URL must parse"),
std::time::Duration::from_millis(100),
std::time::Duration::from_secs(1),
std::option::Option::Some(1),
std::vec![role],
));
}
let pool = crate::HttpTransportPool::new(crate::HttpTransportSettings::new(
endpoints,
crate::HttpRetrySettings::new(1, std::time::Duration::from_millis(1), std::time::Duration::from_millis(2)),
))
.expect("observed fixture pool must build");
let config = crate::SolanaGetTransactionConfig::new(
std::option::Option::Some(crate::SolanaCommitment::Confirmed),
std::option::Option::Some(crate::SolanaTransactionEncoding::Base64),
std::option::Option::Some(0),
);
let observed = pool
.get_transaction_observed(&crate::HttpRoleName::new("default"), "fixture-signature", std::option::Option::Some(&config))
.await
.expect("retry-safe observed getTransaction must succeed on the second endpoint");
assert_eq!(observed.endpoint_name(), "winner-endpoint");
assert_eq!(observed.provider().as_str(), "winner-provider");
let transaction = observed.value().as_ref().expect("winning response must contain a transaction");
assert_eq!(transaction.slot(), 431_000_062);
assert!(matches!(
transaction.transaction(),
crate::SolanaEncodedTransaction::Binary { encoding: crate::SolanaTransactionBinaryEncoding::Base64, .. }
));
let (first_count, first_request) = first_handle.join().expect("first fixture server must join");
assert_eq!(first_count, 1);
assert_eq!(transaction_request_body(first_request.as_str())["method"], serde_json::json!("getTransaction"));
let winner_request = winner_handle.join().expect("winner fixture server must join");
assert_eq!(transaction_request_body(winner_request.as_str())["method"], serde_json::json!("getTransaction"));
}
#[tokio::test(flavor = "current_thread")]
async fn typed_get_transaction_observed_preserves_null_and_redacts_typed_value_debug() {
let (url, handle) = serve_transaction_once(include_str!("../fixtures/http/get_transaction.null.json"));
let pool = transaction_pool_for_url(url.as_str());
let observed = pool
.get_transaction_observed(&crate::HttpRoleName::new("default"), "fixture-signature", std::option::Option::None)
.await
.expect("observed null getTransaction must succeed");
assert!(observed.value().is_none());
assert_eq!(observed.endpoint_name(), "fixture-0");
assert_eq!(observed.provider().as_str(), "fixture");
let rendered = format!("{observed:?}");
assert!(rendered.contains("fixture-0"));
assert!(rendered.contains("fixture"));
assert!(rendered.contains("<available>"));
assert!(!rendered.contains("fixture-signature"));
handle.join().expect("fixture server must join");
}
#[tokio::test(flavor = "current_thread")]
async fn typed_get_transaction_preserves_unsupported_version_rpc_error() {
let (url, handle) = serve_transaction_once(include_str!("../fixtures/http/get_transaction.error_unsupported_version.json"));