mod attempt;
mod auth;
mod breaker;
mod builder;
mod error;
mod observer;
mod policy;
mod rate;
pub use auth::{AuthFuture, AuthRefresher};
pub use breaker::BreakerState;
pub use builder::ResilientClientBuilder;
pub use error::{ClientError, ErrorKind};
pub use observer::{TransportEvent, TransportObserver};
pub use policy::{BreakerMode, CallPolicy, RetryMode, RetryPolicy};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use tokio::sync::Semaphore;
use attempt::AttemptLoop;
use auth::AuthState;
use breaker::CircuitBreaker;
#[derive(Clone)]
pub struct ResilientClient {
http: reqwest::Client,
semaphore: Arc<Semaphore>,
policy: RetryPolicy,
breaker: Option<Arc<CircuitBreaker>>,
auth: Option<Arc<AuthState>>,
in_flight: Arc<AtomicUsize>,
peak_in_flight: Arc<AtomicUsize>,
observer: Option<(Arc<str>, Arc<dyn TransportObserver>)>,
rate: Option<Arc<rate::TokenBucket>>,
}
impl ResilientClient {
pub fn builder() -> ResilientClientBuilder {
ResilientClientBuilder::default()
}
pub async fn execute<F>(&self, build: F) -> Result<reqwest::Response, ClientError>
where
F: Fn(&reqwest::Client) -> reqwest::RequestBuilder,
{
let policy = CallPolicy::bounded(self.policy.clone());
self.execute_with(&policy, build).await
}
pub async fn execute_with<F>(
&self,
policy: &CallPolicy,
build: F,
) -> Result<reqwest::Response, ClientError>
where
F: Fn(&reqwest::Client) -> reqwest::RequestBuilder,
{
AttemptLoop::run(self, policy, &build).await
}
pub fn in_flight(&self) -> usize {
self.in_flight.load(Ordering::SeqCst)
}
pub fn peak_in_flight(&self) -> usize {
self.peak_in_flight.load(Ordering::SeqCst)
}
pub fn breaker_state(&self) -> Option<BreakerState> {
self.breaker.as_ref().map(|b| b.state())
}
pub fn key(&self) -> Option<&str> {
self.observer.as_ref().map(|(k, _)| &**k)
}
pub fn report(&self, event: TransportEvent) {
if let Some((key, obs)) = &self.observer {
obs.on_event(key, event);
}
}
fn enter_flight(&self) -> InFlightGuard {
let now = self.in_flight.fetch_add(1, Ordering::SeqCst) + 1;
self.peak_in_flight.fetch_max(now, Ordering::SeqCst);
InFlightGuard(Arc::clone(&self.in_flight))
}
}
impl std::fmt::Debug for ResilientClient {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ResilientClient")
.field("key", &self.key())
.field("in_flight", &self.in_flight())
.finish_non_exhaustive()
}
}
struct InFlightGuard(Arc<AtomicUsize>);
impl Drop for InFlightGuard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}