use std::sync::Arc;
use std::time::Instant;
use crate::activity::ActivityMetrics;
use crate::adapter::{AdapterError, ExecutionError, OpDispenser, OpResult, WrappingDispenser};
use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
pub const NAME: WrapperName = WrapperName::new("tries");
fn triggers(s: WrapperSubject) -> bool {
let Some(op) = s.op() else {
return false;
};
op.params.contains_key("tries")
}
fn describe_assignment(s: WrapperSubject) -> Option<String> {
let op = s.op()?;
op.params.get("tries").map(|v| format!("tries: {v}"))
}
inventory::submit! {
WrapperRegistration {
name: NAME,
owned_fields: &["tries", "retry_backoff", "retry_backoff_max",
"retry_backoff_ratio", "retry_exemplar_rate",
"retry_exemplar_max_hz", "retry_advisory"],
triggers,
requires_inner: &[],
forbids_outer: &[],
mutually_exclusive_with: &[],
describe_assignment,
levels: &[crate::wrapper_registry::WrapperLevel::Op],
}
}
pub struct TriesDispenser {
inner: Arc<dyn OpDispenser>,
tries: u32,
metrics: Arc<ActivityMetrics>,
backoff_base_ms: u64,
backoff_max_ms: u64,
backoff_ratio: f64,
stop: crate::session_signals::StopView,
op_name: String,
exemplars: crate::exec_events::ExemplarSampler,
advisory: Option<Arc<crate::exec_events::AdvisoryGate>>,
}
impl crate::exec_events::ExecEventSubscriber for TriesDispenser {}
impl TriesDispenser {
#[allow(clippy::too_many_arguments)]
pub fn wrap(
inner: Arc<dyn OpDispenser>,
tries: u32,
metrics: Arc<ActivityMetrics>,
backoff_base_ms: u64,
backoff_max_ms: u64,
backoff_ratio: f64,
stop: crate::session_signals::StopView,
op_name: String,
exemplars: crate::exec_events::ExemplarSampler,
advisory: Option<Arc<crate::exec_events::AdvisoryGate>>,
) -> Arc<dyn OpDispenser> {
Arc::new(Self {
inner,
tries,
metrics,
backoff_base_ms,
backoff_max_ms,
backoff_ratio,
stop,
op_name,
exemplars,
advisory,
})
}
}
async fn portable_sleep_ms(ms: u64) {
let (tx, rx) = futures::channel::oneshot::channel::<()>();
std::thread::spawn(move || {
std::thread::sleep(std::time::Duration::from_millis(ms));
let _ = tx.send(());
});
let _ = rx.await;
}
use crate::exec_events::splitmix64;
fn backoff_wait_ms(base_ms: u64, max_ms: u64, ratio: f64, attempt_no: u32, cycle: u64) -> u64 {
if base_ms == 0 {
return 0;
}
let mult = ratio.max(1.0).powi((attempt_no.saturating_sub(1)) as i32);
let raw = base_ms as f64 * mult;
let capped = raw.min(max_ms as f64).max(1.0) as u64;
let h = splitmix64(cycle ^ ((attempt_no as u64) << 48));
capped / 2 + (h % (capped / 2 + 1))
}
impl WrappingDispenser for TriesDispenser {}
impl OpDispenser for TriesDispenser {
fn execute<'a>(
&'a self,
cycle: u64,
ctx: &'a crate::fixture::ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
if self.tries == 0 {
self.metrics.attempt_total.inc();
self.metrics.attempt_failure.observe(0);
self.metrics.tries_histogram.record(0);
return Err(ExecutionError::Op(AdapterError {
error_name: "tries_zero".into(),
message: "tries: 0 — op is configured to fail without executing".into(),
retryable: false,
}));
}
let mut attempt_no: u32 = 0;
loop {
attempt_no += 1;
let attempt_start = Instant::now();
let outcome: Result<OpResult, ExecutionError> = {
use futures::FutureExt as _;
match std::panic::AssertUnwindSafe(self.inner.execute(cycle, ctx))
.catch_unwind()
.await
{
Ok(r) => r,
Err(payload) => {
let msg = payload
.downcast_ref::<&'static str>()
.map(|s| (*s).to_string())
.or_else(|| payload.downcast_ref::<String>().cloned())
.unwrap_or_else(|| "<non-string panic payload>".into());
Err(ExecutionError::Op(AdapterError {
error_name: "panic".into(),
message: msg,
retryable: false,
}))
}
}
};
let dt = attempt_start.elapsed().as_nanos() as u64;
match outcome {
Ok(result) => {
self.metrics.attempt_total.inc();
self.metrics.attempt_success.observe(dt);
self.metrics.tries_histogram.record(attempt_no as u64);
return Ok(result);
}
Err(e) => {
self.metrics.attempt_total.inc();
self.metrics.attempt_failure.observe(dt);
let retryable = matches!(&e, ExecutionError::Op(ad) if ad.retryable);
if retryable && attempt_no < self.tries {
if self.stop.stopped() {
self.metrics.tries_histogram.record(attempt_no as u64);
return Err(e);
}
if let Some(gate) = &self.advisory
&& let ExecutionError::Op(ad) = &e
&& gate.first_sighting(&ad.error_name)
{
use crate::exec_events::ExecEventSubscriber as _;
self.submit_advisory(
&self.op_name,
cycle,
self.tries,
&ad.error_name,
&ad.message,
);
}
if self.exemplars.enabled()
&& let ExecutionError::Op(ad) = &e
&& let Some(squelched) = self.exemplars.admit(cycle, attempt_no)
{
use crate::exec_events::ExecEventSubscriber as _;
self.submit_exemplar(&crate::exec_events::ExecExemplar {
op_name: &self.op_name,
cycle,
attempt_no,
tries_budget: self.tries,
error_class: &ad.error_name,
message: &ad.message,
will_retry: true,
squelched_since_last: squelched,
});
}
let wait = backoff_wait_ms(
self.backoff_base_ms,
self.backoff_max_ms,
self.backoff_ratio,
attempt_no,
cycle,
);
if wait > 0 {
portable_sleep_ms(wait).await;
}
if self.stop.stopped() {
self.metrics.tries_histogram.record(attempt_no as u64);
return Err(e);
}
continue;
}
self.metrics.tries_histogram.record(attempt_no as u64);
return Err(e);
}
}
}
})
}
fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
Some(self.inner.as_ref())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::adapter::ResultBody;
use crate::fixture::{ExecCtx, ResolvedPulls};
use nmbrs_metrics::labels::Labels;
use std::sync::atomic::{AtomicU32, Ordering};
struct FlakyInner {
fail_first: u32,
calls: AtomicU32,
}
impl OpDispenser for FlakyInner {
fn execute<'a>(
&'a self,
_cycle: u64,
_ctx: &'a ExecCtx<'a>,
) -> std::pin::Pin<
Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
> {
Box::pin(async move {
let n = self.calls.fetch_add(1, Ordering::Relaxed) + 1;
if n <= self.fail_first {
Err(ExecutionError::Op(AdapterError {
error_name: "Timeout".into(),
message: "flaky".into(),
retryable: true,
}))
} else {
Ok(OpResult {
body: None::<Box<dyn ResultBody>>,
skipped: false,
})
}
})
}
}
fn empty_ctx() -> (crate::adapter::ResolvedFields, ResolvedPulls) {
let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
let pulls = ResolvedPulls::empty();
(fields, pulls)
}
#[test]
fn backoff_is_geometric_capped_and_jittered() {
let caps = [100u64, 200, 400, 800, 1600, 3200, 6400, 10_000, 10_000];
for (i, &cap) in caps.iter().enumerate() {
let attempt = (i + 1) as u32;
let w = backoff_wait_ms(100, 10_000, 2.0, attempt, 42);
assert!(
w >= cap / 2 && w <= cap,
"attempt {attempt}: wait {w} out of [{},{cap}]",
cap / 2
);
}
for attempt in 1..=5u32 {
let w = backoff_wait_ms(200, 10_000, 1.0, attempt, 7);
assert!(w >= 100 && w <= 200, "constant backoff drifted: {w}");
}
assert_eq!(backoff_wait_ms(0, 10_000, 2.0, 3, 1), 0);
assert_eq!(
backoff_wait_ms(100, 10_000, 2.0, 4, 99),
backoff_wait_ms(100, 10_000, 2.0, 4, 99)
);
let w = backoff_wait_ms(100, 10_000, 0.1, 5, 3);
assert!(
w >= 50 && w <= 100,
"sub-1.0 ratio should hold at floor: {w}"
);
}
#[tokio::test]
async fn tries_zero_fails_without_executing() {
let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::clear_session_stop_for_test();
let inner = Arc::new(FlakyInner {
fail_first: 0,
calls: AtomicU32::new(0),
});
let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
let d = TriesDispenser::wrap(
inner.clone(),
0,
metrics,
0,
0,
2.0,
crate::session_signals::StopView::default(),
"test_op".to_string(),
crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
None,
);
let (fields, pulls) = empty_ctx();
let ctx = ExecCtx::new(&fields, &pulls);
let err = d.execute(0, &ctx).await.expect_err("tries:0 must fail");
assert_eq!(err.error().error_name, "tries_zero");
assert_eq!(
inner.calls.load(Ordering::Relaxed),
0,
"inner must never be invoked at tries:0"
);
}
#[tokio::test]
async fn shutdown_stops_retries_immediately() {
let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::clear_session_stop_for_test();
let inner = Arc::new(FlakyInner {
fail_first: 50,
calls: AtomicU32::new(0),
});
let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
let d = TriesDispenser::wrap(
inner.clone(),
100,
metrics,
0,
0,
2.0,
crate::session_signals::StopView::default(),
"test_op".to_string(),
crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
None,
);
let (fields, pulls) = empty_ctx();
let ctx = ExecCtx::new(&fields, &pulls);
crate::session_signals::request_stop();
let res = d.execute(0, &ctx).await;
crate::session_signals::clear_session_stop_for_test();
res.expect_err("failure under shutdown must be terminal");
assert_eq!(
inner.calls.load(Ordering::Relaxed),
1,
"no fresh attempts once the session is stopping"
);
}
#[tokio::test]
async fn tries_is_a_total_attempt_budget() {
let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::clear_session_stop_for_test();
let inner = Arc::new(FlakyInner {
fail_first: 2,
calls: AtomicU32::new(0),
});
let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
let d = TriesDispenser::wrap(
inner.clone(),
3,
metrics.clone(),
0,
0,
2.0,
crate::session_signals::StopView::default(),
"test_op".to_string(),
crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
None,
);
let (fields, pulls) = empty_ctx();
let ctx = ExecCtx::new(&fields, &pulls);
d.execute(0, &ctx).await.expect("third attempt succeeds");
assert_eq!(inner.calls.load(Ordering::Relaxed), 3);
assert_eq!(metrics.attempt_total.get(), 3);
}
#[tokio::test]
async fn budget_exhaustion_is_terminal() {
let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
crate::session_signals::clear_session_stop_for_test();
let inner = Arc::new(FlakyInner {
fail_first: 5,
calls: AtomicU32::new(0),
});
let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
let d = TriesDispenser::wrap(
inner.clone(),
2,
metrics,
0,
0,
2.0,
crate::session_signals::StopView::default(),
"test_op".to_string(),
crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
None,
);
let (fields, pulls) = empty_ctx();
let ctx = ExecCtx::new(&fields, &pulls);
let err = d.execute(0, &ctx).await.expect_err("budget spent");
assert_eq!(err.error().error_name, "Timeout");
assert_eq!(
inner.calls.load(Ordering::Relaxed),
2,
"tries:2 = exactly two total attempts"
);
}
}