v0.4.8-pre.010

This commit is contained in:
2026-08-06 17:26:14 +02:00
parent b314bfc244
commit 3167297283
40 changed files with 2556 additions and 1125 deletions

View File

@@ -1,8 +1,22 @@
<!-- file: kb-onchain-transport/CHANGELOG.md -->
<!-- version: 8 -->
<!-- version: 11 -->
# CHANGELOG — kb-onchain-transport
## 0.4.8-pre.010
- documente que les limites de rôles doivent être additionnées pour respecter le quota global dun endpoint public et que le quota par méthode peut être plus strict que le quota global ;
- aligne lexemple `local_devnet` sur un pacing conservateur pour `api.devnet.solana.com` : `3 r/s` pour les lectures et `1 r/s` pour les transactions, avec bursts bornés et cooldown de `10 s` ;
- applique réellement `requests_per_second`, `burst_capacity` et `max_concurrent_requests` au rôle HTTP sélectionné avant chaque tentative réseau ;
- partage le token bucket, le sémaphore de concurrence et le cooldown distant entre tous les clones dun même endpoint ;
- conserve le retry réactif `429` comme seconde ligne de défense lorsque le quota distant est plus strict ou partagé avec dautres clients ;
- remplace les trois opérateurs `?` accidentels du parseur `Retry-After` par des branches explicites conformes aux règles du workspace ;
- expose la borne publique `MAX_COMPLETE_ACCOUNT_DATA_BYTES = 65536` pour réconcilier les pipelines stateful avec le contrat réel de `getAccountInfo` ;
- ajoute un retry borné avec backoff exponentiel sur les statuts HTTP `429` et les erreurs JSON-RPC `429` ;
- respecte `Retry-After` lorsquil contient un délai en secondes et utilise sinon `pause_after_rate_limit_ms` des rôles configurés ;
- conserve le même payload et la même transaction signée lors des retries, ce qui maintient lidempotence de `sendTransaction` ;
## 0.4.6
- alignement de la crate sur la version fonctionnelle bot3 `0.4.6` ;

View File

@@ -1,5 +1,5 @@
<!-- file: kb-onchain-transport/USAGE.md -->
<!-- version: 2 -->
<!-- version: 5 -->
# Utilisation de kb-onchain-transport
@@ -126,6 +126,18 @@ Les types suivants couvrent les opérations dexécution :
La soumission ne remplace pas les contrôles de sécurité, préflights et confirmations opérateur réalisés par les couches supérieures.
## Lectures complètes et limites de débit
`MAX_COMPLETE_ACCOUNT_DATA_BYTES` expose la borne commune des données décodées retournées par une lecture complète `getAccountInfo`. Les couches supérieures doivent utiliser cette constante au lieu daligner une lecture RPC sur la capacité maximale dun décodeur offline.
Le client applique dabord les limites préventives du rôle sélectionné :
- `requests_per_second` alimente un token bucket partagé ;
- `burst_capacity` borne le burst initial et les crédits accumulés ;
- `max_concurrent_requests` borne les requêtes simultanément en vol.
Ces limiteurs sont partagés entre les clones dun même endpoint. Les statuts HTTP `429` et erreurs JSON-RPC `429` restent réessayés avec un backoff borné, car un fournisseur peut appliquer un quota global, dynamique ou partagé que la configuration locale ne peut pas connaître. Le délai configuré par `pause_after_rate_limit_ms` est alors appliqué comme cooldown partagé, sauf lorsquun en-tête `Retry-After` impose une attente supérieure. Le nombre de retries reste limité et une erreur terminale conserve le nombre de tentatives.
## Erreurs et invariants
- les réponses sont validées et adaptées avant exposition aux couches supérieures ;
@@ -145,3 +157,25 @@ La soumission ne remplace pas les contrôles de sécurité, préflights et confi
- la crate traite les transports on-chain ; elle ne gère pas les métadonnées HTTP, IPFS ou Arweave ;
- elle ne décode pas les instructions de programmes ;
- elle ne décide pas seule de lautorisation denvoyer une transaction.
## Limitation préventive et réponses `429`
Le client applique les limites du rôle sélectionné avant chaque tentative :
- `requests_per_second` et `burst_capacity` alimentent un token bucket partagé ;
- `max_concurrent_requests` borne les requêtes simultanées ;
- `pause_after_rate_limit_ms` fournit le cooldown de repli ;
- `Retry-After` est prioritaire lorsquil est fourni par lendpoint.
Ces limites sont propres aux rôles du client. Elles ne créent pas de quotas indépendants côté fournisseur. Lorsque plusieurs rôles utilisent le même endpoint, leur débit et leurs bursts doivent être additionnés pour rester sous la limite globale de cet endpoint.
Pour `api.devnet.solana.com`, le profil dexemple utilise volontairement :
```text
http_queries 3 r/s, burst 3, concurrence 2
http_transactions 1 r/s, burst 1, concurrence 1
cooldown 10000 ms
```
Un warning `retry_http_json_rpc_after_rate_limit` signifie que le serveur a tout de même renvoyé `429` et que la seconde ligne de défense a été activée. La requête nest perdue que si le retry borné se termine en erreur.

