use std::time::Duration;
use serde::Deserialize;
const MAX_RETRY_SECS: f64 = 60.0;
const MAX_RETRIES: u32 = 3;
#[derive(Deserialize, Default)]
struct RateLimitBody {
retry_after: Option<f64>,
}
pub(crate) async fn send_with_retry<F>(
context: &str,
make_req: F,
) -> Result<reqwest::Response, reqwest::Error>
where
F: Fn() -> reqwest::RequestBuilder,
{
let mut attempts = 0u32;
loop {
let resp = make_req().send().await?;
if resp.status() != reqwest::StatusCode::TOO_MANY_REQUESTS {
return resp.error_for_status();
}
let header_secs = resp
.headers()
.get(reqwest::header::RETRY_AFTER)
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse::<f64>().ok());
let body_secs = resp
.json::<RateLimitBody>()
.await
.unwrap_or_default()
.retry_after;
let delay_secs = header_secs.or(body_secs).unwrap_or(1.0).min(MAX_RETRY_SECS);
attempts += 1;
if attempts > MAX_RETRIES {
tracing::warn!(
context,
delay_secs,
attempts,
"rate-limited and retries exhausted"
);
return make_req().send().await?.error_for_status();
}
tracing::warn!(
context,
delay_secs,
attempt = attempts,
max = MAX_RETRIES,
"rate-limited (429), backing off"
);
tokio::time::sleep(Duration::from_secs_f64(delay_secs)).await;
}
}
#[cfg(test)]
mod tests {
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
use super::*;
#[tokio::test]
async fn send_with_retry_succeeds_on_200() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": "1"})))
.mount(&server)
.await;
let client = reqwest::Client::new();
let url = format!("{}/channels/ch1/messages", server.uri());
let resp = send_with_retry("test", || client.post(&url)).await.unwrap();
assert_eq!(resp.status(), 200);
}
#[tokio::test]
async fn send_with_retry_retries_on_429_then_succeeds() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(
ResponseTemplate::new(429)
.append_header("Retry-After", "0")
.set_body_json(serde_json::json!({"retry_after": 0.0})),
)
.up_to_n_times(1)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": "2"})))
.mount(&server)
.await;
let client = reqwest::Client::new();
let url = format!("{}/channels/ch1/messages", server.uri());
let resp = send_with_retry("test", || client.post(&url)).await.unwrap();
assert_eq!(resp.status(), 200);
}
#[tokio::test]
async fn send_with_retry_uses_body_retry_after_when_no_header() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(
ResponseTemplate::new(429).set_body_json(serde_json::json!({"retry_after": 0.0})),
)
.up_to_n_times(3)
.mount(&server)
.await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": "3"})))
.mount(&server)
.await;
let client = reqwest::Client::new();
let url = format!("{}/channels/ch1/messages", server.uri());
let resp = send_with_retry("test", || client.post(&url)).await.unwrap();
assert_eq!(resp.status(), 200);
}
#[tokio::test]
async fn send_with_retry_propagates_non_429_errors() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(ResponseTemplate::new(403))
.mount(&server)
.await;
let client = reqwest::Client::new();
let url = format!("{}/channels/ch1/messages", server.uri());
let result = send_with_retry("test", || client.post(&url)).await;
assert!(result.is_err());
let err = result.unwrap_err();
assert_eq!(err.status(), Some(reqwest::StatusCode::FORBIDDEN));
}
#[tokio::test]
async fn send_with_retry_errors_when_retries_exhausted() {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/channels/ch1/messages"))
.respond_with(
ResponseTemplate::new(429)
.append_header("Retry-After", "0")
.set_body_json(serde_json::json!({"retry_after": 0.0})),
)
.mount(&server)
.await;
let client = reqwest::Client::new();
let url = format!("{}/channels/ch1/messages", server.uri());
let result = send_with_retry("test", || client.post(&url)).await;
assert!(result.is_err());
let err = result.unwrap_err();
assert_eq!(err.status(), Some(reqwest::StatusCode::TOO_MANY_REQUESTS));
}
#[test]
fn rate_limit_body_defaults_to_none() {
let body: RateLimitBody = serde_json::from_str("{}").unwrap();
assert!(body.retry_after.is_none());
}
#[test]
fn rate_limit_body_parses_float() {
let body: RateLimitBody = serde_json::from_str(r#"{"retry_after": 1.5}"#).unwrap();
assert!((body.retry_after.unwrap() - 1.5).abs() < f64::EPSILON);
}
#[test]
fn max_retry_secs_clamps() {
let unclamped: f64 = 120.0;
assert_eq!(
unclamped.min(MAX_RETRY_SECS).to_bits(),
MAX_RETRY_SECS.to_bits()
);
}
}