v0.2.5-pre.009

This commit is contained in:
2026-08-20 09:24:10 +02:00
parent 1125ade4c1
commit e91bf36e9e
83 changed files with 2467 additions and 1801 deletions

View File

@@ -1,5 +1,5 @@
// file: crates/ksp-onchain-transport-lib/src/resilience.rs
// version: 1
// version: 2
const DEFAULT_RATE_LIMIT_COOLDOWN: std::time::Duration = std::time::Duration::from_secs(1);
const MAX_PROVIDER_RETRY_AFTER: std::time::Duration = std::time::Duration::from_secs(60);
@@ -71,61 +71,6 @@ impl HttpRetryDecision {
}
}
/// Evaluates the centralized bounded HTTP retry policy for one audited RPC method.
///
/// `completed_retries` counts retries already performed after the initial attempt. Provider `Retry-After` values are defensively bounded to sixty seconds
/// before they can extend the local exponential backoff. RPC application errors are never converted into transport retries.
#[must_use]
pub fn evaluate_transport_retry(
method: &crate::HttpRpcMethodDescriptor,
settings: &crate::HttpRetrySettings,
cause: crate::HttpRetryCause,
dispatch_state: crate::HttpDispatchState,
completed_retries: u32,
provider_retry_after: std::option::Option<std::time::Duration>,
) -> crate::HttpRetryDecision {
if completed_retries >= settings.max_retries() || !cause.is_retryable() {
return crate::HttpRetryDecision::Stop;
}
if method.transport_retry_class() == crate::TransportRetryClass::NotApplicable {
return crate::HttpRetryDecision::Stop;
}
if method.transport_retry_class() == crate::TransportRetryClass::NeverAfterDispatch && dispatch_state == crate::HttpDispatchState::DispatchedAmbiguous {
return crate::HttpRetryDecision::Stop;
}
let retry_number = completed_retries.saturating_add(1);
let mut delay = retry_backoff(settings, retry_number);
if cause == crate::HttpRetryCause::RateLimited
&& let std::option::Option::Some(provider_delay) = provider_retry_after
{
let bounded_provider_delay = std::cmp::min(provider_delay, MAX_PROVIDER_RETRY_AFTER);
if bounded_provider_delay > delay {
delay = bounded_provider_delay;
}
}
return crate::HttpRetryDecision::RetryAfter(delay);
}
pub(crate) fn retry_backoff(settings: &crate::HttpRetrySettings, retry_number: u32) -> std::time::Duration {
let mut delay = settings.initial_backoff();
if retry_number <= 1 {
return std::cmp::min(delay, settings.max_backoff());
}
let mut step = 1_u32;
while step < retry_number {
let doubled = match delay.checked_mul(2) {
std::option::Option::Some(value) => value,
std::option::Option::None => settings.max_backoff(),
};
delay = std::cmp::min(doubled, settings.max_backoff());
if delay >= settings.max_backoff() {
return settings.max_backoff();
}
step = step.saturating_add(1);
}
return delay;
}
pub(crate) struct HttpRoleRuntime {
limits: crate::HttpRoleLimits,
bucket: std::sync::Mutex<std::option::Option<HttpTokenBucketState>>,
@@ -216,21 +161,21 @@ impl HttpRoleRuntime {
return self.rate_limit_count.load(std::sync::atomic::Ordering::Relaxed);
}
pub(crate) fn try_acquire(self: &std::sync::Arc<Self>, now: std::time::Instant) -> RoleAdmissionAttempt {
pub(crate) fn try_acquire(self: &std::sync::Arc<Self>, now: std::time::Instant) -> crate::RoleAdmissionAttempt {
if let std::option::Option::Some(remaining) = self.cooldown_remaining_at(now) {
let ready_at = match now.checked_add(remaining) {
std::option::Option::Some(value) => value,
std::option::Option::None => now,
};
return RoleAdmissionAttempt::BlockedUntil(ready_at);
return crate::RoleAdmissionAttempt::BlockedUntil(ready_at);
}
let semaphore_permit = match &self.semaphore {
std::option::Option::Some(semaphore) => {
let permit_result = std::sync::Arc::clone(semaphore).try_acquire_owned();
match permit_result {
std::result::Result::Ok(permit) => std::option::Option::Some(permit),
std::result::Result::Err(tokio::sync::TryAcquireError::NoPermits) => return RoleAdmissionAttempt::ConcurrencySaturated,
std::result::Result::Err(tokio::sync::TryAcquireError::Closed) => return RoleAdmissionAttempt::Unavailable,
std::result::Result::Err(tokio::sync::TryAcquireError::NoPermits) => return crate::RoleAdmissionAttempt::ConcurrencySaturated,
std::result::Result::Err(tokio::sync::TryAcquireError::Closed) => return crate::RoleAdmissionAttempt::Unavailable,
}
},
std::option::Option::None => std::option::Option::None,
@@ -239,9 +184,9 @@ impl HttpRoleRuntime {
if let std::option::Option::Some(ready_at) = token_result {
drop(semaphore_permit);
self.notify.notify_one();
return RoleAdmissionAttempt::BlockedUntil(ready_at);
return crate::RoleAdmissionAttempt::BlockedUntil(ready_at);
}
return RoleAdmissionAttempt::Ready(HttpConcurrencyPermit { semaphore_permit, notify: std::sync::Arc::clone(&self.notify) });
return crate::RoleAdmissionAttempt::Ready(HttpConcurrencyPermit { semaphore_permit, notify: std::sync::Arc::clone(&self.notify) });
}
pub(crate) fn record_success(&self) {
@@ -385,6 +330,61 @@ impl HttpTokenBucketState {
}
}
/// Evaluates the centralized bounded HTTP retry policy for one audited RPC method.
///
/// `completed_retries` counts retries already performed after the initial attempt. Provider `Retry-After` values are defensively bounded to sixty seconds
/// before they can extend the local exponential backoff. RPC application errors are never converted into transport retries.
#[must_use]
pub fn evaluate_transport_retry(
method: &crate::HttpRpcMethodDescriptor,
settings: &crate::HttpRetrySettings,
cause: crate::HttpRetryCause,
dispatch_state: crate::HttpDispatchState,
completed_retries: u32,
provider_retry_after: std::option::Option<std::time::Duration>,
) -> crate::HttpRetryDecision {
if completed_retries >= settings.max_retries() || !cause.is_retryable() {
return crate::HttpRetryDecision::Stop;
}
if method.transport_retry_class() == crate::TransportRetryClass::NotApplicable {
return crate::HttpRetryDecision::Stop;
}
if method.transport_retry_class() == crate::TransportRetryClass::NeverAfterDispatch && dispatch_state == crate::HttpDispatchState::DispatchedAmbiguous {
return crate::HttpRetryDecision::Stop;
}
let retry_number = completed_retries.saturating_add(1);
let mut delay = retry_backoff(settings, retry_number);
if cause == crate::HttpRetryCause::RateLimited
&& let std::option::Option::Some(provider_delay) = provider_retry_after
{
let bounded_provider_delay = std::cmp::min(provider_delay, MAX_PROVIDER_RETRY_AFTER);
if bounded_provider_delay > delay {
delay = bounded_provider_delay;
}
}
return crate::HttpRetryDecision::RetryAfter(delay);
}
fn retry_backoff(settings: &crate::HttpRetrySettings, retry_number: u32) -> std::time::Duration {
let mut delay = settings.initial_backoff();
if retry_number <= 1 {
return std::cmp::min(delay, settings.max_backoff());
}
let mut step = 1_u32;
while step < retry_number {
let doubled = match delay.checked_mul(2) {
std::option::Option::Some(value) => value,
std::option::Option::None => settings.max_backoff(),
};
delay = std::cmp::min(doubled, settings.max_backoff());
if delay >= settings.max_backoff() {
return settings.max_backoff();
}
step = step.saturating_add(1);
}
return delay;
}
#[cfg(test)]
#[path = "../unit_tests/resilience.rs"]
mod tests;