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}