Skip to main content

nmbrs_runtime/
exec_events.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Structured exec-event exemplars — the wrapper-facing tap onto the
5//! system's structured event sink (the observer surface, SRD-81/-88).
6//!
7//! Wrappers OPT IN by implementing [`ExecEventSubscriber`], a decorator
8//! service in the established wrapper style (a dyn-safe trait whose
9//! default methods ARE the service — the same shape as
10//! `WrappingDispenser`): implementing it grants the wrapper the
11//! canonical submission surface, and nothing else changes about the
12//! wrapper's construction or registration. The default routing is the
13//! single chokepoint [`submit_exemplar`]: one rendered projection per
14//! event through [`crate::observer::log_tagged`], which fans out
15//! to every installed sink — the durable `session.log` always, the
16//! live display per its level gates. No wrapper hand-rolls its own
17//! event formatting or reaches for a sink directly.
18//!
19//! # Exemplars, not streams
20//!
21//! An [`ExecExemplar`] is a SAMPLED counter-exemplar: a concrete
22//! specimen of an error class that is otherwise visible only as a
23//! counter (e.g. `attempt_failure` inside the retry loop, where the
24//! error policy never sees the message because the attempt recovers).
25//! Sampling is the submitting wrapper's job via [`ExemplarSampler`]:
26//! a fraction (`rate`, default 0.0 = off) decides which caught errors
27//! become exemplars, and a frequency ceiling (`max_hz`) squelches
28//! bursts. Squelched admissions are COUNTED, never dropped silently:
29//! the next emitted exemplar carries `(+N squelched)`, and any
30//! leftover tally is flushed at Debug when the sampler drops.
31
32use std::sync::atomic::{AtomicU64, Ordering};
33use std::time::Instant;
34
35/// One sampled error specimen from an execution wrapper.
36pub struct ExecExemplar<'a> {
37    /// Op template name the error occurred under.
38    pub op_name: &'a str,
39    /// Cycle whose attempt produced the error.
40    pub cycle: u64,
41    /// 1-based attempt number that failed.
42    pub attempt_no: u32,
43    /// The op's total-attempts budget.
44    pub tries_budget: u32,
45    /// Adapter error class (`error_name`).
46    pub error_class: &'a str,
47    /// Full adapter error message.
48    pub message: &'a str,
49    /// True when the failed attempt will be retried (the class of
50    /// error that is otherwise invisible outside counters).
51    pub will_retry: bool,
52    /// Admissions squelched by the frequency ceiling since the last
53    /// emitted exemplar — carried on this line so the squelch is
54    /// visible, never silent.
55    pub squelched_since_last: u64,
56}
57
58/// Render an exemplar to its one-line session projection. Pure —
59/// separated from [`submit_exemplar`] so the format is testable
60/// without an observer.
61pub fn render_exemplar(ex: &ExecExemplar<'_>) -> String {
62    let retry_note = if ex.will_retry {
63        "retrying"
64    } else {
65        "terminal"
66    };
67    let squelch_note = if ex.squelched_since_last > 0 {
68        format!(" (+{} squelched)", ex.squelched_since_last)
69    } else {
70        String::new()
71    };
72    format!(
73        "retry exemplar: op '{}' attempt {}/{} cycle {} ({retry_note}): \
74         [{}] {}{squelch_note}",
75        ex.op_name, ex.attempt_no, ex.tries_budget, ex.cycle, ex.error_class, ex.message,
76    )
77}
78
79/// The canonical submission chokepoint: one projection through the
80/// observer's categorized log surface at Warn (an exemplar IS an
81/// error specimen the operator asked to see).
82pub fn submit_exemplar(ex: &ExecExemplar<'_>) {
83    crate::observer::log_tagged(
84        crate::observer::LogLevel::Warn,
85        crate::observer::EventTag::in_flight(crate::observer::EventCategory::Retry),
86        &render_exemplar(ex),
87    );
88}
89
90/// Decorator service: a wrapper subscribes to the structured event
91/// sink by implementing this trait (dyn-safe; default methods are
92/// the whole service). Override nothing to get the canonical
93/// routing; the trait exists so the subscription is a declared,
94/// greppable property of the wrapper type rather than an ad-hoc
95/// call into logging.
96pub trait ExecEventSubscriber {
97    /// Submit one sampled exemplar to the structured sink.
98    fn submit_exemplar(&self, ex: &ExecExemplar<'_>) {
99        submit_exemplar(ex)
100    }
101
102    /// Submit one first-sighting retry advisory (default-on signal;
103    /// the per-phase [`AdvisoryGate`] bounds it).
104    fn submit_advisory(
105        &self,
106        op_name: &str,
107        cycle: u64,
108        tries_budget: u32,
109        error_class: &str,
110        message: &str,
111    ) {
112        crate::observer::log_tagged(
113            crate::observer::LogLevel::Warn,
114            crate::observer::EventTag::in_flight(crate::observer::EventCategory::Retry),
115            &render_advisory(op_name, cycle, tries_budget, error_class, message),
116        );
117    }
118}
119
120/// Render a first-sighting retry advisory. Pure — testable without
121/// an observer.
122pub fn render_advisory(
123    op_name: &str,
124    cycle: u64,
125    tries_budget: u32,
126    error_class: &str,
127    message: &str,
128) -> String {
129    format!(
130        "retry advisory: op '{op_name}' hit its first retryable \
131         [{error_class}] at cycle {cycle}: {message} — further \
132         occurrences are absorbed by the tries budget ({tries_budget}) \
133         and appear only as att:%/r: chips and attempt_* metrics; \
134         sample live specimens via the retry_exemplar_rate control, \
135         or silence this line with retry_advisory: off"
136    )
137}
138
139/// Per-phase advisory gate: by DEFAULT (no exemplar sampling opted
140/// in) the operator still gets at least SOME signal when the retry
141/// loop starts absorbing errors — one advisory per error class per
142/// phase, capped, so a retry storm identifies itself without
143/// flooding the session output. Shared per activity (like
144/// [`ExemplarConfig`]) so many ops in one phase share the budget.
145pub struct AdvisoryGate {
146    seen: std::sync::Mutex<std::collections::HashSet<String>>,
147    /// Max distinct classes advised per phase; beyond it the gate
148    /// closes (the classes are countable in `errors_total` labels).
149    cap: usize,
150}
151
152impl Default for AdvisoryGate {
153    fn default() -> Self {
154        Self::new()
155    }
156}
157
158impl AdvisoryGate {
159    pub fn new() -> Self {
160        Self {
161            seen: std::sync::Mutex::new(std::collections::HashSet::new()),
162            cap: 3,
163        }
164    }
165
166    /// True exactly once per error class (under the cap) — the
167    /// caller emits the advisory for that sighting.
168    pub fn first_sighting(&self, class: &str) -> bool {
169        let mut seen = self.seen.lock().unwrap_or_else(|e| e.into_inner());
170        if seen.len() >= self.cap && !seen.contains(class) {
171            return false;
172        }
173        seen.insert(class.to_string())
174    }
175}
176
177/// splitmix64 — cheap deterministic hash shared by replayable
178/// sampling decisions (and the tries wrapper's backoff jitter): the
179/// same (cycle, attempt) always makes the same choice, so a replay
180/// reproduces the same exemplars.
181pub(crate) fn splitmix64(mut z: u64) -> u64 {
182    z = z.wrapping_add(0x9E37_79B9_7F4A_7C15);
183    z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
184    z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
185    z ^ (z >> 31)
186}
187
188/// The sampling configuration cell: rate and frequency ceiling as
189/// shared atomics, so a dynamic control can move every sampler
190/// reading the cell with ONE store — push-on-set, no polling, no
191/// per-op control traffic (the `cql_trace_rate` pattern). Readers
192/// pay one atomic load, and only on the retry path.
193///
194/// Scoping: the activity owns one shared cell for every tries
195/// wrapper that does not pin its own values; an op that declares
196/// `retry_exemplar_*` params gets a private cell the controls
197/// deliberately do not move (authored matter wins).
198pub struct ExemplarConfig {
199    /// f64 bits of the sampling fraction ∈ [0, 1]. 0.0 = off.
200    rate_bits: AtomicU64,
201    /// Minimum nanos between emissions (derived from `max_hz` at
202    /// SET time so the read path never divides). 0 = no ceiling.
203    min_interval_nanos: AtomicU64,
204}
205
206impl ExemplarConfig {
207    pub fn new(rate: f64, max_hz: f64) -> Self {
208        let cfg = Self {
209            rate_bits: AtomicU64::new(0),
210            min_interval_nanos: AtomicU64::new(0),
211        };
212        cfg.set_rate(rate);
213        cfg.set_max_hz(max_hz);
214        cfg
215    }
216
217    /// Publish a new sampling fraction (clamped to [0, 1];
218    /// non-finite = off). One atomic store — this IS the dynamic
219    /// control's applier body.
220    pub fn set_rate(&self, rate: f64) {
221        let rate = if rate.is_finite() {
222            rate.clamp(0.0, 1.0)
223        } else {
224            0.0
225        };
226        self.rate_bits.store(rate.to_bits(), Ordering::Release);
227    }
228
229    /// Publish a new frequency ceiling (events/sec; `0` or
230    /// non-finite = uncapped). Converted to an interval here so
231    /// admission never divides.
232    pub fn set_max_hz(&self, max_hz: f64) {
233        let interval = if max_hz.is_finite() && max_hz > 0.0 {
234            (1_000_000_000f64 / max_hz) as u64
235        } else {
236            0
237        };
238        self.min_interval_nanos.store(interval, Ordering::Release);
239    }
240
241    fn rate(&self) -> f64 {
242        f64::from_bits(self.rate_bits.load(Ordering::Acquire))
243    }
244
245    fn min_interval(&self) -> u64 {
246        self.min_interval_nanos.load(Ordering::Acquire)
247    }
248}
249
250/// Sampling + squelch gate for exemplar submission.
251///
252/// Two independent controls compose (read live from the
253/// [`ExemplarConfig`] cell, so a dynamic control moves them
254/// mid-run):
255/// - `rate` ∈ [0, 1] — the fraction of caught errors that become
256///   exemplar candidates. `0.0` (the default) disables the sampler
257///   entirely; the caller's hot path pays one atomic load.
258///   The roll is DETERMINISTIC on (cycle, attempt) so runs replay.
259/// - `max_hz` — ceiling on emitted exemplars per second. Candidates
260///   over the ceiling are squelched and COUNTED; the count drains
261///   onto the next admitted exemplar, and any leftover flushes at
262///   Debug on drop. `0` or non-finite = no ceiling.
263pub struct ExemplarSampler {
264    cfg: std::sync::Arc<ExemplarConfig>,
265    base: Instant,
266    /// Elapsed nanos (since `base`) of the last admitted exemplar,
267    /// +1 so that 0 means "never admitted".
268    last_admit: AtomicU64,
269    squelched: AtomicU64,
270}
271
272impl ExemplarSampler {
273    /// A sampler over its own private cell — the authored-pin form
274    /// (op-level `retry_exemplar_*` params); dynamic controls do
275    /// not move it.
276    pub fn pinned(rate: f64, max_hz: f64) -> Self {
277        Self::shared(std::sync::Arc::new(ExemplarConfig::new(rate, max_hz)))
278    }
279
280    /// A sampler over a shared cell — the default form; the cell's
281    /// owner (the activity) wires it to the `retry_exemplar_rate` /
282    /// `retry_exemplar_max_hz` dynamic controls.
283    pub fn shared(cfg: std::sync::Arc<ExemplarConfig>) -> Self {
284        Self {
285            cfg,
286            base: Instant::now(),
287            last_admit: AtomicU64::new(0),
288            squelched: AtomicU64::new(0),
289        }
290    }
291
292    /// True when any sampling can happen at all — one atomic load,
293    /// paid only on the retry path.
294    pub fn enabled(&self) -> bool {
295        self.cfg.rate() > 0.0
296    }
297
298    /// Decide one caught error. `None` = not sampled (failed the
299    /// roll, or over the frequency ceiling — the latter counted).
300    /// `Some(n)` = admitted, draining `n` squelched admissions to
301    /// report on this exemplar's line.
302    pub fn admit(&self, cycle: u64, attempt_no: u32) -> Option<u64> {
303        let rate = self.cfg.rate();
304        if rate <= 0.0 {
305            return None;
306        }
307        // Deterministic roll on (cycle, attempt) — replayable, and
308        // uniform enough for a sampling fraction.
309        let h = splitmix64(cycle ^ ((attempt_no as u64) << 48) ^ 0xE0E0_5EED);
310        if (h as f64 / u64::MAX as f64) >= rate {
311            return None;
312        }
313        let min_interval = self.cfg.min_interval();
314        if min_interval == 0 {
315            // Uncapped — but still stamp the gate so a ceiling
316            // applied LIVE measures from real emission history
317            // rather than treating the next admission as first.
318            let now = self.base.elapsed().as_nanos() as u64 + 1;
319            self.last_admit.store(now, Ordering::Release);
320            return Some(self.squelched.swap(0, Ordering::AcqRel));
321        }
322        let now = self.base.elapsed().as_nanos() as u64 + 1;
323        loop {
324            let last = self.last_admit.load(Ordering::Acquire);
325            if last != 0 && now.saturating_sub(last) < min_interval {
326                self.squelched.fetch_add(1, Ordering::AcqRel);
327                return None;
328            }
329            if self
330                .last_admit
331                .compare_exchange(last, now, Ordering::AcqRel, Ordering::Acquire)
332                .is_ok()
333            {
334                return Some(self.squelched.swap(0, Ordering::AcqRel));
335            }
336        }
337    }
338}
339
340impl Drop for ExemplarSampler {
341    fn drop(&mut self) {
342        // Leftover squelch tally: surfaced, never silently lost.
343        let leftover = self.squelched.load(Ordering::Acquire);
344        if leftover > 0 {
345            crate::diag!(
346                crate::observer::LogLevel::Debug,
347                "exemplar sampler retired with {leftover} squelched \
348                 admission(s) unreported (frequency ceiling)"
349            );
350        }
351    }
352}
353
354#[cfg(test)]
355mod tests {
356    use super::*;
357
358    /// `rate: 0.0` (the default) admits nothing and stays cheap.
359    #[test]
360    fn zero_rate_is_off() {
361        let s = ExemplarSampler::pinned(0.0, 1000.0);
362        assert!(!s.enabled());
363        for c in 0..1000 {
364            assert!(s.admit(c, 1).is_none());
365        }
366    }
367
368    /// `rate: 1.0` with no ceiling admits every caught error.
369    #[test]
370    fn full_rate_uncapped_admits_all() {
371        let s = ExemplarSampler::pinned(1.0, 0.0);
372        for c in 0..100 {
373            assert_eq!(s.admit(c, 1), Some(0), "cycle {c}");
374        }
375    }
376
377    /// A fractional rate admits roughly its share, deterministically:
378    /// the same (cycle, attempt) keys always make the same choice.
379    #[test]
380    fn fractional_rate_samples_deterministically() {
381        let s1 = ExemplarSampler::pinned(0.25, 0.0);
382        let s2 = ExemplarSampler::pinned(0.25, 0.0);
383        let picks1: Vec<bool> = (0..4000).map(|c| s1.admit(c, 3).is_some()).collect();
384        let picks2: Vec<bool> = (0..4000).map(|c| s2.admit(c, 3).is_some()).collect();
385        assert_eq!(picks1, picks2, "sampling must be replayable");
386        let hits = picks1.iter().filter(|b| **b).count();
387        assert!(
388            (600..=1400).contains(&hits),
389            "0.25 of 4000 should land near 1000, got {hits}"
390        );
391    }
392
393    /// The frequency ceiling squelches bursts, counts what it
394    /// squelched, and drains the count onto the next admission.
395    #[test]
396    fn frequency_ceiling_squelches_and_counts() {
397        // 1 event per 10 seconds: within a fast test, exactly one
398        // admission fits; the rest of the burst is squelched.
399        let s = ExemplarSampler::pinned(1.0, 0.1);
400        assert_eq!(s.admit(0, 1), Some(0), "first admission passes");
401        let mut squelched = 0u64;
402        for c in 1..50 {
403            if s.admit(c, 1).is_none() {
404                squelched += 1;
405            }
406        }
407        assert_eq!(squelched, 49, "burst over the ceiling is squelched");
408        assert_eq!(s.squelched.load(Ordering::Acquire), 49);
409    }
410
411    /// A shared cell moves LIVE samplers push-on-set: flipping the
412    /// rate through the cell (what the dynamic control's applier
413    /// does) enables/disables an already-constructed sampler with
414    /// no reconstruction and no polling.
415    #[test]
416    fn shared_cell_moves_live_samplers_on_set() {
417        let cfg = std::sync::Arc::new(ExemplarConfig::new(0.0, 0.0));
418        let s = ExemplarSampler::shared(cfg.clone());
419        assert!(!s.enabled(), "starts off");
420        assert!(s.admit(1, 1).is_none());
421
422        cfg.set_rate(1.0); // the control applier's one atomic store
423        assert!(s.enabled(), "flips on push-on-set");
424        assert_eq!(s.admit(1, 1), Some(0));
425
426        cfg.set_max_hz(0.001); // ceiling: next admissions squelch
427        assert!(s.admit(2, 1).is_none());
428        assert!(s.admit(3, 1).is_none());
429
430        cfg.set_rate(0.0); // and off again, live
431        assert!(!s.enabled());
432    }
433
434    /// One advisory per class per phase, capped at 3 classes —
435    /// a storm identifies itself without flooding the output.
436    #[test]
437    fn advisory_gate_is_once_per_class_and_capped() {
438        let g = AdvisoryGate::new();
439        assert!(g.first_sighting("Overload"));
440        assert!(!g.first_sighting("Overload"), "once per class");
441        assert!(g.first_sighting("Timeout"));
442        assert!(g.first_sighting("Unavailable"));
443        assert!(!g.first_sighting("FourthClass"), "cap closes the gate");
444        assert!(!g.first_sighting("Overload"), "seen classes stay closed");
445    }
446
447    /// The advisory names the class, the budget, and both paths
448    /// forward (sampling control, opt-out).
449    #[test]
450    fn advisory_line_is_actionable() {
451        let line = render_advisory("insert", 42, 21, "Overload", "in_flight=9 > 8");
452        assert!(line.contains("op 'insert'"), "{line}");
453        assert!(line.contains("[Overload]"), "{line}");
454        assert!(line.contains("cycle 42"), "{line}");
455        assert!(line.contains("(21)"), "{line}");
456        assert!(line.contains("retry_exemplar_rate"), "{line}");
457        assert!(line.contains("retry_advisory: off"), "{line}");
458    }
459
460    /// The rendered line carries every field an operator needs to
461    /// act on the specimen, including the squelch tally.
462    #[test]
463    fn rendered_line_is_self_describing() {
464        let line = render_exemplar(&ExecExemplar {
465            op_name: "insert",
466            cycle: 12345,
467            attempt_no: 3,
468            tries_budget: 21,
469            error_class: "Overload",
470            message: "simulated overload: in_flight=9 > 8",
471            will_retry: true,
472            squelched_since_last: 7,
473        });
474        assert!(line.contains("op 'insert'"), "{line}");
475        assert!(line.contains("attempt 3/21"), "{line}");
476        assert!(line.contains("cycle 12345"), "{line}");
477        assert!(line.contains("retrying"), "{line}");
478        assert!(line.contains("[Overload]"), "{line}");
479        assert!(line.contains("in_flight=9 > 8"), "{line}");
480        assert!(line.contains("(+7 squelched)"), "{line}");
481    }
482}