use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
pub struct ExecExemplar<'a> {
pub op_name: &'a str,
pub cycle: u64,
pub attempt_no: u32,
pub tries_budget: u32,
pub error_class: &'a str,
pub message: &'a str,
pub will_retry: bool,
pub squelched_since_last: u64,
}
pub fn render_exemplar(ex: &ExecExemplar<'_>) -> String {
let retry_note = if ex.will_retry {
"retrying"
} else {
"terminal"
};
let squelch_note = if ex.squelched_since_last > 0 {
format!(" (+{} squelched)", ex.squelched_since_last)
} else {
String::new()
};
format!(
"retry exemplar: op '{}' attempt {}/{} cycle {} ({retry_note}): \
[{}] {}{squelch_note}",
ex.op_name, ex.attempt_no, ex.tries_budget, ex.cycle, ex.error_class, ex.message,
)
}
pub fn submit_exemplar(ex: &ExecExemplar<'_>) {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Retry),
&render_exemplar(ex),
);
}
pub trait ExecEventSubscriber {
fn submit_exemplar(&self, ex: &ExecExemplar<'_>) {
submit_exemplar(ex)
}
fn submit_advisory(
&self,
op_name: &str,
cycle: u64,
tries_budget: u32,
error_class: &str,
message: &str,
) {
crate::observer::log_tagged(
crate::observer::LogLevel::Warn,
crate::observer::EventTag::in_flight(crate::observer::EventCategory::Retry),
&render_advisory(op_name, cycle, tries_budget, error_class, message),
);
}
}
pub fn render_advisory(
op_name: &str,
cycle: u64,
tries_budget: u32,
error_class: &str,
message: &str,
) -> String {
format!(
"retry advisory: op '{op_name}' hit its first retryable \
[{error_class}] at cycle {cycle}: {message} — further \
occurrences are absorbed by the tries budget ({tries_budget}) \
and appear only as att:%/r: chips and attempt_* metrics; \
sample live specimens via the retry_exemplar_rate control, \
or silence this line with retry_advisory: off"
)
}
pub struct AdvisoryGate {
seen: std::sync::Mutex<std::collections::HashSet<String>>,
cap: usize,
}
impl Default for AdvisoryGate {
fn default() -> Self {
Self::new()
}
}
impl AdvisoryGate {
pub fn new() -> Self {
Self {
seen: std::sync::Mutex::new(std::collections::HashSet::new()),
cap: 3,
}
}
pub fn first_sighting(&self, class: &str) -> bool {
let mut seen = self.seen.lock().unwrap_or_else(|e| e.into_inner());
if seen.len() >= self.cap && !seen.contains(class) {
return false;
}
seen.insert(class.to_string())
}
}
pub(crate) fn splitmix64(mut z: u64) -> u64 {
z = z.wrapping_add(0x9E37_79B9_7F4A_7C15);
z = (z ^ (z >> 30)).wrapping_mul(0xBF58_476D_1CE4_E5B9);
z = (z ^ (z >> 27)).wrapping_mul(0x94D0_49BB_1331_11EB);
z ^ (z >> 31)
}
pub struct ExemplarConfig {
rate_bits: AtomicU64,
min_interval_nanos: AtomicU64,
}
impl ExemplarConfig {
pub fn new(rate: f64, max_hz: f64) -> Self {
let cfg = Self {
rate_bits: AtomicU64::new(0),
min_interval_nanos: AtomicU64::new(0),
};
cfg.set_rate(rate);
cfg.set_max_hz(max_hz);
cfg
}
pub fn set_rate(&self, rate: f64) {
let rate = if rate.is_finite() {
rate.clamp(0.0, 1.0)
} else {
0.0
};
self.rate_bits.store(rate.to_bits(), Ordering::Release);
}
pub fn set_max_hz(&self, max_hz: f64) {
let interval = if max_hz.is_finite() && max_hz > 0.0 {
(1_000_000_000f64 / max_hz) as u64
} else {
0
};
self.min_interval_nanos.store(interval, Ordering::Release);
}
fn rate(&self) -> f64 {
f64::from_bits(self.rate_bits.load(Ordering::Acquire))
}
fn min_interval(&self) -> u64 {
self.min_interval_nanos.load(Ordering::Acquire)
}
}
pub struct ExemplarSampler {
cfg: std::sync::Arc<ExemplarConfig>,
base: Instant,
last_admit: AtomicU64,
squelched: AtomicU64,
}
impl ExemplarSampler {
pub fn pinned(rate: f64, max_hz: f64) -> Self {
Self::shared(std::sync::Arc::new(ExemplarConfig::new(rate, max_hz)))
}
pub fn shared(cfg: std::sync::Arc<ExemplarConfig>) -> Self {
Self {
cfg,
base: Instant::now(),
last_admit: AtomicU64::new(0),
squelched: AtomicU64::new(0),
}
}
pub fn enabled(&self) -> bool {
self.cfg.rate() > 0.0
}
pub fn admit(&self, cycle: u64, attempt_no: u32) -> Option<u64> {
let rate = self.cfg.rate();
if rate <= 0.0 {
return None;
}
let h = splitmix64(cycle ^ ((attempt_no as u64) << 48) ^ 0xE0E0_5EED);
if (h as f64 / u64::MAX as f64) >= rate {
return None;
}
let min_interval = self.cfg.min_interval();
if min_interval == 0 {
let now = self.base.elapsed().as_nanos() as u64 + 1;
self.last_admit.store(now, Ordering::Release);
return Some(self.squelched.swap(0, Ordering::AcqRel));
}
let now = self.base.elapsed().as_nanos() as u64 + 1;
loop {
let last = self.last_admit.load(Ordering::Acquire);
if last != 0 && now.saturating_sub(last) < min_interval {
self.squelched.fetch_add(1, Ordering::AcqRel);
return None;
}
if self
.last_admit
.compare_exchange(last, now, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
{
return Some(self.squelched.swap(0, Ordering::AcqRel));
}
}
}
}
impl Drop for ExemplarSampler {
fn drop(&mut self) {
let leftover = self.squelched.load(Ordering::Acquire);
if leftover > 0 {
crate::diag!(
crate::observer::LogLevel::Debug,
"exemplar sampler retired with {leftover} squelched \
admission(s) unreported (frequency ceiling)"
);
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn zero_rate_is_off() {
let s = ExemplarSampler::pinned(0.0, 1000.0);
assert!(!s.enabled());
for c in 0..1000 {
assert!(s.admit(c, 1).is_none());
}
}
#[test]
fn full_rate_uncapped_admits_all() {
let s = ExemplarSampler::pinned(1.0, 0.0);
for c in 0..100 {
assert_eq!(s.admit(c, 1), Some(0), "cycle {c}");
}
}
#[test]
fn fractional_rate_samples_deterministically() {
let s1 = ExemplarSampler::pinned(0.25, 0.0);
let s2 = ExemplarSampler::pinned(0.25, 0.0);
let picks1: Vec<bool> = (0..4000).map(|c| s1.admit(c, 3).is_some()).collect();
let picks2: Vec<bool> = (0..4000).map(|c| s2.admit(c, 3).is_some()).collect();
assert_eq!(picks1, picks2, "sampling must be replayable");
let hits = picks1.iter().filter(|b| **b).count();
assert!(
(600..=1400).contains(&hits),
"0.25 of 4000 should land near 1000, got {hits}"
);
}
#[test]
fn frequency_ceiling_squelches_and_counts() {
let s = ExemplarSampler::pinned(1.0, 0.1);
assert_eq!(s.admit(0, 1), Some(0), "first admission passes");
let mut squelched = 0u64;
for c in 1..50 {
if s.admit(c, 1).is_none() {
squelched += 1;
}
}
assert_eq!(squelched, 49, "burst over the ceiling is squelched");
assert_eq!(s.squelched.load(Ordering::Acquire), 49);
}
#[test]
fn shared_cell_moves_live_samplers_on_set() {
let cfg = std::sync::Arc::new(ExemplarConfig::new(0.0, 0.0));
let s = ExemplarSampler::shared(cfg.clone());
assert!(!s.enabled(), "starts off");
assert!(s.admit(1, 1).is_none());
cfg.set_rate(1.0); assert!(s.enabled(), "flips on push-on-set");
assert_eq!(s.admit(1, 1), Some(0));
cfg.set_max_hz(0.001); assert!(s.admit(2, 1).is_none());
assert!(s.admit(3, 1).is_none());
cfg.set_rate(0.0); assert!(!s.enabled());
}
#[test]
fn advisory_gate_is_once_per_class_and_capped() {
let g = AdvisoryGate::new();
assert!(g.first_sighting("Overload"));
assert!(!g.first_sighting("Overload"), "once per class");
assert!(g.first_sighting("Timeout"));
assert!(g.first_sighting("Unavailable"));
assert!(!g.first_sighting("FourthClass"), "cap closes the gate");
assert!(!g.first_sighting("Overload"), "seen classes stay closed");
}
#[test]
fn advisory_line_is_actionable() {
let line = render_advisory("insert", 42, 21, "Overload", "in_flight=9 > 8");
assert!(line.contains("op 'insert'"), "{line}");
assert!(line.contains("[Overload]"), "{line}");
assert!(line.contains("cycle 42"), "{line}");
assert!(line.contains("(21)"), "{line}");
assert!(line.contains("retry_exemplar_rate"), "{line}");
assert!(line.contains("retry_advisory: off"), "{line}");
}
#[test]
fn rendered_line_is_self_describing() {
let line = render_exemplar(&ExecExemplar {
op_name: "insert",
cycle: 12345,
attempt_no: 3,
tries_budget: 21,
error_class: "Overload",
message: "simulated overload: in_flight=9 > 8",
will_retry: true,
squelched_since_last: 7,
});
assert!(line.contains("op 'insert'"), "{line}");
assert!(line.contains("attempt 3/21"), "{line}");
assert!(line.contains("cycle 12345"), "{line}");
assert!(line.contains("retrying"), "{line}");
assert!(line.contains("[Overload]"), "{line}");
assert!(line.contains("in_flight=9 > 8"), "{line}");
assert!(line.contains("(+7 squelched)"), "{line}");
}
}