use super::*;
use std::ops::ControlFlow;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SuccessClass {
Strict,
AllowRedirects,
}
pub fn retry_http_blocking<F, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
success_class: SuccessClass,
send: F,
error_msg: M,
) -> anyhow::Result<(reqwest::StatusCode, String)>
where
F: FnMut(u32) -> Result<reqwest::blocking::Response, reqwest::Error>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
retry_http_blocking_deadline(rlog, policy, None, success_class, send, error_msg)
}
pub fn retry_http_blocking_deadline<F, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
deadline: Option<std::time::Instant>,
success_class: SuccessClass,
mut send: F,
error_msg: M,
) -> anyhow::Result<(reqwest::StatusCode, String)>
where
F: FnMut(u32) -> Result<reqwest::blocking::Response, reqwest::Error>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
use anyhow::Context as _;
retry_sync_deadline(rlog, policy, deadline, |attempt| {
match send(attempt) {
Ok(resp) => {
let status = resp.status();
let succeeded = match success_class {
SuccessClass::Strict => status.is_success(),
SuccessClass::AllowRedirects => status.is_success() || status.is_redirection(),
};
let body = resp
.text()
.unwrap_or_else(|e| format!("<failed to read body: {e}>"));
if succeeded {
Ok((status, body))
} else {
let msg = error_msg(status, &body);
let inner = anyhow::anyhow!("{msg}");
let wrapped = anyhow::Error::new(HttpError::new(
std::io::Error::other(inner.to_string()),
status.as_u16(),
))
.context(inner);
if is_retriable(wrapped.as_ref()) {
Err(ControlFlow::Continue(wrapped))
} else {
Err(ControlFlow::Break(wrapped))
}
}
}
Err(e) => {
let err = anyhow::Error::new(HttpError::from_response(e, None))
.context(format!("{}: HTTP transport error", rlog.desc()));
if is_retriable(err.as_ref()) {
Err(ControlFlow::Continue(err))
} else {
Err(ControlFlow::Break(err))
}
}
}
})
.with_context(|| format!("{}: exhausted retry attempts", rlog.desc()))
}
pub fn retry_http_blocking_bytes<F, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
success_class: SuccessClass,
send: F,
error_msg: M,
) -> anyhow::Result<(reqwest::StatusCode, Vec<u8>)>
where
F: FnMut(u32) -> Result<reqwest::blocking::Response, reqwest::Error>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
retry_http_blocking_bytes_deadline(rlog, policy, None, success_class, send, error_msg)
}
pub fn retry_http_blocking_bytes_deadline<F, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
deadline: Option<std::time::Instant>,
success_class: SuccessClass,
mut send: F,
error_msg: M,
) -> anyhow::Result<(reqwest::StatusCode, Vec<u8>)>
where
F: FnMut(u32) -> Result<reqwest::blocking::Response, reqwest::Error>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
use anyhow::Context as _;
retry_sync_deadline(rlog, policy, deadline, |attempt| match send(attempt) {
Ok(resp) => {
let status = resp.status();
let succeeded = match success_class {
SuccessClass::Strict => status.is_success(),
SuccessClass::AllowRedirects => status.is_success() || status.is_redirection(),
};
let bytes = resp
.bytes()
.map(|b| b.to_vec())
.unwrap_or_else(|e| format!("<failed to read body: {e}>").into_bytes());
if succeeded {
Ok((status, bytes))
} else {
let body_text = String::from_utf8_lossy(&bytes).into_owned();
let msg = error_msg(status, &body_text);
let inner = anyhow::anyhow!("{msg}");
let wrapped = anyhow::Error::new(HttpError::new(
std::io::Error::other(inner.to_string()),
status.as_u16(),
))
.context(inner);
if is_retriable(wrapped.as_ref()) {
Err(ControlFlow::Continue(wrapped))
} else {
Err(ControlFlow::Break(wrapped))
}
}
}
Err(e) => {
let err = anyhow::Error::new(HttpError::from_response(e, None))
.context(format!("{}: HTTP transport error", rlog.desc()));
if is_retriable(err.as_ref()) {
Err(ControlFlow::Continue(err))
} else {
Err(ControlFlow::Break(err))
}
}
})
.with_context(|| format!("{}: exhausted retry attempts", rlog.desc()))
}
pub async fn retry_http_async<F, Fut, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
success_class: SuccessClass,
send: F,
error_msg: M,
) -> anyhow::Result<reqwest::Response>
where
F: FnMut(u32) -> Fut,
Fut: std::future::Future<Output = Result<reqwest::Response, reqwest::Error>>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
retry_http_async_deadline(rlog, policy, None, success_class, send, error_msg).await
}
pub async fn retry_http_async_deadline<F, Fut, M>(
rlog: RetryLog<'_>,
policy: &RetryPolicy,
deadline: Option<std::time::Instant>,
success_class: SuccessClass,
mut send: F,
error_msg: M,
) -> anyhow::Result<reqwest::Response>
where
F: FnMut(u32) -> Fut,
Fut: std::future::Future<Output = Result<reqwest::Response, reqwest::Error>>,
M: Fn(reqwest::StatusCode, &str) -> String,
{
use anyhow::Context as _;
retry_async_deadline(rlog, policy, deadline, |attempt| {
let fut = send(attempt);
let error_msg = &error_msg;
async move {
match fut.await {
Ok(resp) => {
let status = resp.status();
let succeeded = match success_class {
SuccessClass::Strict => status.is_success(),
SuccessClass::AllowRedirects => {
status.is_success() || status.is_redirection()
}
};
if succeeded {
Ok(resp)
} else {
let body = resp
.text()
.await
.unwrap_or_else(|e| format!("<failed to read body: {e}>"));
let msg = error_msg(status, &body);
let inner = anyhow::anyhow!("{msg}");
let wrapped = anyhow::Error::new(HttpError::new(
std::io::Error::other(inner.to_string()),
status.as_u16(),
))
.context(inner);
if is_retriable(wrapped.as_ref()) {
Err(ControlFlow::Continue(wrapped))
} else {
Err(ControlFlow::Break(wrapped))
}
}
}
Err(e) => {
let err = anyhow::Error::new(HttpError::from_response(e, None))
.context(format!("{}: HTTP transport error", rlog.desc()));
if is_retriable(err.as_ref()) {
Err(ControlFlow::Continue(err))
} else {
Err(ControlFlow::Break(err))
}
}
}
}
})
.await
.with_context(|| format!("{}: exhausted retry attempts", rlog.desc()))
}
pub fn classify_http_sync(
result: reqwest::Result<reqwest::blocking::Response>,
) -> Result<reqwest::blocking::Response, ControlFlow<anyhow::Error, anyhow::Error>> {
use anyhow::anyhow;
match result {
Ok(resp) => {
let status = resp.status();
if status.is_success() || status.is_redirection() {
Ok(resp)
} else if status.is_server_error() {
Err(ControlFlow::Continue(anyhow!(
"HTTP {} {}",
status.as_u16(),
status.canonical_reason().unwrap_or("server error")
)))
} else {
Err(ControlFlow::Break(anyhow!(
"HTTP {} {}",
status.as_u16(),
status.canonical_reason().unwrap_or("client error")
)))
}
}
Err(e) => Err(ControlFlow::Continue(anyhow!(e))),
}
}