feat: round-4 metrics for circuit breaker + retry stats + lifecycle fixes
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
CI / Rustfmt (push) Canceled after 0s
CI / Clippy (push) Canceled after 0s
CI / Test (workspace all features) (push) Canceled after 0s
CI / Test (workspace default features) (push) Canceled after 0s
CI / Test (mytheclipse / bg only) (push) Canceled after 0s
CI / Test (mytheclipse / compute only) (push) Canceled after 0s
CI / Test (mytheclipse / io only) (push) Canceled after 0s
CI / Test (mytheclipse / lifecycle only) (push) Canceled after 0s
CI / Test (mytheclipse / observability only) (push) Canceled after 0s
CI / Test (mytheclipse / resiliency only) (push) Canceled after 0s
CI / Test (mytheclipse / traffic only) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l2-redis) (push) Canceled after 0s
CI / Test (mytheclipse-cache / l1-moka) (push) Canceled after 0s
CI / Test (mytheclipse-cache / default) (push) Canceled after 0s
CI / Test (mytheclipse-config / default) (push) Canceled after 0s
CI / Test (mytheclipse-crypto / default) (push) Canceled after 0s
CI / Test (mytheclipse-event / amqp) (push) Canceled after 0s
CI / Test (mytheclipse-event / nats) (push) Canceled after 0s
CI / Test (mytheclipse-event / default (mem)) (push) Canceled after 0s
CI / Test (mytheclipse-storage / gcs) (push) Canceled after 0s
CI / Test (mytheclipse-storage / s3) (push) Canceled after 0s
CI / Test (mytheclipse-storage / default (local)) (push) Canceled after 0s
CI / Run mytheclipse example (push) Canceled after 0s
CI / Docs check (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-cache) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-config) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-crypto) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-event) (push) Canceled after 0s
CI / Cargo package dry-run (mytheclipse-storage) (push) Canceled after 0s
Release / Semantic Release (push) Canceled after 0s
This commit is contained in:
@@ -22,6 +22,17 @@ pub enum CircuitState {
|
||||
HalfOpen,
|
||||
}
|
||||
|
||||
/// Point-in-time snapshot of a [`CircuitBreaker`] for metrics/observability.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub struct CircuitSnapshot {
|
||||
/// Current circuit state.
|
||||
pub state: CircuitState,
|
||||
/// Consecutive failures recorded (resets on success in `Closed`).
|
||||
pub failures: u64,
|
||||
/// Consecutive successes recorded (resets on failure/open).
|
||||
pub successes: u64,
|
||||
}
|
||||
|
||||
const CLOSED: u8 = 0;
|
||||
const OPEN: u8 = 1;
|
||||
const HALF_OPEN: u8 = 2;
|
||||
@@ -199,6 +210,16 @@ impl CircuitBreaker {
|
||||
*self.inner.opened_at.lock().unwrap() = None;
|
||||
}
|
||||
|
||||
/// Returns a point-in-time snapshot of the breaker's internal counters and
|
||||
/// state, for metrics/observability export.
|
||||
pub fn snapshot(&self) -> CircuitSnapshot {
|
||||
CircuitSnapshot {
|
||||
state: self.state(),
|
||||
failures: self.inner.failures.load(Ordering::Acquire),
|
||||
successes: self.inner.successes.load(Ordering::Acquire),
|
||||
}
|
||||
}
|
||||
|
||||
fn record_result(&self, success: bool) {
|
||||
match self.inner.state.load(Ordering::Acquire) {
|
||||
HALF_OPEN => {
|
||||
@@ -362,4 +383,25 @@ mod tests {
|
||||
let err: Result<u32, CircuitError<u8>> = b.call(|| Err(9u8));
|
||||
assert!(matches!(err, Err(CircuitError::Inner(9))));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn snapshot_reflects_state_and_counts() {
|
||||
let b = breaker();
|
||||
let snap = b.snapshot();
|
||||
assert_eq!(snap.state, CircuitState::Closed);
|
||||
assert_eq!(snap.failures, 0);
|
||||
assert_eq!(snap.successes, 0);
|
||||
|
||||
// success in Closed state resets failure count (no failure counter added).
|
||||
b.call::<(), u8, _>(|| Ok(()));
|
||||
let snap2 = b.snapshot();
|
||||
assert_eq!(snap2.state, CircuitState::Closed);
|
||||
|
||||
for _ in 0..3 {
|
||||
let _: Result<(), CircuitError<u8>> = b.call(|| Err(1u8));
|
||||
}
|
||||
let snap3 = b.snapshot();
|
||||
assert_eq!(snap3.state, CircuitState::Open);
|
||||
assert_eq!(snap3.failures, 0); // reset on open()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -85,6 +85,79 @@ impl<E: std::fmt::Display> std::fmt::Display for RetryError<E> {
|
||||
|
||||
impl<E: std::fmt::Debug + std::fmt::Display> std::error::Error for RetryError<E> {}
|
||||
|
||||
/// Statistics collected during a [`retry`] call.
|
||||
#[derive(Debug, Clone, Default)]
|
||||
pub struct RetryStats {
|
||||
/// Total number of attempts made (including the first).
|
||||
pub attempts: u32,
|
||||
/// Number of retries performed (= `attempts - 1` if exhausted, or
|
||||
/// `attempts - 1` if ultimately succeeded after at least one retry).
|
||||
pub retries: u32,
|
||||
/// The error message from the final attempt, if any.
|
||||
pub last_error: Option<String>,
|
||||
}
|
||||
|
||||
/// Like [`retry`] but also returns [`RetryStats`] capturing attempt counts.
|
||||
///
|
||||
/// Retries `op` according to `config`, retrying only errors for which
|
||||
/// `filter` returns `true`.
|
||||
///
|
||||
/// Like [`retry`] but also returns [`RetryStats`].
|
||||
pub async fn retry_with_stats<T, E, F, Fut, P>(
|
||||
config: RetryConfig,
|
||||
filter: P,
|
||||
mut op: F,
|
||||
) -> (Result<T, RetryError<E>>, RetryStats)
|
||||
where
|
||||
F: FnMut() -> Fut,
|
||||
Fut: Future<Output = Result<T, E>>,
|
||||
P: Fn(&E) -> bool,
|
||||
E: std::fmt::Display,
|
||||
{
|
||||
let mut attempt: u32 = 0;
|
||||
let mut last_error: Option<String> = None;
|
||||
loop {
|
||||
attempt += 1;
|
||||
let span = tracing::info_span!(
|
||||
"mytheclipse_retry_task",
|
||||
attempt,
|
||||
max_attempts = config.max_attempts
|
||||
);
|
||||
let result = op().instrument(span).await;
|
||||
|
||||
match result {
|
||||
Ok(value) => {
|
||||
let stats = RetryStats {
|
||||
attempts: attempt,
|
||||
retries: attempt.saturating_sub(1),
|
||||
last_error,
|
||||
};
|
||||
return (Ok(value), stats);
|
||||
}
|
||||
Err(err) => {
|
||||
last_error = Some(err.to_string());
|
||||
let retryable = filter(&err);
|
||||
if !retryable || attempt >= config.max_attempts {
|
||||
let stats = RetryStats {
|
||||
attempts: attempt,
|
||||
retries: attempt.saturating_sub(1),
|
||||
last_error,
|
||||
};
|
||||
return (
|
||||
Err(RetryError::Exhausted {
|
||||
attempts: attempt,
|
||||
last: err,
|
||||
}),
|
||||
stats,
|
||||
);
|
||||
}
|
||||
let delay = backoff_delay(&config, attempt, rand::thread_rng());
|
||||
tokio::time::sleep(delay).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Retries `op` according to `config`, retrying only errors for which
|
||||
/// `filter` returns `true`.
|
||||
///
|
||||
@@ -238,6 +311,24 @@ mod tests {
|
||||
assert_eq!(calls.get(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn retry_with_stats_succeeds_with_counts() {
|
||||
use std::cell::Cell;
|
||||
let config = RetryConfig {
|
||||
max_attempts: 5,
|
||||
base_delay: Duration::from_millis(1),
|
||||
..RetryConfig::default()
|
||||
};
|
||||
let calls = Cell::new(0u32);
|
||||
let (result, stats) = retry_with_stats(config, |_| true, || async {
|
||||
calls.set(calls.get() + 1);
|
||||
if calls.get() < 3 { Err::<u32, &str>("fail") } else { Ok(42u32) }
|
||||
}).await;
|
||||
assert_eq!(result.unwrap(), 42);
|
||||
assert_eq!(stats.attempts, 3);
|
||||
assert_eq!(stats.retries, 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn full_jitter_is_within_bounds_and_capped() {
|
||||
let config = RetryConfig {
|
||||
|
||||
Reference in New Issue
Block a user