nmbrs_runtime/wrappers/tries.rs
1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Tries wrapper — the CONDITIONAL innermost op-level wrapper (SRD-82 Part
5//! 3b). `tries:` is its sigil: the TOTAL number of attempts an op may make.
6//!
7//! - **No `tries` in scope** → the wrapper is not constructed; the op runs
8//! single-attempt (the outermost error-handler wrapper records the
9//! single-attempt tallies).
10//! - **`tries: 1`** → identical to no wrapper (single attempt); the cascade
11//! skips construction.
12//! - **`tries: 0`** → the op FAILS WITHOUT EXECUTING: every cycle yields a
13//! synthesised `tries_zero` op error, routed through the op's error
14//! policy like any terminal failure. The explicit "never run this" knob.
15//! - **`tries: N ≥ 2`** → up to `N` total attempts; a retryable attempt
16//! failure re-runs the inner op until the budget is spent.
17//!
18//! When constructed it owns the whole ATTEMPT→RESULT boundary: it runs the
19//! inner op (adapter) one-or-more times, owns the `attempt_*` counters,
20//! catches a per-attempt panic, and returns exactly ONE terminal outcome to
21//! the layers above. Everything above it — traversal, result-binding,
22//! metrics, the outermost error-handler wrapper — sees a single result,
23//! never the retries.
24//!
25//! Retryability is the ADAPTER's signal: an `ExecutionError::Op` whose
26//! `retryable` flag is set (CQL timeouts/overloads are). The `errors:`
27//! policy is deliberately NOT consulted here — the two surfaces are
28//! orthogonal; the policy's `retry` verb participates only by INJECTING a
29//! `tries` budget at dispenser build (the SRD-82 Part 3b bridge), never by
30//! steering the loop per-cycle.
31
32use std::sync::Arc;
33use std::time::Instant;
34
35use crate::activity::ActivityMetrics;
36use crate::adapter::{AdapterError, ExecutionError, OpDispenser, OpResult, WrappingDispenser};
37use crate::wrapper_registry::{WrapperName, WrapperRegistration, WrapperSubject};
38
39pub const NAME: WrapperName = WrapperName::new("tries");
40
41/// Registry trigger: the op's own `tries:` field. The full activation set is
42/// wider — an in-scope `tries` wire, the inherited phase/root `tries`, or an
43/// `errors:` retry-verb injection — but those resolve against the kernel /
44/// config at dispenser build (the cascade), which a `ParsedOp` predicate
45/// cannot see. The registry entry drives field validation + telemetry; the
46/// hand-placed innermost construction is authoritative.
47fn triggers(s: WrapperSubject) -> bool {
48 let Some(op) = s.op() else {
49 return false;
50 };
51 op.params.contains_key("tries")
52}
53
54fn describe_assignment(s: WrapperSubject) -> Option<String> {
55 let op = s.op()?;
56 op.params.get("tries").map(|v| format!("tries: {v}"))
57}
58
59inventory::submit! {
60 WrapperRegistration {
61 name: NAME,
62 // The wrapper owns its sigil plus the standalone companion
63 // knobs consumed at wrapper build: retry pacing and the
64 // retry-error exemplar sampler (rate fraction, default 0.0
65 // = off; max_hz emission ceiling).
66 owned_fields: &["tries", "retry_backoff", "retry_backoff_max",
67 "retry_backoff_ratio", "retry_exemplar_rate",
68 "retry_exemplar_max_hz", "retry_advisory"],
69 triggers,
70 requires_inner: &[],
71 forbids_outer: &[],
72 mutually_exclusive_with: &[],
73 describe_assignment,
74 levels: &[crate::wrapper_registry::WrapperLevel::Op],
75 }
76}
77
78/// Wraps an inner dispenser with a bounded attempt loop over retryable
79/// attempt failures, owning the `attempt_*` metrics and the per-attempt
80/// panic catch. Constructed only when a `tries` budget of `0` or `≥ 2`
81/// resolves for the op (see the module doc; `1` / absent skip construction).
82pub struct TriesDispenser {
83 inner: Arc<dyn OpDispenser>,
84 /// TOTAL attempts the op may make. `0` = fail without executing.
85 /// (`1` never reaches construction — the cascade skips the wrapper.)
86 tries: u32,
87 /// Activity-level metrics — the attempt tallies live here so the whole
88 /// phase shares one attempt counter regardless of which op path ran.
89 metrics: Arc<ActivityMetrics>,
90 /// Backoff pacing between retryable attempts (compaction-demo
91 /// diagnosis: an immediate-`continue` retry loop hammers a dying
92 /// server, and `tries: 20 × timeout: 60s` makes it look like a
93 /// silent stall). Attempt `k`'s wait is `backoff_base_ms *
94 /// backoff_ratio^(k-1)`, capped at `backoff_max_ms`, with
95 /// deterministic jitter in [50%, 100%] derived from (cycle, attempt)
96 /// so runs stay replayable. `base == 0` disables pacing. Set per op
97 /// via the `tries:` map form (`backoff: {ratio, min, max}`) or the
98 /// standalone `retry_backoff` / `retry_backoff_max` /
99 /// `retry_backoff_ratio` params (defaults 100ms / 10s / 2.0).
100 backoff_base_ms: u64,
101 backoff_max_ms: u64,
102 /// Geometric growth factor per retry. `2.0` doubles each wait; `1.0`
103 /// holds it constant at the floor. Clamped to `>= 1.0` at use so a
104 /// misconfigured ratio never shrinks the backoff.
105 backoff_ratio: f64,
106 /// SRD-92 cooperative-stop view (the same aspect handed to the
107 /// `while:` wrapper): once any layered stop source fires — session
108 /// Ctrl-C, activity fault, walk halt, daemon-group completion —
109 /// a retryable failure returns terminally instead of starting
110 /// fresh attempts through the drain window. Injected at wrap time
111 /// so the wrapper never reaches for globals itself.
112 stop: crate::session_signals::StopView,
113 /// Op template name — identifies the specimen in exemplar lines.
114 op_name: String,
115 /// Counter-exemplar sampling for errors caught in the retry loop
116 /// (`retry_exemplar_rate` / `retry_exemplar_max_hz` params;
117 /// default rate 0.0 = off). The error policy never sees a
118 /// retried-then-recovered attempt, so without sampling those
119 /// messages are visible only as `attempt_failure` counts.
120 exemplars: crate::exec_events::ExemplarSampler,
121 /// Default-on per-phase retry advisory (`exec_events`): the first
122 /// sighting of each error class in the retry loop emits one
123 /// advisory line through the shared per-activity gate — a retry
124 /// storm identifies itself without flooding the output. `None`
125 /// when the op opted out (`retry_advisory: off`).
126 advisory: Option<Arc<crate::exec_events::AdvisoryGate>>,
127}
128
129impl crate::exec_events::ExecEventSubscriber for TriesDispenser {}
130
131impl TriesDispenser {
132 /// Wrap `inner` with a total-attempts budget, retry pacing, and
133 /// optional retry-error exemplar sampling.
134 #[allow(clippy::too_many_arguments)]
135 pub fn wrap(
136 inner: Arc<dyn OpDispenser>,
137 tries: u32,
138 metrics: Arc<ActivityMetrics>,
139 backoff_base_ms: u64,
140 backoff_max_ms: u64,
141 backoff_ratio: f64,
142 stop: crate::session_signals::StopView,
143 op_name: String,
144 exemplars: crate::exec_events::ExemplarSampler,
145 advisory: Option<Arc<crate::exec_events::AdvisoryGate>>,
146 ) -> Arc<dyn OpDispenser> {
147 Arc::new(Self {
148 inner,
149 tries,
150 metrics,
151 backoff_base_ms,
152 backoff_max_ms,
153 backoff_ratio,
154 stop,
155 op_name,
156 exemplars,
157 advisory,
158 })
159 }
160}
161
162/// Runtime-agnostic async sleep. `tokio::time::sleep` binds its
163/// timer to the runtime owning the current thread — fibers run on
164/// shared pool threads, so under multiple runtimes (the lib-test
165/// harness; any embedder) the timer can land on a runtime that
166/// shuts down mid-sleep. A detached sleeper thread + oneshot has
167/// no such coupling, and backoff is the degraded path — a
168/// short-lived thread is noise next to the ≥50ms wait it serves.
169async fn portable_sleep_ms(ms: u64) {
170 let (tx, rx) = futures::channel::oneshot::channel::<()>();
171 std::thread::spawn(move || {
172 std::thread::sleep(std::time::Duration::from_millis(ms));
173 let _ = tx.send(());
174 });
175 let _ = rx.await;
176}
177
178use crate::exec_events::splitmix64;
179
180/// The jittered wait (ms) before retry `attempt_no` (1-based). Geometric
181/// growth `base * ratio^(attempt-1)`, capped at `max`, then deterministic
182/// jitter into `[50%, 100%]` of the capped value keyed on `(cycle,
183/// attempt)` so a replay reproduces the same schedule. `base == 0` disables
184/// pacing (returns 0). `ratio` is clamped to `>= 1.0` so a misconfigured
185/// value never shrinks the wait; a huge exponent saturates `powi` to +inf,
186/// which the `.min(max)` folds to the cap (no integer overflow).
187fn backoff_wait_ms(base_ms: u64, max_ms: u64, ratio: f64, attempt_no: u32, cycle: u64) -> u64 {
188 if base_ms == 0 {
189 return 0;
190 }
191 let mult = ratio.max(1.0).powi((attempt_no.saturating_sub(1)) as i32);
192 let raw = base_ms as f64 * mult;
193 let capped = raw.min(max_ms as f64).max(1.0) as u64;
194 let h = splitmix64(cycle ^ ((attempt_no as u64) << 48));
195 capped / 2 + (h % (capped / 2 + 1))
196}
197
198impl WrappingDispenser for TriesDispenser {}
199
200impl OpDispenser for TriesDispenser {
201 fn execute<'a>(
202 &'a self,
203 cycle: u64,
204 ctx: &'a crate::fixture::ExecCtx<'a>,
205 ) -> std::pin::Pin<
206 Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
207 > {
208 Box::pin(async move {
209 // `tries: 0` — the op is configured to fail without executing.
210 // Accounted as one failed (zero-length) attempt so the att:%
211 // display stays truthful, and routed through the error policy by
212 // the outermost error-handler wrapper like any terminal failure.
213 if self.tries == 0 {
214 self.metrics.attempt_total.inc();
215 self.metrics.attempt_failure.observe(0);
216 self.metrics.tries_histogram.record(0);
217 return Err(ExecutionError::Op(AdapterError {
218 error_name: "tries_zero".into(),
219 message: "tries: 0 — op is configured to fail without executing".into(),
220 retryable: false,
221 }));
222 }
223 // `attempt_no` counts attempts (1-based); total attempts run is
224 // at most `tries`. It is also the value recorded into
225 // `tries_histogram` per op.
226 let mut attempt_no: u32 = 0;
227 loop {
228 attempt_no += 1;
229 let attempt_start = Instant::now();
230
231 // Per-attempt panic catch: an adapter that unwinds out of
232 // `execute` (an FFI driver on a broken connection) becomes a
233 // synthesised `panic` op error rather than killing the fiber.
234 // It is NOT adapter-retryable, so it counts as a failed attempt
235 // and (absent a retryable classification) terminates without
236 // spinning — the fiber survives regardless.
237 let outcome: Result<OpResult, ExecutionError> = {
238 use futures::FutureExt as _;
239 match std::panic::AssertUnwindSafe(self.inner.execute(cycle, ctx))
240 .catch_unwind()
241 .await
242 {
243 Ok(r) => r,
244 Err(payload) => {
245 let msg = payload
246 .downcast_ref::<&'static str>()
247 .map(|s| (*s).to_string())
248 .or_else(|| payload.downcast_ref::<String>().cloned())
249 .unwrap_or_else(|| "<non-string panic payload>".into());
250 Err(ExecutionError::Op(AdapterError {
251 error_name: "panic".into(),
252 message: msg,
253 retryable: false,
254 }))
255 }
256 }
257 };
258 let dt = attempt_start.elapsed().as_nanos() as u64;
259
260 match outcome {
261 // Attempt instruments count at RESOLUTION, same
262 // discipline as the result instruments: an attempt
263 // exists in the tallies only once it has returned,
264 // so `attempt_total == attempt_success +
265 // attempt_failure` holds at every read and a rate
266 // over `attempt_total` carries no in-flight skew.
267 Ok(result) => {
268 self.metrics.attempt_total.inc();
269 self.metrics.attempt_success.observe(dt);
270 self.metrics.tries_histogram.record(attempt_no as u64);
271 return Ok(result);
272 }
273 Err(e) => {
274 self.metrics.attempt_total.inc();
275 self.metrics.attempt_failure.observe(dt);
276 // Retry only an adapter-retryable OP error, within
277 // the total-attempts budget. Adapter-level errors are
278 // never retried here (they are connection-level, not
279 // per-op).
280 let retryable = matches!(&e, ExecutionError::Op(ad) if ad.retryable);
281 if retryable && attempt_no < self.tries {
282 // Session shutdown: the cooperative drain
283 // waits for the CURRENT attempt only — a
284 // fresh retry is new work, and against a
285 // struggling server it spins through the
286 // whole grace window. Checked again after
287 // the backoff so a Ctrl-C during the wait
288 // also lands.
289 if self.stop.stopped() {
290 self.metrics.tries_histogram.record(attempt_no as u64);
291 return Err(e);
292 }
293 // Default-on advisory: the FIRST sighting of
294 // each error class in this phase announces
295 // itself once — the retry loop is otherwise
296 // silent about what it absorbs.
297 if let Some(gate) = &self.advisory
298 && let ExecutionError::Op(ad) = &e
299 && gate.first_sighting(&ad.error_name)
300 {
301 use crate::exec_events::ExecEventSubscriber as _;
302 self.submit_advisory(
303 &self.op_name,
304 cycle,
305 self.tries,
306 &ad.error_name,
307 &ad.message,
308 );
309 }
310 // Counter-exemplar sampling: this error will be
311 // retried, so the error policy never sees it —
312 // a sampled specimen goes to the structured
313 // sink instead (see `exec_events`).
314 if self.exemplars.enabled()
315 && let ExecutionError::Op(ad) = &e
316 && let Some(squelched) = self.exemplars.admit(cycle, attempt_no)
317 {
318 use crate::exec_events::ExecEventSubscriber as _;
319 self.submit_exemplar(&crate::exec_events::ExecExemplar {
320 op_name: &self.op_name,
321 cycle,
322 attempt_no,
323 tries_budget: self.tries,
324 error_class: &ad.error_name,
325 message: &ad.message,
326 will_retry: true,
327 squelched_since_last: squelched,
328 });
329 }
330 let wait = backoff_wait_ms(
331 self.backoff_base_ms,
332 self.backoff_max_ms,
333 self.backoff_ratio,
334 attempt_no,
335 cycle,
336 );
337 if wait > 0 {
338 portable_sleep_ms(wait).await;
339 }
340 if self.stop.stopped() {
341 self.metrics.tries_histogram.record(attempt_no as u64);
342 return Err(e);
343 }
344 continue;
345 }
346 // Terminal: hand the failure up to the result level.
347 self.metrics.tries_histogram.record(attempt_no as u64);
348 return Err(e);
349 }
350 }
351 }
352 })
353 }
354
355 fn inner_dispenser(&self) -> Option<&dyn OpDispenser> {
356 Some(self.inner.as_ref())
357 }
358}
359
360#[cfg(test)]
361mod tests {
362 use super::*;
363 use crate::adapter::ResultBody;
364 use crate::fixture::{ExecCtx, ResolvedPulls};
365 use nmbrs_metrics::labels::Labels;
366 use std::sync::atomic::{AtomicU32, Ordering};
367
368 /// Inner stub: fails with a retryable error until `fail_first` attempts
369 /// have been consumed, then succeeds. Counts invocations.
370 struct FlakyInner {
371 fail_first: u32,
372 calls: AtomicU32,
373 }
374
375 impl OpDispenser for FlakyInner {
376 fn execute<'a>(
377 &'a self,
378 _cycle: u64,
379 _ctx: &'a ExecCtx<'a>,
380 ) -> std::pin::Pin<
381 Box<dyn std::future::Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>,
382 > {
383 Box::pin(async move {
384 let n = self.calls.fetch_add(1, Ordering::Relaxed) + 1;
385 if n <= self.fail_first {
386 Err(ExecutionError::Op(AdapterError {
387 error_name: "Timeout".into(),
388 message: "flaky".into(),
389 retryable: true,
390 }))
391 } else {
392 Ok(OpResult {
393 body: None::<Box<dyn ResultBody>>,
394 skipped: false,
395 })
396 }
397 })
398 }
399 }
400
401 fn empty_ctx() -> (crate::adapter::ResolvedFields, ResolvedPulls) {
402 let fields = crate::adapter::ResolvedFields::new(vec![], vec![]);
403 let pulls = ResolvedPulls::empty();
404 (fields, pulls)
405 }
406
407 /// The backoff schedule grows geometrically by `ratio`, stays within
408 /// `[50%, 100%]` of the capped value, honours the `max` cap, holds
409 /// constant at `ratio == 1.0`, and disables at `base == 0`.
410 #[test]
411 fn backoff_is_geometric_capped_and_jittered() {
412 // ratio 2.0, base 100, max 10s: the CAP for attempt k is
413 // 100 * 2^(k-1) until it hits 10_000, and the wait lands in
414 // [cap/2, cap].
415 let caps = [100u64, 200, 400, 800, 1600, 3200, 6400, 10_000, 10_000];
416 for (i, &cap) in caps.iter().enumerate() {
417 let attempt = (i + 1) as u32;
418 let w = backoff_wait_ms(100, 10_000, 2.0, attempt, 42);
419 assert!(
420 w >= cap / 2 && w <= cap,
421 "attempt {attempt}: wait {w} out of [{},{cap}]",
422 cap / 2
423 );
424 }
425 // ratio 1.0 holds the wait at the floor every attempt.
426 for attempt in 1..=5u32 {
427 let w = backoff_wait_ms(200, 10_000, 1.0, attempt, 7);
428 assert!(w >= 100 && w <= 200, "constant backoff drifted: {w}");
429 }
430 // base 0 disables pacing entirely.
431 assert_eq!(backoff_wait_ms(0, 10_000, 2.0, 3, 1), 0);
432 // Deterministic: same (cycle, attempt) → same wait (replayable).
433 assert_eq!(
434 backoff_wait_ms(100, 10_000, 2.0, 4, 99),
435 backoff_wait_ms(100, 10_000, 2.0, 4, 99)
436 );
437 // A misconfigured ratio < 1.0 is clamped, never shrinks the wait.
438 let w = backoff_wait_ms(100, 10_000, 0.1, 5, 3);
439 assert!(
440 w >= 50 && w <= 100,
441 "sub-1.0 ratio should hold at floor: {w}"
442 );
443 }
444
445 /// `tries: 0` fails WITHOUT invoking the inner op.
446 #[tokio::test]
447 async fn tries_zero_fails_without_executing() {
448 let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
449 .lock()
450 .unwrap_or_else(|e| e.into_inner());
451 crate::session_signals::clear_session_stop_for_test();
452 let inner = Arc::new(FlakyInner {
453 fail_first: 0,
454 calls: AtomicU32::new(0),
455 });
456 let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
457 let d = TriesDispenser::wrap(
458 inner.clone(),
459 0,
460 metrics,
461 0,
462 0,
463 2.0,
464 crate::session_signals::StopView::default(),
465 "test_op".to_string(),
466 crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
467 None,
468 );
469 let (fields, pulls) = empty_ctx();
470 let ctx = ExecCtx::new(&fields, &pulls);
471 let err = d.execute(0, &ctx).await.expect_err("tries:0 must fail");
472 assert_eq!(err.error().error_name, "tries_zero");
473 assert_eq!(
474 inner.calls.load(Ordering::Relaxed),
475 0,
476 "inner must never be invoked at tries:0"
477 );
478 }
479
480 /// Session shutdown ends the retry loop: a retryable failure with
481 /// budget remaining returns terminally instead of spinning through
482 /// the cooperative-drain window on a struggling server.
483 #[tokio::test]
484 async fn shutdown_stops_retries_immediately() {
485 let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
486 .lock()
487 .unwrap_or_else(|e| e.into_inner());
488 crate::session_signals::clear_session_stop_for_test();
489 // Would fail retryably 50 times — but stop is raised, so the
490 // very first failure must be terminal.
491 let inner = Arc::new(FlakyInner {
492 fail_first: 50,
493 calls: AtomicU32::new(0),
494 });
495 let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
496 let d = TriesDispenser::wrap(
497 inner.clone(),
498 100,
499 metrics,
500 0,
501 0,
502 2.0,
503 crate::session_signals::StopView::default(),
504 "test_op".to_string(),
505 crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
506 None,
507 );
508 let (fields, pulls) = empty_ctx();
509 let ctx = ExecCtx::new(&fields, &pulls);
510 crate::session_signals::request_stop();
511 let res = d.execute(0, &ctx).await;
512 crate::session_signals::clear_session_stop_for_test();
513 res.expect_err("failure under shutdown must be terminal");
514 assert_eq!(
515 inner.calls.load(Ordering::Relaxed),
516 1,
517 "no fresh attempts once the session is stopping"
518 );
519 }
520
521 /// `tries: 3` retries a retryable failure up to 3 TOTAL attempts and
522 /// succeeds when the third works.
523 #[tokio::test]
524 async fn tries_is_a_total_attempt_budget() {
525 let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
526 .lock()
527 .unwrap_or_else(|e| e.into_inner());
528 crate::session_signals::clear_session_stop_for_test();
529 let inner = Arc::new(FlakyInner {
530 fail_first: 2,
531 calls: AtomicU32::new(0),
532 });
533 let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
534 let d = TriesDispenser::wrap(
535 inner.clone(),
536 3,
537 metrics.clone(),
538 0,
539 0,
540 2.0,
541 crate::session_signals::StopView::default(),
542 "test_op".to_string(),
543 crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
544 None,
545 );
546 let (fields, pulls) = empty_ctx();
547 let ctx = ExecCtx::new(&fields, &pulls);
548 d.execute(0, &ctx).await.expect("third attempt succeeds");
549 assert_eq!(inner.calls.load(Ordering::Relaxed), 3);
550 assert_eq!(metrics.attempt_total.get(), 3);
551 }
552
553 /// The budget is TOTAL attempts: `tries: 2` against an op that needs a
554 /// third attempt fails after exactly 2 invocations.
555 #[tokio::test]
556 async fn budget_exhaustion_is_terminal() {
557 let _guard = crate::session_signals::STOP_GLOBAL_TEST_LOCK
558 .lock()
559 .unwrap_or_else(|e| e.into_inner());
560 crate::session_signals::clear_session_stop_for_test();
561 let inner = Arc::new(FlakyInner {
562 fail_first: 5,
563 calls: AtomicU32::new(0),
564 });
565 let metrics = Arc::new(ActivityMetrics::new(&Labels::empty()));
566 let d = TriesDispenser::wrap(
567 inner.clone(),
568 2,
569 metrics,
570 0,
571 0,
572 2.0,
573 crate::session_signals::StopView::default(),
574 "test_op".to_string(),
575 crate::exec_events::ExemplarSampler::pinned(0.0, 0.0),
576 None,
577 );
578 let (fields, pulls) = empty_ctx();
579 let ctx = ExecCtx::new(&fields, &pulls);
580 let err = d.execute(0, &ctx).await.expect_err("budget spent");
581 assert_eq!(err.error().error_name, "Timeout");
582 assert_eq!(
583 inner.calls.load(Ordering::Relaxed),
584 2,
585 "tries:2 = exactly two total attempts"
586 );
587 }
588}