View File

@@ -1,5 +1,5 @@
// file: kb-onchain-transport/src/constants.rs
// version: 7
// version: 8
//! Local constants for the `kb-onchain-transport` crate.
@@ -13,8 +13,12 @@ pub(crate) const TESTNET_GENESIS_HASH: &str = "4uhcVJyU9pJkvQyS88uRDiswHXSCkY3zQ
pub(crate) const MAINNET_GENESIS_HASH: &str = "5eykt4UsFv8P8NJdTREpY1vzqKqZKvdpKuc147dw2N9d";
/// Local defensive maximum for one base64 message or transaction request value.
pub(crate) const MAX_EXECUTION_RPC_BASE64_LENGTH: usize = 65_536;
/// Local defensive maximum for decoded account data returned by one execution RPC request.
pub(crate) const MAX_EXECUTION_ACCOUNT_DATA_BYTES: usize = 65_536;
/// Maximum decoded account bytes returned by one complete execution RPC account read.
pub const MAX_COMPLETE_ACCOUNT_DATA_BYTES: usize = 65_536;
/// Maximum retries after one HTTP or JSON-RPC rate-limit response.
pub(crate) const MAX_HTTP_RATE_LIMIT_RETRIES: u32 = 4;
/// Fallback pause when a matching endpoint role exposes no rate-limit delay.
pub(crate) const DEFAULT_HTTP_RATE_LIMIT_PAUSE_MS: u64 = 1_500;
/// Local defensive maximum for account snapshots requested from one simulation.
pub(crate) const MAX_SIMULATION_ACCOUNT_COUNT: usize = 128;
/// Maximum signatures accepted by one `getSignatureStatuses` request.

View File

