Merge branch 'wt/review-002-cli01-retry-after-budget'
This commit is contained in:
@@ -69,7 +69,7 @@ use reqwest_retry::RetryTransientMiddleware;
|
||||
use reqwest_retry::{Jitter, RetryDecision, RetryPolicy};
|
||||
use thiserror::Error;
|
||||
|
||||
use super::retry_after::RetryAfterMiddleware;
|
||||
use super::retry_after::{anchor_budget, RetryAfterMiddleware};
|
||||
|
||||
/// Maximum number of URLs tracked by the inlined `Retry-After`
|
||||
/// middleware (LRU-bounded; see `retry_after.rs`).
|
||||
@@ -388,6 +388,7 @@ fn is_idempotent(method: &reqwest::Method) -> bool {
|
||||
|
||||
struct RetryGateMiddleware {
|
||||
retry: Arc<RetryTransientMiddleware<TotalRetryBudget<ExponentialBackoff>>>,
|
||||
max_total_retry_duration: Duration,
|
||||
}
|
||||
|
||||
impl RetryGateMiddleware {
|
||||
@@ -404,6 +405,7 @@ impl RetryGateMiddleware {
|
||||
inner: policy,
|
||||
},
|
||||
)),
|
||||
max_total_retry_duration,
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -417,6 +419,7 @@ impl reqwest_middleware::Middleware for RetryGateMiddleware {
|
||||
next: reqwest_middleware::Next<'_>,
|
||||
) -> reqwest_middleware::Result<reqwest::Response> {
|
||||
if is_idempotent(req.method()) {
|
||||
anchor_budget(extensions, self.max_total_retry_duration);
|
||||
self.retry.handle(req, extensions, next).await
|
||||
} else {
|
||||
next.run(req, extensions).await
|
||||
@@ -594,15 +597,17 @@ fn build_client_with_pems(
|
||||
builder = builder.identity(identity);
|
||||
}
|
||||
let reqwest_client = builder.build().map_err(HttpClientBuildError::Build)?;
|
||||
let retry_after = RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
DEFAULT_RETRY_AFTER_CAPACITY,
|
||||
config.retry_after_ceiling,
|
||||
config.max_total_retry_duration,
|
||||
);
|
||||
let client = reqwest_middleware::ClientBuilder::new(reqwest_client)
|
||||
.with(RetryGateMiddleware::new(
|
||||
&config.retry,
|
||||
config.max_total_retry_duration,
|
||||
))
|
||||
.with(RetryAfterMiddleware::with_capacity_and_ceiling(
|
||||
DEFAULT_RETRY_AFTER_CAPACITY,
|
||||
config.retry_after_ceiling,
|
||||
))
|
||||
.with(retry_after)
|
||||
.build();
|
||||
Ok(client)
|
||||
}
|
||||
|
||||
+3
-1
@@ -6,5 +6,7 @@
|
||||
mod http_client;
|
||||
mod retry_after;
|
||||
|
||||
pub use http_client::{ClientCertConfig, HttpClientBuildError, HttpClientConfig, SharedHttpClient};
|
||||
pub use http_client::{
|
||||
ClientCertConfig, HttpClientBuildError, HttpClientConfig, RetryConfig, SharedHttpClient,
|
||||
};
|
||||
pub use retry_after::RetryAfterMiddleware;
|
||||
|
||||
+288
-27
@@ -5,13 +5,19 @@
|
||||
//! unbounded `HashMap<Url, SystemTime>` storage can be bounded for a
|
||||
//! long-running process.
|
||||
//!
|
||||
//! Policy (review-001 FWD-05 / FWD-11):
|
||||
//! Policy (review-001 FWD-05 / FWD-11, review-002 CLI-01):
|
||||
//!
|
||||
//! - Every parsed deadline is clamped to a ceiling
|
||||
//! ([`RetryAfterMiddleware::with_capacity_and_ceiling`]; the shared
|
||||
//! client wires `HttpClientConfig::retry_after_ceiling`, default
|
||||
//! 300 s) — a hostile upstream cannot park the client on a 10-year
|
||||
//! deadline.
|
||||
//! - A logical request's sleeps are further capped by the retry budget
|
||||
//! ([`BudgetClock`], anchored once per request by the retry gate and
|
||||
//! shared by every attempt through the request's `Extensions`): the
|
||||
//! sleep is truncated to the budget's `hard_stop`, and a throttle
|
||||
//! refresh cannot push the deadline past that stop (CLI-01 — a retry
|
||||
//! storm can no longer extend the wall clock ceiling-per-attempt).
|
||||
//! - Deadlines are recorded against the *effective* (post-redirect)
|
||||
//! URL — the host actually being throttled — not the pre-redirect
|
||||
//! request URL.
|
||||
@@ -24,7 +30,7 @@
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Mutex;
|
||||
use std::time::{Duration, SystemTime};
|
||||
use std::time::{Duration, Instant, SystemTime};
|
||||
|
||||
use http::Extensions;
|
||||
use reqwest::{Request, Response, StatusCode};
|
||||
@@ -47,6 +53,75 @@ fn clamp_deadline_to_ceiling(deadline: SystemTime, ceiling: Duration) -> Option<
|
||||
Some(deadline.min(capped))
|
||||
}
|
||||
|
||||
/// Wall-clock budget for one logical request (review-002 CLI-01).
|
||||
///
|
||||
/// Anchored once per request — by the outer retry gate, into the
|
||||
/// request's `Extensions` — and shared by every retry attempt of that
|
||||
/// request. Attempts made without the gate in the stack (test stacks,
|
||||
/// standalone `RetryAfterMiddleware` users) self-anchor on first use.
|
||||
///
|
||||
/// `hard_stop` anchors a coarse `Instant` monotonic clock; the
|
||||
/// `SystemTime` projection survives a wall-clock step during the
|
||||
/// request's lifetime.
|
||||
#[derive(Clone, Debug)]
|
||||
pub(crate) struct BudgetClock {
|
||||
hard_stop: SystemTime,
|
||||
anchored: Instant,
|
||||
}
|
||||
|
||||
impl BudgetClock {
|
||||
pub(crate) fn start(max_total_retry_duration: Duration) -> Self {
|
||||
Self {
|
||||
hard_stop: SystemTime::now()
|
||||
.checked_add(max_total_retry_duration)
|
||||
.unwrap_or(SystemTime::now()),
|
||||
anchored: Instant::now(),
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn for_duration(anchor: SystemTime, max_total_retry_duration: Duration) -> Self {
|
||||
Self {
|
||||
hard_stop: anchor
|
||||
.checked_add(max_total_retry_duration)
|
||||
.unwrap_or(anchor),
|
||||
anchored: Instant::now(),
|
||||
}
|
||||
}
|
||||
|
||||
fn hard_stop(&self) -> SystemTime {
|
||||
let offset = self
|
||||
.hard_stop
|
||||
.duration_since(self.anchor_wall())
|
||||
.unwrap_or(Duration::ZERO);
|
||||
self.anchor_wall()
|
||||
.checked_add(offset)
|
||||
.unwrap_or(self.hard_stop)
|
||||
}
|
||||
|
||||
fn anchor_wall(&self) -> SystemTime {
|
||||
let now = SystemTime::now();
|
||||
now.checked_sub(self.anchored.elapsed()).unwrap_or(now)
|
||||
}
|
||||
}
|
||||
|
||||
fn budget_from_extensions(extensions: &Extensions) -> Option<BudgetClock> {
|
||||
if let Some(clock) = extensions.get::<BudgetClock>() {
|
||||
return Some(clock.clone());
|
||||
}
|
||||
extensions
|
||||
.get::<RequestStartTime>()
|
||||
.map(|start| BudgetClock::for_duration(start.0, Duration::ZERO))
|
||||
}
|
||||
|
||||
pub(crate) fn anchor_budget(extensions: &mut Extensions, budget: Duration) {
|
||||
if extensions.get::<BudgetClock>().is_none() {
|
||||
extensions.insert(BudgetClock::start(budget));
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
struct RequestStartTime(SystemTime);
|
||||
|
||||
fn parse_retry_after_with_ceiling(value: &str, ceiling: Duration) -> Option<SystemTime> {
|
||||
let trimmed = value.trim();
|
||||
let parsed = if let Ok(secs) = trimmed.parse::<u64>() {
|
||||
@@ -69,10 +144,12 @@ pub struct RetryAfterMiddleware {
|
||||
deadlines: Mutex<HashMap<Url, SystemTime>>,
|
||||
capacity: usize,
|
||||
ceiling: Duration,
|
||||
budget: Duration,
|
||||
}
|
||||
|
||||
impl RetryAfterMiddleware {
|
||||
/// A middleware with the default 300 s `Retry-After` ceiling.
|
||||
/// A middleware with the default 300 s `Retry-After` ceiling and no
|
||||
/// retry-budget coupling.
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
@@ -89,19 +166,46 @@ impl RetryAfterMiddleware {
|
||||
///
|
||||
/// Panics when `capacity` is 0.
|
||||
pub fn with_capacity_and_ceiling(capacity: usize, ceiling: Duration) -> Self {
|
||||
Self::with_capacity_ceiling_and_budget(capacity, ceiling, Duration::ZERO)
|
||||
}
|
||||
|
||||
/// A middleware whose per-request sleeps and throttle refreshes are
|
||||
/// additionally capped by `budget` (the shared client wires
|
||||
/// `HttpClientConfig::max_total_retry_duration`).
|
||||
///
|
||||
/// # Panics
|
||||
///
|
||||
/// Panics when `capacity` is 0.
|
||||
pub fn with_capacity_ceiling_and_budget(
|
||||
capacity: usize,
|
||||
ceiling: Duration,
|
||||
budget: Duration,
|
||||
) -> Self {
|
||||
Self {
|
||||
deadlines: Mutex::new(HashMap::with_capacity(capacity.min(128))),
|
||||
capacity,
|
||||
ceiling,
|
||||
budget,
|
||||
}
|
||||
}
|
||||
|
||||
fn record(&self, url: Url, deadline: SystemTime) {
|
||||
fn record(&self, url: Url, deadline: Option<SystemTime>) {
|
||||
let mut deadlines = self.deadlines.lock().unwrap_or_else(|e| e.into_inner());
|
||||
if !deadlines.contains_key(&url) && deadlines.len() >= self.capacity {
|
||||
self.evict(&mut deadlines);
|
||||
}
|
||||
deadlines.insert(url, deadline);
|
||||
match deadline {
|
||||
Some(deadline) => {
|
||||
let earliest = deadlines
|
||||
.get(&url)
|
||||
.map(|existing| (*existing).min(deadline))
|
||||
.unwrap_or(deadline);
|
||||
deadlines.insert(url, earliest);
|
||||
}
|
||||
None => {
|
||||
deadlines.remove(&url);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn evict(&self, deadlines: &mut HashMap<Url, SystemTime>) {
|
||||
@@ -129,24 +233,38 @@ impl RetryAfterMiddleware {
|
||||
deadlines.get(url).copied()
|
||||
}
|
||||
|
||||
async fn maybe_sleep_for(&self, url: &Url) {
|
||||
async fn maybe_sleep_for(&self, url: &Url, extensions: &Extensions) {
|
||||
let Some(deadline) = self.deadline_for(url) else {
|
||||
return;
|
||||
};
|
||||
let Ok(remaining) = deadline.duration_since(SystemTime::now()) else {
|
||||
let now = SystemTime::now();
|
||||
let Ok(remaining) = deadline.duration_since(now) else {
|
||||
return;
|
||||
};
|
||||
if remaining.is_zero() {
|
||||
return;
|
||||
}
|
||||
let hard_stop = match budget_from_extensions(extensions) {
|
||||
Some(clock) => clock.hard_stop(),
|
||||
None if self.budget.is_zero() => {
|
||||
let jitter = max_sleep_jitter(remaining);
|
||||
let wait = remaining
|
||||
.checked_sub(jitter)
|
||||
.unwrap_or(Duration::from_millis(1));
|
||||
tokio::time::sleep(wait).await;
|
||||
return;
|
||||
}
|
||||
None => SystemTime::now() + self.budget,
|
||||
};
|
||||
let budget_left = hard_stop.duration_since(now).unwrap_or(Duration::ZERO);
|
||||
let remaining = remaining.min(budget_left);
|
||||
if remaining.is_zero() {
|
||||
return;
|
||||
}
|
||||
let wait = remaining
|
||||
.checked_sub(max_sleep_jitter(remaining))
|
||||
.unwrap_or(Duration::from_millis(1));
|
||||
tokio::time::sleep(wait).await;
|
||||
}
|
||||
|
||||
fn record_if_throttled(&self, url: Url, response: &Response) {
|
||||
fn record_if_throttled(&self, url: Url, response: &Response, extensions: &Extensions) {
|
||||
let status = response.status();
|
||||
if is_throttled(status.as_u16()) {
|
||||
if let Some(retry_after) = response
|
||||
@@ -155,7 +273,18 @@ impl RetryAfterMiddleware {
|
||||
.and_then(|value| value.to_str().ok())
|
||||
{
|
||||
if let Some(deadline) = parse_retry_after_with_ceiling(retry_after, self.ceiling) {
|
||||
self.record(url, deadline);
|
||||
let stored = match budget_from_extensions(extensions) {
|
||||
Some(clock) => {
|
||||
let clamped = deadline.min(clock.hard_stop());
|
||||
(clamped > SystemTime::now()).then_some(clamped)
|
||||
}
|
||||
None if self.budget.is_zero() => Some(deadline),
|
||||
None => {
|
||||
let clamped = deadline.min(SystemTime::now() + self.budget);
|
||||
(clamped > SystemTime::now()).then_some(clamped)
|
||||
}
|
||||
};
|
||||
self.record(url, stored);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -176,7 +305,18 @@ impl RetryAfterMiddleware {
|
||||
|
||||
#[cfg(test)]
|
||||
fn record_test(&self, url: Url, deadline: SystemTime) {
|
||||
self.record(url, deadline);
|
||||
self.record(url, Some(deadline));
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn record_with_budget_test(&self, url: Url, budget: Duration, header: Option<&str>) {
|
||||
let mut extensions = Extensions::new();
|
||||
extensions.insert(BudgetClock::start(budget));
|
||||
let response = crate::client::retry_after::tests::synthetic_response(
|
||||
StatusCode::TOO_MANY_REQUESTS,
|
||||
header,
|
||||
);
|
||||
self.record_if_throttled(url, &response, &extensions);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -203,22 +343,22 @@ impl Middleware for RetryAfterMiddleware {
|
||||
next: Next<'_>,
|
||||
) -> Result<Response> {
|
||||
let req_url = req.url().clone();
|
||||
self.maybe_sleep_for(&req_url).await;
|
||||
self.maybe_sleep_for(&req_url, extensions).await;
|
||||
let response = next.run(req, extensions).await?;
|
||||
self.record_if_throttled(response.url().clone(), &response);
|
||||
self.record_if_throttled(response.url().clone(), &response, extensions);
|
||||
Ok(response)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
pub(crate) mod tests {
|
||||
use super::*;
|
||||
|
||||
fn url(s: &str) -> Url {
|
||||
Url::parse(s).unwrap()
|
||||
}
|
||||
|
||||
fn synthetic_response(status: StatusCode, retry_after: Option<&str>) -> Response {
|
||||
pub(crate) fn synthetic_response(status: StatusCode, retry_after: Option<&str>) -> Response {
|
||||
let mut builder = http::Response::builder().status(status);
|
||||
if let Some(value) = retry_after {
|
||||
builder = builder.header(RETRY_AFTER_HEADER, value);
|
||||
@@ -231,7 +371,7 @@ mod tests {
|
||||
let mw = RetryAfterMiddleware::with_capacity_and_ceiling(8, Duration::from_secs(300));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("5"));
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
let deadline = mw.deadline_for_test(&target).expect("seconds parse");
|
||||
let now = SystemTime::now();
|
||||
assert!(deadline > now);
|
||||
@@ -246,7 +386,7 @@ mod tests {
|
||||
StatusCode::SERVICE_UNAVAILABLE,
|
||||
Some("Wed, 21 Oct 2099 07:28:00 GMT"),
|
||||
);
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
let deadline = mw.deadline_for_test(&target).expect("HTTP-date parses");
|
||||
let ceiling = SystemTime::now() + Duration::from_secs(300);
|
||||
assert!(
|
||||
@@ -264,7 +404,7 @@ mod tests {
|
||||
StatusCode::TOO_MANY_REQUESTS,
|
||||
Some("Wed, 21 Oct 2015 07:28:00 GMT"),
|
||||
);
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
assert!(
|
||||
mw.deadline_for_test(&target).is_none(),
|
||||
"a deadline already in the past must not be recorded"
|
||||
@@ -283,7 +423,7 @@ mod tests {
|
||||
let mw = RetryAfterMiddleware::with_capacity_and_ceiling(8, Duration::from_secs(300));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("315360000"));
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
let deadline = mw
|
||||
.deadline_for_test(&target)
|
||||
.expect("a huge Retry-After still records (clamped)");
|
||||
@@ -300,7 +440,7 @@ mod tests {
|
||||
let mw = RetryAfterMiddleware::with_capacity_and_ceiling(8, Duration::from_secs(5));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("3600"));
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
let deadline = mw
|
||||
.deadline_for_test(&target)
|
||||
.expect("clamped deadline kept");
|
||||
@@ -377,7 +517,7 @@ mod tests {
|
||||
let origin = url("https://api.example.com/v1/chat");
|
||||
let redirector = url("https://redirector.example.com/429");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("5"));
|
||||
mw.record_if_throttled(origin.clone(), &response);
|
||||
mw.record_if_throttled(origin.clone(), &response, &Extensions::new());
|
||||
assert!(
|
||||
mw.deadline_for_test(&origin).is_some(),
|
||||
"the effective (post-redirect) URL carries the deadline"
|
||||
@@ -390,7 +530,7 @@ mod tests {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity(8));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::OK, Some("5"));
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
assert!(mw.deadline_for_test(&target).is_none());
|
||||
}
|
||||
|
||||
@@ -399,7 +539,7 @@ mod tests {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity(8));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, None);
|
||||
mw.record_if_throttled(target.clone(), &response);
|
||||
mw.record_if_throttled(target.clone(), &response, &Extensions::new());
|
||||
assert!(mw.deadline_for_test(&target).is_none());
|
||||
}
|
||||
|
||||
@@ -412,7 +552,7 @@ mod tests {
|
||||
SystemTime::now() + Duration::from_millis(50),
|
||||
);
|
||||
let started = SystemTime::now();
|
||||
mw.maybe_sleep_for(&target).await;
|
||||
mw.maybe_sleep_for(&target, &Extensions::new()).await;
|
||||
let elapsed = SystemTime::now().duration_since(started).unwrap();
|
||||
assert!(
|
||||
elapsed >= Duration::from_millis(37),
|
||||
@@ -428,7 +568,7 @@ mod tests {
|
||||
let remaining = Duration::from_secs(4);
|
||||
mw.record_test(target.clone(), SystemTime::now() + remaining);
|
||||
let started = SystemTime::now();
|
||||
mw.maybe_sleep_for(&target).await;
|
||||
mw.maybe_sleep_for(&target, &Extensions::new()).await;
|
||||
let elapsed = SystemTime::now().duration_since(started).unwrap();
|
||||
let max_jitter = max_sleep_jitter(remaining);
|
||||
assert!(
|
||||
@@ -451,4 +591,125 @@ mod tests {
|
||||
"jitter is capped at 2 s even for long waits"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn sleep_is_truncated_to_the_retry_budget() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
8,
|
||||
Duration::from_secs(300),
|
||||
Duration::from_secs(1),
|
||||
));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
mw.record_test(target.clone(), SystemTime::now() + Duration::from_secs(60));
|
||||
let mut extensions = Extensions::new();
|
||||
extensions.insert(BudgetClock::start(Duration::from_secs(1)));
|
||||
let started = Instant::now();
|
||||
mw.maybe_sleep_for(&target, &extensions).await;
|
||||
let elapsed = started.elapsed();
|
||||
assert!(
|
||||
elapsed <= Duration::from_millis(970),
|
||||
"the sleep must be truncated to the remaining budget, took {elapsed:?}"
|
||||
);
|
||||
assert!(elapsed >= Duration::from_millis(700));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn budget_spent_skips_the_sleep_entirely() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
8,
|
||||
Duration::from_secs(300),
|
||||
Duration::from_secs(1),
|
||||
));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
mw.record_test(target.clone(), SystemTime::now() + Duration::from_secs(60));
|
||||
let expired = BudgetClock::for_duration(
|
||||
SystemTime::now() - Duration::from_secs(5),
|
||||
Duration::from_secs(1),
|
||||
);
|
||||
let mut extensions = Extensions::new();
|
||||
extensions.insert(expired);
|
||||
let started = Instant::now();
|
||||
mw.maybe_sleep_for(&target, &extensions).await;
|
||||
assert!(
|
||||
started.elapsed() <= Duration::from_millis(60),
|
||||
"a spent budget must skip the Retry-After sleep, took {:?}",
|
||||
started.elapsed()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn throttle_refresh_is_clamped_to_the_budget_hard_stop() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
8,
|
||||
Duration::from_secs(300),
|
||||
Duration::from_secs(1),
|
||||
));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let response = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("300"));
|
||||
let mut extensions = Extensions::new();
|
||||
extensions.insert(BudgetClock::start(Duration::from_secs(1)));
|
||||
mw.record_if_throttled(target.clone(), &response, &extensions);
|
||||
let deadline = mw
|
||||
.deadline_for_test(&target)
|
||||
.expect("a live deadline records (clamped to the hard stop)");
|
||||
assert!(
|
||||
deadline <= SystemTime::now() + Duration::from_millis(1100),
|
||||
"a retry-storm refresh must not push the deadline past the budget, got {deadline:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn throttle_storm_cannot_extend_the_first_seen_deadline() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
8,
|
||||
Duration::from_secs(300),
|
||||
Duration::from_secs(300),
|
||||
));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
let first = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("10"));
|
||||
mw.record_if_throttled(target.clone(), &first, &Extensions::new());
|
||||
let first_deadline = mw.deadline_for_test(&target).expect("first deadline");
|
||||
let storm = synthetic_response(StatusCode::TOO_MANY_REQUESTS, Some("300"));
|
||||
mw.record_if_throttled(target.clone(), &storm, &Extensions::new());
|
||||
assert_eq!(
|
||||
mw.deadline_for_test(&target),
|
||||
Some(first_deadline),
|
||||
"a later higher Retry-After must not extend the first-seen deadline"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn budget_exhaustion_drops_the_throttle_entry() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity_ceiling_and_budget(
|
||||
8,
|
||||
Duration::from_secs(300),
|
||||
Duration::from_secs(300),
|
||||
));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
mw.record_with_budget_test(target.clone(), Duration::from_secs(300), None);
|
||||
assert!(
|
||||
mw.deadline_for_test(&target).is_none(),
|
||||
"a refresh arriving after the budget is spent must drop the entry"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn missing_header_refresh_after_budget_exhaustion_removes_the_deadline() {
|
||||
let mw = std::sync::Arc::new(RetryAfterMiddleware::with_capacity(8));
|
||||
let target = url("https://api.example.com/v1/chat");
|
||||
mw.record_test(target.clone(), SystemTime::now() + Duration::from_secs(10));
|
||||
mw.record_with_budget_test(target.clone(), Duration::from_secs(300), None);
|
||||
assert!(mw.deadline_for_test(&target).is_some());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn budget_clock_hard_stop_survives_a_wall_clock_step() {
|
||||
let clock = BudgetClock::start(Duration::from_secs(30));
|
||||
std::thread::sleep(Duration::from_millis(30));
|
||||
let projected = clock.hard_stop();
|
||||
assert!(
|
||||
projected <= SystemTime::now() + Duration::from_secs(30),
|
||||
"hard_stop must track monotonic progress against a stepped wall clock, got {projected:?}"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,126 @@
|
||||
//! Budget-aware `Retry-After` wire coverage (review-002 CLI-01): a
|
||||
//! counting responder that always answers `429` with a large
|
||||
//! `Retry-After` must not extend the caller's wall time past
|
||||
//! `max_total_retry_duration` + one attempt's request time. The
|
||||
//! middleware-level truncation is pinned here through the real
|
||||
//! `SharedHttpClient` stack; the per-URL throttle map semantics for
|
||||
//! separate logical requests stay in `src/client/retry_after.rs`.
|
||||
|
||||
use std::sync::atomic::{AtomicU32, Ordering};
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
use alkhttp::client::{HttpClientConfig, SharedHttpClient};
|
||||
|
||||
struct CountingResponder {
|
||||
addr: std::net::SocketAddr,
|
||||
hits: Arc<AtomicU32>,
|
||||
}
|
||||
|
||||
impl CountingResponder {
|
||||
async fn spawn(status_line: &'static str, extra_headers: &'static str) -> Self {
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
|
||||
.await
|
||||
.expect("responder binds");
|
||||
let addr = listener.local_addr().expect("responder address");
|
||||
let hits = Arc::new(AtomicU32::new(0));
|
||||
let hits_loop = Arc::clone(&hits);
|
||||
tokio::spawn(async move {
|
||||
loop {
|
||||
let Ok((mut sock, _)) = listener.accept().await else {
|
||||
return;
|
||||
};
|
||||
hits_loop.fetch_add(1, Ordering::SeqCst);
|
||||
let body = "";
|
||||
let response = format!(
|
||||
"{status_line}\r\ncontent-length: {}\r\nconnection: close\r\n{extra_headers}\r\n",
|
||||
body.len()
|
||||
);
|
||||
let _ = tokio::io::AsyncWriteExt::write_all(&mut sock, response.as_bytes()).await;
|
||||
tokio::io::AsyncWriteExt::shutdown(&mut sock).await.ok();
|
||||
}
|
||||
});
|
||||
Self { addr, hits }
|
||||
}
|
||||
|
||||
fn hits(&self) -> u32 {
|
||||
self.hits.load(Ordering::SeqCst)
|
||||
}
|
||||
}
|
||||
|
||||
fn throttling_config() -> HttpClientConfig {
|
||||
HttpClientConfig {
|
||||
request_timeout: Some(Duration::from_secs(5)),
|
||||
connect_timeout: Some(Duration::from_secs(2)),
|
||||
read_timeout: Some(Duration::from_secs(2)),
|
||||
retry: alkhttp::client::RetryConfig {
|
||||
max_retries: 50,
|
||||
initial_backoff: Duration::from_millis(10),
|
||||
max_retry_interval: Duration::from_millis(50),
|
||||
},
|
||||
max_total_retry_duration: Duration::from_secs(1),
|
||||
retry_after_ceiling: Duration::from_secs(300),
|
||||
..HttpClientConfig::default()
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn always_throttled_request_is_bounded_by_the_retry_budget() {
|
||||
let responder =
|
||||
CountingResponder::spawn("HTTP/1.1 429 Too Many Requests", "retry-after: 300\r\n").await;
|
||||
let http = SharedHttpClient::new(throttling_config()).expect("client builds");
|
||||
let started = Instant::now();
|
||||
let response = http
|
||||
.client()
|
||||
.get(format!("http://{}/throttled", responder.addr))
|
||||
.send()
|
||||
.await
|
||||
.expect("429 is a delivered response, not a transport error");
|
||||
let wall = started.elapsed();
|
||||
assert_eq!(response.status(), 429);
|
||||
assert!(
|
||||
responder.hits() >= 2,
|
||||
"the request retries into the throttle (hits: {})",
|
||||
responder.hits()
|
||||
);
|
||||
assert!(
|
||||
wall < Duration::from_secs(3),
|
||||
"wall time must be bounded by max_total_retry_duration + one attempt, took {wall:?} over {} hits",
|
||||
responder.hits()
|
||||
);
|
||||
assert!(
|
||||
responder.hits() <= 60,
|
||||
"the attempt cap is generous but not infinite, got {}",
|
||||
responder.hits()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn separate_logical_request_still_honors_the_recorded_window() {
|
||||
let responder =
|
||||
CountingResponder::spawn("HTTP/1.1 429 Too Many Requests", "retry-after: 1\r\n").await;
|
||||
let http = SharedHttpClient::new(throttling_config()).expect("client builds");
|
||||
let url = format!("http://{}/window", responder.addr);
|
||||
|
||||
let first_started = Instant::now();
|
||||
let first = http.client().get(&url).send().await.expect("first 429");
|
||||
let first_wall = first_started.elapsed();
|
||||
assert_eq!(first.status(), 429);
|
||||
assert!(
|
||||
first_wall >= Duration::from_millis(700),
|
||||
"the first logical request sleeps within its own budget, took {first_wall:?}"
|
||||
);
|
||||
|
||||
let second_started = Instant::now();
|
||||
let second = http.client().get(&url).send().await.expect("second 429");
|
||||
let second_wall = second_started.elapsed();
|
||||
assert_eq!(second.status(), 429);
|
||||
assert!(
|
||||
second_wall <= Duration::from_millis(1100),
|
||||
"a fresh logical request honors the recorded throttle window (shortened by the budget clamp), took {second_wall:?}"
|
||||
);
|
||||
assert!(
|
||||
responder.hits() >= 2,
|
||||
"both logical requests reached the upstream"
|
||||
);
|
||||
}
|
||||
Reference in New Issue
Block a user