v0.2.11-pre.003-fix.001
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
// file: crates/ksp-offchain-transport-lib/src/http_admission.rs
|
||||
// version: 1
|
||||
// version: 2
|
||||
|
||||
//! Provider-neutral local request admission and rate-limit cooldown primitives.
|
||||
|
||||
@@ -20,14 +20,14 @@ pub(crate) enum HttpAdmissionPolicy {
|
||||
Unlimited,
|
||||
}
|
||||
|
||||
impl HttpAdmissionPolicy {
|
||||
impl crate::HttpAdmissionPolicy {
|
||||
/// Creates a validated fixed-window admission policy.
|
||||
pub(crate) fn fixed(requests: u32, window: std::time::Duration, burst: std::option::Option<u32>) -> ksp_core_lib::Result<Self> {
|
||||
let burst = match burst {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => 1,
|
||||
};
|
||||
if requests == 0 || window.is_zero() || window > MAX_RATE_LIMIT_WINDOW || burst == 0 || burst > requests {
|
||||
if requests == 0 || window.is_zero() || window > MAX_RATE_LIMIT_WINDOW || burst == 0 {
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_HTTP_RATE_LIMIT_INVALID, "HTTP request-admission policy is invalid")
|
||||
.with_context("field", "rate_limit"),
|
||||
@@ -49,28 +49,28 @@ pub(crate) enum HttpAdmissionDecision {
|
||||
pub(crate) struct HttpAdmissionController {
|
||||
cooldown_until: std::sync::Mutex<std::option::Option<std::time::Instant>>,
|
||||
fallback_cooldown: std::time::Duration,
|
||||
policy: HttpAdmissionPolicy,
|
||||
policy: crate::HttpAdmissionPolicy,
|
||||
token_bucket: std::sync::Mutex<std::option::Option<HttpTokenBucketState>>,
|
||||
}
|
||||
|
||||
impl HttpAdmissionController {
|
||||
impl crate::HttpAdmissionController {
|
||||
/// Creates one limiter from a provider-owned local admission policy.
|
||||
pub(crate) fn new(policy: HttpAdmissionPolicy, fallback_cooldown: std::option::Option<std::time::Duration>) -> ksp_core_lib::Result<Self> {
|
||||
pub(crate) fn new(policy: crate::HttpAdmissionPolicy, fallback_cooldown: std::option::Option<std::time::Duration>) -> ksp_core_lib::Result<Self> {
|
||||
let fallback_cooldown = match fallback_cooldown {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => DEFAULT_RATE_LIMIT_COOLDOWN,
|
||||
};
|
||||
if fallback_cooldown.is_zero() || fallback_cooldown > HTTP_MAX_RETRY_AFTER {
|
||||
if fallback_cooldown.is_zero() || fallback_cooldown > crate::HTTP_MAX_RETRY_AFTER {
|
||||
return std::result::Result::Err(
|
||||
ksp_core_lib::Error::new(crate::ERROR_CODE_HTTP_RATE_LIMIT_INVALID, "HTTP fallback cooldown is outside the supported bounds")
|
||||
.with_context("field", "fallback_cooldown"),
|
||||
);
|
||||
}
|
||||
let token_bucket = match policy {
|
||||
HttpAdmissionPolicy::Fixed { requests, window, burst } => {
|
||||
crate::HttpAdmissionPolicy::Fixed { requests, window, burst } => {
|
||||
std::option::Option::Some(HttpTokenBucketState::new(requests, window, burst, std::time::Instant::now()))
|
||||
},
|
||||
HttpAdmissionPolicy::Dynamic | HttpAdmissionPolicy::Unlimited => std::option::Option::None,
|
||||
crate::HttpAdmissionPolicy::Dynamic | crate::HttpAdmissionPolicy::Unlimited => std::option::Option::None,
|
||||
};
|
||||
return std::result::Result::Ok(Self {
|
||||
cooldown_until: std::sync::Mutex::new(std::option::Option::None),
|
||||
@@ -82,19 +82,19 @@ impl HttpAdmissionController {
|
||||
|
||||
/// Returns the configured local policy.
|
||||
#[must_use]
|
||||
pub(crate) const fn policy(&self) -> HttpAdmissionPolicy {
|
||||
pub(crate) const fn policy(&self) -> crate::HttpAdmissionPolicy {
|
||||
return self.policy;
|
||||
}
|
||||
|
||||
/// Tries to admit one request immediately without sleeping.
|
||||
pub(crate) fn try_admit(&self) -> HttpAdmissionDecision {
|
||||
pub(crate) fn try_admit(&self) -> crate::HttpAdmissionDecision {
|
||||
return self.try_admit_at(std::time::Instant::now());
|
||||
}
|
||||
|
||||
/// Records a provider 429 and extends cooldown using a bounded `Retry-After` value when present.
|
||||
pub(crate) fn record_rate_limited(&self, provider_retry_after: std::option::Option<std::time::Duration>) -> std::time::Duration {
|
||||
let provider_delay = match provider_retry_after {
|
||||
std::option::Option::Some(value) => std::cmp::min(value, HTTP_MAX_RETRY_AFTER),
|
||||
std::option::Option::Some(value) => std::cmp::min(value, crate::HTTP_MAX_RETRY_AFTER),
|
||||
std::option::Option::None => std::time::Duration::ZERO,
|
||||
};
|
||||
let effective = std::cmp::max(self.fallback_cooldown, provider_delay);
|
||||
@@ -113,9 +113,9 @@ impl HttpAdmissionController {
|
||||
return self.cooldown_remaining_at(std::time::Instant::now());
|
||||
}
|
||||
|
||||
fn try_admit_at(&self, now: std::time::Instant) -> HttpAdmissionDecision {
|
||||
fn try_admit_at(&self, now: std::time::Instant) -> crate::HttpAdmissionDecision {
|
||||
if let std::option::Option::Some(remaining) = self.cooldown_remaining_at(now) {
|
||||
return HttpAdmissionDecision::Deferred(remaining);
|
||||
return crate::HttpAdmissionDecision::Deferred(remaining);
|
||||
}
|
||||
let lock_result = self.token_bucket.lock();
|
||||
let mut token_bucket = match lock_result {
|
||||
@@ -124,11 +124,11 @@ impl HttpAdmissionController {
|
||||
};
|
||||
let state = match token_bucket.as_mut() {
|
||||
std::option::Option::Some(value) => value,
|
||||
std::option::Option::None => return HttpAdmissionDecision::Ready,
|
||||
std::option::Option::None => return crate::HttpAdmissionDecision::Ready,
|
||||
};
|
||||
return match state.try_consume_at(now) {
|
||||
std::option::Option::Some(delay) => HttpAdmissionDecision::Deferred(delay),
|
||||
std::option::Option::None => HttpAdmissionDecision::Ready,
|
||||
std::option::Option::Some(delay) => crate::HttpAdmissionDecision::Deferred(delay),
|
||||
std::option::Option::None => crate::HttpAdmissionDecision::Ready,
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user