@@ -1,5 +1,5 @@
// file: kb-onchain-transport/src/execution_rpc.rs
// version: 8
// version: 9
//! Typed Solana JSON-RPC adapters used by execution orchestration.
@@ -537,10 +537,10 @@ impl crate::GetAccountInfoConfig {
min_context_slot: std::option::Option<u64>,
max_data_bytes: usize,
) -> kb_core::Result<Self> {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES {
return std::result::Result::Err(kb_core::Error::config(format!(
"getAccountInfo complete data limit must be between 1 and {} bytes",
crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES
crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES
)));
}
return std::result::Result::Ok(Self {
@@ -569,10 +569,10 @@ impl crate::GetAccountInfoConfig {
pubkey: &kb_lib::MdPubkey,
) -> kb_core::Result<std::vec::Vec<serde_json::Value>> {
if let std::option::Option::Some(max_data_bytes) = self.max_data_bytes {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES {
return std::result::Result::Err(kb_core::Error::config(format!(
"getAccountInfo complete data limit must be between 1 and {} bytes",
crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES
crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES
)));
}
}
@@ -1344,10 +1344,10 @@ fn decode_account_data(
}
},
std::option::Option::Some(max_data_bytes) => {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES {
if max_data_bytes == 0 || max_data_bytes > crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES {
return std::result::Result::Err(kb_core::Error::json(format!(
"getAccountInfo complete data limit must be between 1 and {} bytes",
crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES
crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES
)));
}
if space > max_data_bytes as u64 {
@@ -2534,7 +2534,7 @@ mod tests {
assert!(crate::GetAccountInfoConfig::confirmed_with_data(0).is_err());
assert!(
crate::GetAccountInfoConfig::confirmed_with_data(
crate::MAX_EXECUTION_ACCOUNT_DATA_BYTES + 1,
crate::MAX_COMPLETE_ACCOUNT_DATA_BYTES + 1,
)
.is_err()
);

View File

@@ -1,5 +1,5 @@
// file: kb-onchain-transport/src/http_client.rs
// version: 11
// version: 13
//! HTTP JSON-RPC client for standard Solana RPC endpoints.
@@ -30,12 +30,121 @@ pub struct HttpPoolClientSnapshot {
pub status: std::string::String,
}
#[derive(Debug)]
struct HttpRequestLimitState {
available_tokens: f64,
blocked_until: std::time::Instant,
last_refill: std::time::Instant,
}
#[derive(Debug)]
struct HttpRequestLimiter {
burst_capacity: u32,
pause_after_rate_limit_ms: u64,
priority: u32,
request_kinds: std::vec::Vec<std::string::String>,
requests_per_second: u32,
role: std::string::String,
semaphore: std::sync::Arc<tokio::sync::Semaphore>,
state: tokio::sync::Mutex<HttpRequestLimitState>,
}
impl HttpRequestLimiter {
fn new(config: &kb_config::EndpointRoleConfig) -> Self {
let now = std::time::Instant::now();
let burst_capacity = config.burst_capacity.max(1);
let max_concurrent_requests = config.max_concurrent_requests.max(1);
return Self {
burst_capacity,
pause_after_rate_limit_ms: config.pause_after_rate_limit_ms,
priority: config.priority,
request_kinds: config.request_kinds.clone(),
requests_per_second: config.requests_per_second.max(1),
role: config.role.clone(),
semaphore: std::sync::Arc::new(tokio::sync::Semaphore::new(
max_concurrent_requests as usize,
)),
state: tokio::sync::Mutex::new(HttpRequestLimitState {
available_tokens: f64::from(burst_capacity),
blocked_until: now,
last_refill: now,
}),
};
}
fn handles_request(&self, request_kind: &str) -> bool {
for configured_kind in &self.request_kinds {
if configured_kind == request_kind || configured_kind == "*" {
return true;
}
}
return false;
}
async fn acquire(&self) -> kb_core::Result<tokio::sync::OwnedSemaphorePermit> {
loop {
let wait_duration = {
let mut state = self.state.lock().await;
let now = std::time::Instant::now();
if now < state.blocked_until {
state.blocked_until.duration_since(now)
} else {
let elapsed_seconds = now.duration_since(state.last_refill).as_secs_f64();
let refill = elapsed_seconds * f64::from(self.requests_per_second);
state.available_tokens =
(state.available_tokens + refill).min(f64::from(self.burst_capacity));
state.last_refill = now;
if state.available_tokens >= 1.0 {
state.available_tokens -= 1.0;
std::time::Duration::ZERO
} else {
let missing_tokens = 1.0 - state.available_tokens;
std::time::Duration::from_secs_f64(
missing_tokens / f64::from(self.requests_per_second),
)
}
}
};
if wait_duration.is_zero() {
break;
}
tokio::time::sleep(wait_duration).await;
}
let permit_result = self.semaphore.clone().acquire_owned().await;
return match permit_result {
std::result::Result::Ok(permit) => std::result::Result::Ok(permit),
std::result::Result::Err(error) => {
std::result::Result::Err(kb_core::Error::invalid_state(format!(
"http request limiter '{}' is closed: {error}",
self.role
)))
},
};
}
async fn block_for(&self, pause_ms: u64) {
let mut state = self.state.lock().await;
let candidate =
std::time::Instant::now().checked_add(std::time::Duration::from_millis(pause_ms));
let blocked_until = match candidate {
std::option::Option::Some(value) => value,
std::option::Option::None => std::time::Instant::now(),
};
if blocked_until > state.blocked_until {
state.blocked_until = blocked_until;
}
return;
}
}
/// HTTP JSON-RPC client bound to one configured endpoint.
#[derive(Clone, Debug)]
pub struct HttpClient {
endpoint: kb_config::HttpEndpointConfig,
client: reqwest::Client,
next_request_id: std::sync::Arc<std::sync::atomic::AtomicU64>,
request_limiters: std::sync::Arc<std::vec::Vec<std::sync::Arc<HttpRequestLimiter>>>,
selected_role: std::option::Option<std::string::String>,
}
impl crate::HttpClient {
@@ -65,14 +174,29 @@ impl crate::HttpClient {
)));
},
};
tracing::debug!(target: crate::TRACING_TARGET, action = "create_http_client", endpoint_name = %endpoint.name, provider = %endpoint.provider, role_count = endpoint.roles.len(), "HTTP client created");
let mut request_limiters = std::vec::Vec::new();
for role in &endpoint.roles {
if role.enabled {
request_limiters.push(std::sync::Arc::new(HttpRequestLimiter::new(role)));
}
}
tracing::debug!(target: crate::TRACING_TARGET, action = "create_http_client", endpoint_name = %endpoint.name, provider = %endpoint.provider, role_count = endpoint.roles.len(), limiter_count = request_limiters.len(), "HTTP client created");
return std::result::Result::Ok(Self {
endpoint,
client,
next_request_id: std::sync::Arc::new(std::sync::atomic::AtomicU64::new(1)),
request_limiters: std::sync::Arc::new(request_limiters),
selected_role: std::option::Option::None,
});
}
/// Returns a clone bound to the exact role selected by the endpoint pool.
pub(crate) fn with_required_role(&self, required_role: &str) -> Self {
let mut selected = self.clone();
selected.selected_role = std::option::Option::Some(required_role.to_string());
return selected;
}
/// Returns the endpoint name.
pub fn endpoint_name(&self) -> &str {
return self.endpoint.name.as_str();
@@ -199,65 +323,195 @@ impl crate::HttpClient {
let request_id = self.next_request_id.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
let parameter_count = params.len();
let method_class = crate::HttpClient::classify_method(method.as_str());
tracing::debug!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, method_class = ?method_class, parameter_count, "send HTTP JSON-RPC request");
let request_kind = crate::request_kind_from_method(method.as_str());
let request = crate::JsonRpcRequest::new_with_u64_id(request_id, method.clone(), params);
let response_result =
self.client.post(self.endpoint.url.as_str()).json(&request).send().await;
let response = match response_result {
std::result::Result::Ok(response) => response,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, error = %error, "HTTP JSON-RPC transport failed");
let request_limiter = self.request_limiter(request_kind.as_str());
let base_pause_ms = self.configured_rate_limit_pause_ms(request_kind.as_str());
let mut retry_index = 0_u32;
loop {
let request_permit = match &request_limiter {
std::option::Option::Some(limiter) => match limiter.acquire().await {
std::result::Result::Ok(permit) => std::option::Option::Some(permit),
std::result::Result::Err(error) => return std::result::Result::Err(error),
},
std::option::Option::None => std::option::Option::None,
};
tracing::debug!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, method_class = ?method_class, parameter_count, retry_index, selected_role = ?self.selected_role.as_deref(), "send HTTP JSON-RPC request");
let response_result =
self.client.post(self.endpoint.url.as_str()).json(&request).send().await;
let response = match response_result {
std::result::Result::Ok(response) => response,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, retry_index, error = %error, "HTTP JSON-RPC transport failed");
return std::result::Result::Err(kb_core::Error::http(format!(
"http json-rpc request '{}' failed on endpoint '{}': {error}",
method, self.endpoint.name
)));
},
};
let status = response.status();
let retry_after_ms = crate::HttpClient::retry_after_millis(response.headers());
let text_result = response.text().await;
let text = match text_result {
std::result::Result::Ok(text) => text,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, retry_index, error = %error, "HTTP JSON-RPC response body read failed");
return std::result::Result::Err(kb_core::Error::http(format!(
"cannot read http json-rpc response '{}' from endpoint '{}': {error}",
method, self.endpoint.name
)));
},
};
drop(request_permit);
if status == reqwest::StatusCode::TOO_MANY_REQUESTS
&& retry_index < crate::MAX_HTTP_RATE_LIMIT_RETRIES
{
retry_index = retry_index.saturating_add(1);
let configured_pause_ms =
crate::HttpClient::rate_limit_backoff_ms(base_pause_ms, retry_index);
let pause_ms = match retry_after_ms {
std::option::Option::Some(value) => std::cmp::max(value, configured_pause_ms),
std::option::Option::None => configured_pause_ms,
};
tracing::warn!(target: crate::TRACING_TARGET, action = "retry_http_json_rpc_after_rate_limit", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, retry_index, pause_ms, response_byte_length = text.len(), "HTTP JSON-RPC endpoint rate limited the request; retrying with bounded backoff");
if let std::option::Option::Some(limiter) = &request_limiter {
limiter.block_for(pause_ms).await;
}
tokio::time::sleep(std::time::Duration::from_millis(pause_ms)).await;
continue;
}
if !status.is_success() {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, retry_index, response_byte_length = text.len(), "HTTP JSON-RPC endpoint returned non-success status");
return std::result::Result::Err(kb_core::Error::http(format!(
"http json-rpc request '{}' failed on endpoint '{}': {error}",
method, self.endpoint.name
"http json-rpc endpoint '{}' returned status {} after {} retries: {}",
self.endpoint.name, status, retry_index, text
)));
},
};
let status = response.status();
let text_result = response.text().await;
let text = match text_result {
std::result::Result::Ok(text) => text,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, error = %error, "HTTP JSON-RPC response body read failed");
return std::result::Result::Err(kb_core::Error::http(format!(
"cannot read http json-rpc response '{}' from endpoint '{}': {error}",
method, self.endpoint.name
)));
},
};
if !status.is_success() {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, response_byte_length = text.len(), "HTTP JSON-RPC endpoint returned non-success status");
return std::result::Result::Err(kb_core::Error::http(format!(
"http json-rpc endpoint '{}' returned status {}: {}",
self.endpoint.name, status, text
)));
}
let parsed = match crate::parse_json_rpc_text(&text) {
std::result::Result::Ok(parsed) => parsed,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, retry_index, response_byte_length = text.len(), error = %error, "HTTP JSON-RPC response parsing failed");
return std::result::Result::Err(error);
},
};
match parsed {
crate::JsonRpcResponse::Success(success) => {
tracing::debug!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, retry_index, response_byte_length = text.len(), outcome = "success", "HTTP JSON-RPC request completed");
return std::result::Result::Ok(success.result);
},
crate::JsonRpcResponse::Error(error_response)
if error_response.error.code == 429
&& retry_index < crate::MAX_HTTP_RATE_LIMIT_RETRIES =>
{
retry_index = retry_index.saturating_add(1);
let configured_pause_ms =
crate::HttpClient::rate_limit_backoff_ms(base_pause_ms, retry_index);
let pause_ms = match retry_after_ms {
std::option::Option::Some(value) => {
std::cmp::max(value, configured_pause_ms)
},
std::option::Option::None => configured_pause_ms,
};
tracing::warn!(target: crate::TRACING_TARGET, action = "retry_http_json_rpc_after_rate_limit", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, rpc_error_code = error_response.error.code, rpc_error_message = %error_response.error.message, retry_index, pause_ms, "HTTP JSON-RPC endpoint returned a rate-limit RPC error; retrying with bounded backoff");
if let std::option::Option::Some(limiter) = &request_limiter {
limiter.block_for(pause_ms).await;
}
tokio::time::sleep(std::time::Duration::from_millis(pause_ms)).await;
continue;
},
crate::JsonRpcResponse::Error(error_response) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, retry_index, rpc_error_code = error_response.error.code, rpc_error_message = %error_response.error.message, "HTTP JSON-RPC endpoint returned an RPC error");
return std::result::Result::Err(kb_core::Error::http(format!(
"json-rpc error {} from '{}' after {} retries: {}",
error_response.error.code,
self.endpoint.name,
retry_index,
error_response.error.message
)));
},
crate::JsonRpcResponse::Notification(_) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, retry_index, "HTTP JSON-RPC response was an unexpected notification");
return std::result::Result::Err(kb_core::Error::http(
"http json-rpc response cannot be a notification".to_string(),
));
},
}
}
let parsed = match crate::parse_json_rpc_text(&text) {
std::result::Result::Ok(parsed) => parsed,
std::result::Result::Err(error) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, response_byte_length = text.len(), error = %error, "HTTP JSON-RPC response parsing failed");
return std::result::Result::Err(error);
},
}
fn request_limiter(
&self,
request_kind: &str,
) -> std::option::Option<std::sync::Arc<HttpRequestLimiter>> {
if let std::option::Option::Some(selected_role) = &self.selected_role {
for limiter in self.request_limiters.iter() {
if limiter.role.as_str() == selected_role.as_str()
&& limiter.handles_request(request_kind)
{
return std::option::Option::Some(limiter.clone());
}
}
}
let mut selected: std::option::Option<std::sync::Arc<HttpRequestLimiter>> =
std::option::Option::None;
for limiter in self.request_limiters.iter() {
if !limiter.handles_request(request_kind) {
continue;
}
let replace = match &selected {
std::option::Option::Some(current) => limiter.priority < current.priority,
std::option::Option::None => true,
};
if replace {
selected = std::option::Option::Some(limiter.clone());
}
}
if selected.is_some() {
return selected;
}
for limiter in self.request_limiters.iter() {
let replace = match &selected {
std::option::Option::Some(current) => limiter.priority < current.priority,
std::option::Option::None => true,
};
if replace {
selected = std::option::Option::Some(limiter.clone());
}
}
return selected;
}
fn configured_rate_limit_pause_ms(&self, request_kind: &str) -> u64 {
let limiter = self.request_limiter(request_kind);
return match limiter {
std::option::Option::Some(value) => value.pause_after_rate_limit_ms,
std::option::Option::None => crate::DEFAULT_HTTP_RATE_LIMIT_PAUSE_MS,
};
return match parsed {
crate::JsonRpcResponse::Success(success) => {
tracing::debug!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, http_status = %status, response_byte_length = text.len(), outcome = "success", "HTTP JSON-RPC request completed");
std::result::Result::Ok(success.result)
},
crate::JsonRpcResponse::Error(error_response) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, rpc_error_code = error_response.error.code, rpc_error_message = %error_response.error.message, "HTTP JSON-RPC endpoint returned an RPC error");
std::result::Result::Err(kb_core::Error::http(format!(
"json-rpc error {} from '{}': {}",
error_response.error.code, self.endpoint.name, error_response.error.message
)))
},
crate::JsonRpcResponse::Notification(_) => {
tracing::error!(target: crate::TRACING_TARGET, action = "execute_http_json_rpc", endpoint_name = %self.endpoint.name, provider = %self.endpoint.provider, request_id, method = %method, "HTTP JSON-RPC response was an unexpected notification");
std::result::Result::Err(kb_core::Error::http(
"http json-rpc response cannot be a notification".to_string(),
))
},
}
fn rate_limit_backoff_ms(base_pause_ms: u64, retry_index: u32) -> u64 {
let shift = retry_index.saturating_sub(1).min(6);
let multiplier = match 1_u64.checked_shl(shift) {
std::option::Option::Some(value) => value,
std::option::Option::None => 64,
};
return base_pause_ms.saturating_mul(multiplier).min(60_000);
}
fn retry_after_millis(headers: &reqwest::header::HeaderMap) -> std::option::Option<u64> {
let value = match headers.get(reqwest::header::RETRY_AFTER) {
std::option::Option::Some(value) => value,
std::option::Option::None => return std::option::Option::None,
};
let text = match value.to_str() {
std::result::Result::Ok(text) => text,
std::result::Result::Err(_) => return std::option::Option::None,
};
let seconds = match text.parse::<u64>() {
std::result::Result::Ok(seconds) => seconds,
std::result::Result::Err(_) => return std::option::Option::None,
};
return std::option::Option::Some(seconds.saturating_mul(1_000).min(60_000));
}
}
@@ -380,4 +634,84 @@ mod tests {
crate::HttpMethodClass::GeneralRpc
);
}
#[test]
fn configured_rate_limit_pause_uses_matching_role_and_wildcard() {
let client = match crate::HttpClient::new(endpoint(true)) {
std::result::Result::Ok(client) => client,
std::result::Result::Err(error) => panic!("client creation failed: {error}"),
};
assert_eq!(client.configured_rate_limit_pause_ms("get_version"), 1500);
assert_eq!(client.configured_rate_limit_pause_ms("send_transaction"), 1500);
}
#[test]
fn rate_limit_backoff_is_exponential_and_bounded() {
assert_eq!(crate::HttpClient::rate_limit_backoff_ms(1500, 1), 1500);
assert_eq!(crate::HttpClient::rate_limit_backoff_ms(1500, 2), 3000);
assert_eq!(crate::HttpClient::rate_limit_backoff_ms(1500, 3), 6000);
assert_eq!(crate::HttpClient::rate_limit_backoff_ms(3000, 10), 60_000);
}
#[test]
fn selected_role_uses_its_exact_configured_limiter() {
let mut endpoint = endpoint(true);
endpoint.roles[0].priority = 10;
endpoint.roles[0].requests_per_second = 8;
endpoint.roles[0].burst_capacity = 16;
endpoint.roles[0].max_concurrent_requests = 8;
endpoint.roles[0].pause_after_rate_limit_ms = 1500;
endpoint.roles[1].priority = 20;
endpoint.roles[1].requests_per_second = 2;
endpoint.roles[1].burst_capacity = 2;
endpoint.roles[1].max_concurrent_requests = 2;
endpoint.roles[1].pause_after_rate_limit_ms = 3000;
let client = match crate::HttpClient::new(endpoint) {
std::result::Result::Ok(client) => client,
std::result::Result::Err(error) => panic!("client creation failed: {error}"),
};
let selected = client.with_required_role("http_heavy");
let limiter = match selected.request_limiter("get_block") {
std::option::Option::Some(limiter) => limiter,
std::option::Option::None => panic!("selected role limiter missing"),
};
assert_eq!(limiter.role.as_str(), "http_heavy");
assert_eq!(limiter.requests_per_second, 2);
assert_eq!(limiter.burst_capacity, 2);
assert_eq!(limiter.semaphore.available_permits(), 2);
assert_eq!(limiter.pause_after_rate_limit_ms, 3000);
}
#[test]
fn unbound_client_prefers_the_lowest_priority_matching_role() {
let mut endpoint = endpoint(true);
endpoint.roles[0].priority = 20;
endpoint.roles[2].priority = 30;
let client = match crate::HttpClient::new(endpoint) {
std::result::Result::Ok(client) => client,
std::result::Result::Err(error) => panic!("client creation failed: {error}"),
};
let limiter = match client.request_limiter("get_version") {
std::option::Option::Some(limiter) => limiter,
std::option::Option::None => panic!("matching limiter missing"),
};
assert_eq!(limiter.role.as_str(), "http_queries");
}
#[test]
fn retry_after_seconds_are_converted_to_bounded_milliseconds() {
let mut headers = reqwest::header::HeaderMap::new();
headers
.insert(reqwest::header::RETRY_AFTER, reqwest::header::HeaderValue::from_static("3"));
assert_eq!(
crate::HttpClient::retry_after_millis(&headers),
std::option::Option::Some(3000)
);
headers
.insert(reqwest::header::RETRY_AFTER, reqwest::header::HeaderValue::from_static("999"));
assert_eq!(
crate::HttpClient::retry_after_millis(&headers),
std::option::Option::Some(60_000)
);
}
}

View File

@@ -1,5 +1,5 @@
// file: kb-onchain-transport/src/http_pool.rs
// version: 7
// version: 8
//! HTTP endpoint pool and role-based routing.
@@ -79,8 +79,9 @@ impl crate::HttpEndpointPool {
let index = (start_index + offset) % client_count;
let client = self.clients[index].clone();
if client.can_handle(required_role, request_kind) {
tracing::debug!(target: crate::TRACING_TARGET, action = "select_http_endpoint", required_role, request_kind, endpoint_name = %client.endpoint_name(), provider = %client.provider(), "selected HTTP endpoint");
return std::result::Result::Ok(client);
let selected = client.with_required_role(required_role);
tracing::debug!(target: crate::TRACING_TARGET, action = "select_http_endpoint", required_role, request_kind, endpoint_name = %selected.endpoint_name(), provider = %selected.provider(), "selected HTTP endpoint");
return std::result::Result::Ok(selected);
}
offset += 1;
}

View File

@@ -1,5 +1,5 @@
// file: kb-onchain-transport/src/lib.rs
// version: 7
// version: 8
#![forbid(unsafe_code)]
#![deny(unreachable_pub)]
@@ -34,6 +34,8 @@ mod ws_session;
pub use self::client::RpcEndpoint;
/// Minimal Solana RPC client abstraction.
pub use self::client::SolanaRpcClient;
/// Maximum decoded account bytes accepted by one complete account read.
pub use self::constants::MAX_COMPLETE_ACCOUNT_DATA_BYTES;
/// Endpoint role snapshot shared by HTTP and WebSocket pools.
pub use self::endpoint_role::EndpointRoleSnapshot;
/// Converts a JSON-RPC method name into a stable request kind.
@@ -461,6 +463,8 @@ pub use self::ws_session::WsSubscriptionSnapshot;
/// Acknowledgement returned after a WebSocket unsubscribe request.
pub use self::ws_session::WsUnsubscribeAck;
/// Internal DEFAULT_HTTP_RATE_LIMIT_PAUSE_MS contract.
pub(crate) use self::constants::DEFAULT_HTTP_RATE_LIMIT_PAUSE_MS;
/// Internal DEVNET_GENESIS_HASH contract.
pub(crate) use self::constants::DEVNET_GENESIS_HASH;
/// Internal MAINNET_GENESIS_HASH contract.
@@ -471,10 +475,10 @@ pub(crate) use self::constants::MAX_BLOCK_RANGE;
pub(crate) use self::constants::MAX_CONFIRMATION_ATTEMPTS;
/// Internal MAX_CONFIRMATION_POLL_INTERVAL_MS contract.
pub(crate) use self::constants::MAX_CONFIRMATION_POLL_INTERVAL_MS;
/// Internal MAX_EXECUTION_ACCOUNT_DATA_BYTES contract.
pub(crate) use self::constants::MAX_EXECUTION_ACCOUNT_DATA_BYTES;
/// Internal MAX_EXECUTION_RPC_BASE64_LENGTH contract.
pub(crate) use self::constants::MAX_EXECUTION_RPC_BASE64_LENGTH;
/// Internal MAX_HTTP_RATE_LIMIT_RETRIES contract.
pub(crate) use self::constants::MAX_HTTP_RATE_LIMIT_RETRIES;
/// Internal MAX_MEMCMP_BASE58_LENGTH contract.
pub(crate) use self::constants::MAX_MEMCMP_BASE58_LENGTH;
/// Internal MAX_MEMCMP_BASE64_LENGTH contract.