use std::collections::{HashMap, VecDeque};
use std::future::Future;
use std::pin::Pin;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use arc_swap::ArcSwap;
use crate::cadence::{CadenceLayer, CadenceTree, Cadences};
use crate::labels::Labels;
use crate::scheduler::Reporter;
use crate::snapshot::{CloseReason, MetricSet};
pub const HISTORY_RING_CAP: usize = 32;
pub const COUNTER_RETAIN_FLOOR: Duration = Duration::from_secs(60);
pub const HARD_RING_CAP: usize = 200_000;
pub const DEFAULT_SUBSCRIPTION_CHANNEL_CAPACITY: usize = 8;
struct CadenceWindow {
cadence: Duration,
accumulated: Duration,
prebuffer: Option<MetricSet>,
latest: Option<Arc<MetricSet>>,
ring: VecDeque<Arc<MetricSet>>,
retain: Duration,
hist_retain: Duration,
}
impl CadenceWindow {
fn new(cadence: Duration) -> Self {
let hist_retain = cadence.saturating_mul(HISTORY_RING_CAP as u32);
let retain = COUNTER_RETAIN_FLOOR.max(hist_retain);
Self {
cadence,
accumulated: Duration::ZERO,
prebuffer: None,
latest: None,
ring: VecDeque::new(),
retain,
hist_retain,
}
}
fn evict_and_compact(&mut self) {
let Some(newest) = self.ring.back().map(|w| w.captured_at()) else {
return;
};
if let Some(cutoff) = newest.checked_sub(self.retain) {
while self.ring.front().is_some_and(|w| w.captured_at() < cutoff) {
self.ring.pop_front();
}
}
if self.retain > self.hist_retain
&& let Some(hist_cutoff) = newest.checked_sub(self.hist_retain)
{
for slot in self.ring.iter_mut() {
if slot.captured_at() >= hist_cutoff {
break; }
if slot.has_distributions() {
*slot = Arc::new(slot.without_distributions());
}
}
}
while self.ring.len() > HARD_RING_CAP {
self.ring.pop_front();
}
}
fn ingest(&mut self, snapshot: MetricSet) -> Option<Arc<MetricSet>> {
self.accumulated += snapshot.interval();
match &mut self.prebuffer {
None => self.prebuffer = Some(snapshot),
Some(buf) => {
let merged = MetricSet::coalesce(&[buf.clone(), snapshot]);
*buf = merged;
}
}
if self.accumulated >= self.cadence {
let mut closed = self.prebuffer.take().expect("prebuffer present after fold");
closed.set_interval(self.cadence);
let arc = Arc::new(closed);
self.latest = Some(arc.clone());
self.ring.push_back(arc.clone());
self.evict_and_compact();
self.accumulated = Duration::ZERO;
Some(arc)
} else {
None
}
}
fn force_close(&mut self, reason: CloseReason) -> Option<Arc<MetricSet>> {
let mut buf = self.prebuffer.take()?;
if buf.is_empty() {
return None;
}
buf.mark_partial();
buf.mark_close(reason);
let arc = Arc::new(buf);
self.latest = Some(arc.clone());
self.ring.push_back(arc.clone());
self.evict_and_compact();
self.accumulated = Duration::ZERO;
Some(arc)
}
fn latest(&self) -> Option<Arc<MetricSet>> {
self.latest.clone()
}
fn ring(&self) -> impl Iterator<Item = &Arc<MetricSet>> {
self.ring.iter()
}
fn prebuffer_clone(&self) -> Option<MetricSet> {
self.prebuffer.clone()
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub struct SubscriberId(u64);
#[derive(Debug)]
pub enum SubscribeError {
UnknownCadence(Duration),
SpawnFailed(String),
}
impl std::fmt::Display for SubscribeError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::UnknownCadence(d) => write!(f, "cadence {d:?} is not a layer of this reporter"),
Self::SpawnFailed(e) => write!(f, "failed to spawn subscription dispatch thread: {e}"),
}
}
}
impl std::error::Error for SubscribeError {}
pub type ContextWrap = Arc<
dyn Fn(Pin<Box<dyn Future<Output = ()> + Send>>) -> Pin<Box<dyn Future<Output = ()> + Send>>
+ Send
+ Sync,
>;
#[derive(Clone)]
pub struct SubscriptionOpts {
pub channel_capacity: usize,
pub timeout: Option<Duration>,
pub on_timeout: Option<TimeoutCallback>,
pub context_wrap: Option<ContextWrap>,
}
impl Default for SubscriptionOpts {
fn default() -> Self {
Self {
channel_capacity: DEFAULT_SUBSCRIPTION_CHANNEL_CAPACITY,
timeout: None,
on_timeout: None,
context_wrap: None,
}
}
}
pub type TimeoutCallback = Arc<dyn Fn(TimeoutEvent) + Send + Sync>;
#[derive(Clone, Debug)]
pub struct TimeoutEvent {
pub subscriber_id: SubscriberId,
pub cadence: Duration,
pub snapshot_age: Duration,
pub consecutive_drops: u64,
}
struct SubscriptionState {
last_delivered: Mutex<Instant>,
consecutive_drops: AtomicU64,
timeout_fired: AtomicBool,
sent: AtomicU64,
delivered: AtomicU64,
finished: AtomicBool,
}
impl SubscriptionState {
fn new() -> Self {
Self {
last_delivered: Mutex::new(Instant::now()),
consecutive_drops: AtomicU64::new(0),
timeout_fired: AtomicBool::new(false),
sent: AtomicU64::new(0),
delivered: AtomicU64::new(0),
finished: AtomicBool::new(false),
}
}
fn mark_delivered(&self) {
if let Ok(mut t) = self.last_delivered.lock() {
*t = Instant::now();
}
self.consecutive_drops.store(0, Ordering::Relaxed);
self.timeout_fired.store(false, Ordering::Relaxed);
}
fn pending(&self) -> u64 {
self.sent
.load(Ordering::Relaxed)
.saturating_sub(self.delivered.load(Ordering::Relaxed))
}
}
struct Subscription {
id: SubscriberId,
cadence: Duration,
reporter: Arc<tokio::sync::Mutex<Box<dyn Reporter>>>,
context_wrap: Option<ContextWrap>,
state: Arc<SubscriptionState>,
opts: SubscriptionOpts,
}
pub(crate) fn block_compensated<R>(f: impl FnOnce() -> R) -> R {
if tokio::runtime::Handle::try_current().is_ok() {
tokio::task::block_in_place(f)
} else {
f()
}
}
pub fn component_path_of(labels: &Labels) -> String {
let mut parts: Vec<String> = labels.iter().map(|(k, v)| format!("{k}={v}")).collect();
parts.sort_unstable();
parts.join(",")
}
#[derive(Clone)]
pub(crate) struct WindowReaderView {
pub(crate) latest: Option<Arc<MetricSet>>,
pub(crate) prebuffer: Option<Arc<MetricSet>>,
pub(crate) ring: Arc<Vec<Arc<MetricSet>>>,
}
impl Default for WindowReaderView {
fn default() -> Self {
Self {
latest: None,
prebuffer: None,
ring: Arc::new(Vec::new()),
}
}
}
#[derive(Default)]
struct ReaderState {
component_labels: HashMap<String, Labels>,
windows: HashMap<(String, usize), Arc<WindowReaderView>>,
}
enum Cmd {
Ingest {
path: String,
labels: Labels,
snapshot: MetricSet,
ack: Option<crossbeam_channel::Sender<()>>,
},
ClosePath {
path: String,
ack: Option<crossbeam_channel::Sender<()>>,
reason: CloseReason,
},
ShutdownFlushAll {
ack: FlushAck,
reason: CloseReason,
},
Barrier { ack: crossbeam_channel::Sender<()> },
}
enum FlushAck {
Sync(crossbeam_channel::Sender<()>),
Async(tokio::sync::oneshot::Sender<()>),
}
impl FlushAck {
fn signal(self) {
match self {
FlushAck::Sync(s) => {
let _ = s.send(());
}
FlushAck::Async(s) => {
let _ = s.send(());
}
}
}
}
pub struct CadenceReporter {
layers: Vec<CadenceLayer>,
declared: Cadences,
cmd_tx: tokio::sync::mpsc::UnboundedSender<Cmd>,
state: Arc<ArcSwap<ReaderState>>,
subscriptions: Arc<Mutex<HashMap<SubscriberId, Subscription>>>,
next_subscriber_id: AtomicU64,
started_at: Instant,
owner_task: Mutex<Option<tokio::task::JoinHandle<()>>>,
}
impl CadenceReporter {
pub fn new(tree: CadenceTree) -> Self {
let layers = tree.layers().to_vec();
let declared = tree.declared().clone();
let state: Arc<ArcSwap<ReaderState>> =
Arc::new(ArcSwap::from_pointee(ReaderState::default()));
let subscriptions: Arc<Mutex<HashMap<SubscriberId, Subscription>>> =
Arc::new(Mutex::new(HashMap::new()));
let (cmd_tx, cmd_rx) = tokio::sync::mpsc::unbounded_channel::<Cmd>();
let owner_layers = layers.clone();
let owner_state = state.clone();
let owner_subs = subscriptions.clone();
let owner_task = tokio::spawn(run_owner(cmd_rx, owner_layers, owner_state, owner_subs));
Self {
layers,
declared,
cmd_tx,
state,
subscriptions,
next_subscriber_id: AtomicU64::new(1),
started_at: Instant::now(),
owner_task: Mutex::new(Some(owner_task)),
}
}
pub fn layers(&self) -> &[CadenceLayer] {
&self.layers
}
pub fn declared_cadences(&self) -> &Cadences {
&self.declared
}
pub fn started_at(&self) -> Instant {
self.started_at
}
pub fn ingest(&self, labels: &Labels, snapshot: MetricSet) {
let path = component_path_of(labels);
let _ = self.cmd_tx.send(Cmd::Ingest {
path,
labels: labels.clone(),
snapshot,
ack: None,
});
}
pub fn flush_for_tests(&self) {
let (ack_tx, ack_rx) = crossbeam_channel::bounded::<()>(1);
let _ = self.cmd_tx.send(Cmd::Barrier { ack: ack_tx });
block_compensated(|| {
let _ = ack_rx.recv();
});
}
pub fn subscribe(
self: &Arc<Self>,
cadence: Duration,
reporter: Box<dyn Reporter>,
mut opts: SubscriptionOpts,
) -> Result<SubscriberId, SubscribeError> {
if !self.layers.iter().any(|l| l.interval == cadence) {
return Err(SubscribeError::UnknownCadence(cadence));
}
if opts.timeout.is_none() {
opts.timeout = Some(cadence.saturating_mul(2));
}
let id = SubscriberId(self.next_subscriber_id.fetch_add(1, Ordering::Relaxed));
let context_wrap = opts.context_wrap.clone();
let sub = Subscription {
id,
cadence,
reporter: Arc::new(tokio::sync::Mutex::new(reporter)),
context_wrap,
state: Arc::new(SubscriptionState::new()),
opts,
};
self.subscriptions
.lock()
.unwrap_or_else(|e| e.into_inner())
.insert(id, sub);
Ok(id)
}
pub fn unsubscribe(&self, id: SubscriberId) {
let sub = {
let mut map = self.subscriptions.lock().unwrap_or_else(|e| e.into_inner());
map.remove(&id)
};
if let Some(sub) = sub {
if let Ok(mut g) = sub.reporter.try_lock() {
g.flush();
}
}
}
pub async fn shutdown(&self) {
self.quiesce(Duration::from_secs(5)).await;
let subs = {
let mut map = self.subscriptions.lock().unwrap_or_else(|e| e.into_inner());
std::mem::take(&mut *map)
};
for (_id, sub) in subs {
if let Ok(mut g) = sub.reporter.try_lock() {
g.flush();
}
}
}
pub fn close_path(&self, labels: &Labels) {
self.close_path_reason(labels, CloseReason::Quiesce);
}
fn close_path_reason(&self, labels: &Labels, reason: CloseReason) {
let path = component_path_of(labels);
let _ = self.cmd_tx.send(Cmd::ClosePath {
path,
ack: None,
reason,
});
}
pub fn scope_close(&self, labels: &Labels, mut partial_delta: MetricSet) {
partial_delta.mark_partial();
partial_delta.mark_close(CloseReason::ScopeClose);
if !partial_delta.is_empty() {
self.ingest(labels, partial_delta);
}
self.close_path_reason(labels, CloseReason::ScopeClose);
}
pub async fn shutdown_flush(&self) {
let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
let _ = self.cmd_tx.send(Cmd::ShutdownFlushAll {
ack: FlushAck::Async(ack_tx),
reason: CloseReason::Shutdown,
});
let _ = ack_rx.await;
}
pub async fn quiesce(&self, max_wait: Duration) {
let (ack_tx, ack_rx) = tokio::sync::oneshot::channel();
let _ = self.cmd_tx.send(Cmd::ShutdownFlushAll {
ack: FlushAck::Async(ack_tx),
reason: CloseReason::Quiesce,
});
let _ = ack_rx.await;
let start = Instant::now();
loop {
let pending: u64 = {
let map = self.subscriptions.lock().unwrap_or_else(|e| e.into_inner());
map.values().map(|s| s.state.pending()).sum()
};
if pending == 0 || start.elapsed() >= max_wait {
break;
}
tokio::time::sleep(Duration::from_millis(1)).await;
}
}
pub fn latest(&self, labels: &Labels, cadence: Duration) -> Option<Arc<MetricSet>> {
let path = component_path_of(labels);
let idx = self.layer_index(cadence)?;
let state = self.state.load_full();
state.windows.get(&(path, idx))?.latest.clone()
}
pub fn prebuffer(&self, labels: &Labels, cadence: Duration) -> Option<MetricSet> {
let path = component_path_of(labels);
let idx = self.layer_index(cadence)?;
let state = self.state.load_full();
state
.windows
.get(&(path, idx))?
.prebuffer
.as_ref()
.map(|arc| (**arc).clone())
}
pub fn ring(&self, labels: &Labels, cadence: Duration) -> Vec<Arc<MetricSet>> {
let path = component_path_of(labels);
let Some(idx) = self.layer_index(cadence) else {
return Vec::new();
};
let state = self.state.load_full();
state
.windows
.get(&(path, idx))
.map(|w| (*w.ring).clone())
.unwrap_or_default()
}
pub(crate) fn window_view(
&self,
labels: &Labels,
cadence: Duration,
) -> Option<Arc<WindowReaderView>> {
let path = component_path_of(labels);
let idx = self.layer_index(cadence)?;
self.state.load_full().windows.get(&(path, idx)).cloned()
}
pub fn component_labels(&self) -> Vec<Labels> {
let state = self.state.load_full();
state.component_labels.values().cloned().collect()
}
fn layer_index(&self, cadence: Duration) -> Option<usize> {
self.layers.iter().position(|l| l.interval == cadence)
}
}
pub trait MetricSink: Send + Sync {
fn submit(&self, labels: &Labels, snapshot: MetricSet);
fn flush(&self) {}
}
impl MetricSink for CadenceReporter {
fn submit(&self, labels: &Labels, snapshot: MetricSet) {
self.ingest(labels, snapshot);
}
}
pub struct FanOutSink {
sinks: Vec<Arc<dyn MetricSink>>,
}
impl FanOutSink {
pub fn new(sinks: Vec<Arc<dyn MetricSink>>) -> Self {
Self { sinks }
}
}
impl MetricSink for FanOutSink {
fn submit(&self, labels: &Labels, snapshot: MetricSet) {
for sink in &self.sinks {
sink.submit(labels, snapshot.clone());
}
}
fn flush(&self) {
for sink in &self.sinks {
sink.flush();
}
}
}
impl Drop for CadenceReporter {
fn drop(&mut self) {
let (dummy_tx, dummy_rx) = tokio::sync::mpsc::unbounded_channel::<Cmd>();
drop(dummy_rx);
let _ = std::mem::replace(&mut self.cmd_tx, dummy_tx);
if let Ok(mut guard) = self.owner_task.lock()
&& let Some(task) = guard.take()
{
task.abort();
}
let subs = self
.subscriptions
.lock()
.map(|mut g| std::mem::take(&mut *g))
.unwrap_or_default();
for (_id, sub) in subs {
if let Ok(mut g) = sub.reporter.try_lock() {
g.flush();
}
}
}
}
async fn run_owner(
mut cmd_rx: tokio::sync::mpsc::UnboundedReceiver<Cmd>,
layers: Vec<CadenceLayer>,
state_pub: Arc<ArcSwap<ReaderState>>,
subscriptions: Arc<Mutex<HashMap<SubscriberId, Subscription>>>,
) {
let mut windows: HashMap<(String, usize), CadenceWindow> = HashMap::new();
let mut component_labels: HashMap<String, Labels> = HashMap::new();
'outer: loop {
let mut cmd = match cmd_rx.recv().await {
Some(c) => c,
None => break 'outer, };
let mut acks: Vec<FlushAck> = Vec::new();
loop {
match cmd {
Cmd::Ingest {
path,
labels,
snapshot,
ack,
} => {
component_labels.entry(path.clone()).or_insert(labels);
let closed_by_cadence = ingest_cascade(&mut windows, &layers, path, snapshot);
fanout_owner(&subscriptions, &closed_by_cadence);
if let Some(a) = ack {
acks.push(FlushAck::Sync(a));
}
}
Cmd::ClosePath { path, ack, reason } => {
let closed_by_cadence =
close_path_cascade(&mut windows, &layers, &path, reason);
fanout_owner(&subscriptions, &closed_by_cadence);
if let Some(a) = ack {
acks.push(FlushAck::Sync(a));
}
}
Cmd::ShutdownFlushAll { ack, reason } => {
let paths: Vec<String> = {
let mut set: std::collections::HashSet<String> =
std::collections::HashSet::new();
for (p, _) in windows.keys() {
set.insert(p.clone());
}
set.into_iter().collect()
};
let mut all_closed = Vec::new();
for path in &paths {
all_closed.extend(close_path_cascade(&mut windows, &layers, path, reason));
}
fanout_owner(&subscriptions, &all_closed);
acks.push(ack);
}
Cmd::Barrier { ack } => {
acks.push(FlushAck::Sync(ack));
}
}
cmd = match cmd_rx.try_recv() {
Ok(c) => c,
Err(_) => break, };
}
publish_reader_state(&windows, &component_labels, &state_pub);
for ack in acks.drain(..) {
ack.signal();
}
}
publish_reader_state(&windows, &component_labels, &state_pub);
}
fn ingest_cascade(
windows: &mut HashMap<(String, usize), CadenceWindow>,
layers: &[CadenceLayer],
path: String,
snapshot: MetricSet,
) -> Vec<(Duration, Arc<MetricSet>)> {
let mut closed_by_cadence: Vec<(Duration, Arc<MetricSet>)> = Vec::new();
let smallest_idx = 0usize;
let mut to_propagate: Option<(usize, Arc<MetricSet>)> = None;
let key = (path.clone(), smallest_idx);
let entry = windows
.entry(key)
.or_insert_with(|| CadenceWindow::new(layers[smallest_idx].interval));
if let Some(closed) = entry.ingest(snapshot) {
closed_by_cadence.push((layers[smallest_idx].interval, closed.clone()));
to_propagate = Some((smallest_idx + 1, closed));
}
while let Some((idx, snapshot_arc)) = to_propagate.take() {
if idx >= layers.len() {
break;
}
let key = (path.clone(), idx);
let entry = windows
.entry(key)
.or_insert_with(|| CadenceWindow::new(layers[idx].interval));
if let Some(closed) = entry.ingest((*snapshot_arc).clone()) {
closed_by_cadence.push((layers[idx].interval, closed.clone()));
to_propagate = Some((idx + 1, closed));
}
}
closed_by_cadence
}
fn close_path_cascade(
windows: &mut HashMap<(String, usize), CadenceWindow>,
layers: &[CadenceLayer],
path: &str,
reason: CloseReason,
) -> Vec<(Duration, Arc<MetricSet>)> {
let mut closed: Vec<(Duration, Arc<MetricSet>)> = Vec::new();
let mut to_propagate: Option<(usize, Arc<MetricSet>)> = None;
for idx in 0..layers.len() {
let key = (path.to_string(), idx);
if let Some((carry_idx, carry)) = to_propagate.take()
&& carry_idx == idx
{
let entry = windows
.entry(key.clone())
.or_insert_with(|| CadenceWindow::new(layers[idx].interval));
let _ = entry.ingest((*carry).clone());
}
if let Some(window) = windows.get_mut(&key)
&& let Some(snap) = window.force_close(reason)
{
closed.push((layers[idx].interval, snap.clone()));
if idx + 1 < layers.len() {
to_propagate = Some((idx + 1, snap));
}
}
}
closed
}
fn fanout_owner(
subscriptions: &Arc<Mutex<HashMap<SubscriberId, Subscription>>>,
closed: &[(Duration, Arc<MetricSet>)],
) {
if closed.is_empty() {
return;
}
let Ok(map) = subscriptions.lock() else {
return;
};
for (cadence, snapshot) in closed {
for sub in map.values() {
if sub.cadence != *cadence {
continue;
}
if sub.state.finished.load(Ordering::Relaxed) {
continue;
}
if sub.state.pending() >= sub.opts.channel_capacity as u64 {
let drops = sub.state.consecutive_drops.fetch_add(1, Ordering::Relaxed) + 1;
if let Some(timeout) = sub.opts.timeout {
let last = sub
.state
.last_delivered
.lock()
.map(|g| *g)
.unwrap_or_else(|_| Instant::now());
let age = last.elapsed();
if age >= timeout && !sub.state.timeout_fired.swap(true, Ordering::Relaxed) {
if let Some(cb) = &sub.opts.on_timeout {
cb(TimeoutEvent {
subscriber_id: sub.id,
cadence: sub.cadence,
snapshot_age: age,
consecutive_drops: drops,
});
} else {
crate::diag::warn(&format!(
"metrics subscription {:?} at cadence {:?} has stalled for \
{age:?} ({drops} consecutive drops)",
sub.id, sub.cadence,
));
}
}
}
continue;
}
sub.state.sent.fetch_add(1, Ordering::Relaxed);
let reporter = sub.reporter.clone();
let state = sub.state.clone();
let snapshot = snapshot.clone();
let fut = async move {
let mut g = reporter.lock().await;
g.report(&snapshot);
if g.finished() {
state.finished.store(true, Ordering::Relaxed);
}
drop(g);
state.mark_delivered(); state.delivered.fetch_add(1, Ordering::Relaxed);
};
let fut: Pin<Box<dyn Future<Output = ()> + Send>> = match &sub.context_wrap {
Some(w) => w(Box::pin(fut)),
None => Box::pin(fut),
};
tokio::spawn(fut);
}
}
}
fn publish_reader_state(
windows: &HashMap<(String, usize), CadenceWindow>,
component_labels: &HashMap<String, Labels>,
state_pub: &Arc<ArcSwap<ReaderState>>,
) {
let mut win_views: HashMap<(String, usize), Arc<WindowReaderView>> =
HashMap::with_capacity(windows.len());
for (key, win) in windows.iter() {
let view = WindowReaderView {
latest: win.latest(),
prebuffer: win.prebuffer_clone().map(Arc::new),
ring: Arc::new(win.ring().cloned().collect()),
};
win_views.insert(key.clone(), Arc::new(view));
}
let new_state = ReaderState {
component_labels: component_labels.clone(),
windows: win_views,
};
state_pub.store(Arc::new(new_state));
}
#[cfg(test)]
mod tests {
use super::*;
use crate::snapshot::MetricValue;
fn counter_set(interval: Duration, value: u64) -> MetricSet {
let mut s = MetricSet::new(interval);
s.insert_counter("ops", Labels::default(), value, Instant::now());
s
}
fn first_counter_total(snap: &MetricSet) -> u64 {
let f = snap.family("ops").expect("ops family");
let m = f.metrics().next().expect("series");
match m.point().unwrap().value() {
MetricValue::Counter(c) => c.cumulative,
_ => panic!("not a counter"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ingest_promotes_at_smallest_cadence_boundary() {
let cadences =
Cadences::new(&[Duration::from_millis(100), Duration::from_millis(400)]).unwrap();
let tree = CadenceTree::plan_default(cadences);
let reporter = CadenceReporter::new(tree);
let labels = Labels::of("phase", "load");
for v in [5, 10, 15, 20] {
reporter.ingest(&labels, counter_set(Duration::from_millis(100), v));
}
reporter.flush_for_tests();
let latest_100 = reporter
.latest(&labels, Duration::from_millis(100))
.expect("100ms cadence should have a latest");
assert_eq!(
first_counter_total(&latest_100),
20,
"last 100ms window's cumulative"
);
let latest_400 = reporter
.latest(&labels, Duration::from_millis(400))
.expect("400ms cadence should have promoted after 4 ticks");
assert_eq!(
first_counter_total(&latest_400),
20,
"promoted 400ms window holds the latest cumulative"
);
}
fn counter_set_at(captured_at: Instant, interval: Duration, value: u64) -> MetricSet {
let mut s = MetricSet::at(captured_at, interval);
s.insert_counter("ops", Labels::default(), value, captured_at);
s
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ring_retains_by_time_not_slot_count() {
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "x");
let base = Instant::now();
for i in 0..(HISTORY_RING_CAP + 5) {
let at = base + Duration::from_millis(i as u64);
reporter.ingest(
&labels,
counter_set_at(at, Duration::from_millis(50), (i as u64) + 1),
);
}
reporter.flush_for_tests();
let ring = reporter.ring(&labels, Duration::from_millis(50));
assert_eq!(
ring.len(),
HISTORY_RING_CAP + 5,
"time-bounded ring keeps all recent windows, not just HISTORY_RING_CAP",
);
let newest_total = first_counter_total(ring.last().unwrap());
assert_eq!(newest_total, (HISTORY_RING_CAP as u64) + 5);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn ring_evicts_windows_past_the_counter_horizon() {
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "x");
let base = Instant::now();
reporter.ingest(&labels, counter_set_at(base, Duration::from_millis(50), 1));
for i in 1..=5u64 {
let at = base + Duration::from_secs(100) + Duration::from_millis(i);
reporter.ingest(
&labels,
counter_set_at(at, Duration::from_millis(50), 1 + i),
);
}
reporter.flush_for_tests();
let ring = reporter.ring(&labels, Duration::from_millis(50));
assert_eq!(
ring.len(),
5,
"the window past the retain horizon is evicted"
);
assert_eq!(
first_counter_total(ring.first().unwrap()),
2,
"oldest retained window is the first of the recent cluster",
);
}
fn counter_and_hist_at(captured_at: Instant, interval: Duration, value: u64) -> MetricSet {
use hdrhistogram::Histogram as HdrHistogram;
let mut h = HdrHistogram::<u64>::new(3).unwrap();
h.record(value).unwrap();
let mut s = MetricSet::at(captured_at, interval);
s.insert_counter("ops", Labels::default(), value, captured_at);
s.insert_histogram("latency", Labels::default(), h, captured_at);
s
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn distributions_stripped_past_hist_horizon_but_counters_kept() {
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "h");
let base = Instant::now();
reporter.ingest(
&labels,
counter_and_hist_at(base, Duration::from_millis(50), 1),
);
reporter.ingest(
&labels,
counter_and_hist_at(base + Duration::from_secs(2), Duration::from_millis(50), 2),
);
reporter.flush_for_tests();
let ring = reporter.ring(&labels, Duration::from_millis(50));
assert_eq!(ring.len(), 2, "both windows are within the counter horizon");
let old = ring.first().unwrap();
assert!(
old.family("ops").is_some(),
"old window keeps its cumulative counter"
);
assert!(
old.family("latency").is_none(),
"old window's HDR histogram is stripped past the distribution horizon"
);
let recent = ring.last().unwrap();
assert!(
recent.family("latency").is_some(),
"the recent window keeps its histogram"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn partial_flush_is_retained_as_a_distinct_window() {
let cadences = Cadences::new(&[Duration::from_secs(1)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "p");
let base = Instant::now();
reporter.ingest(&labels, counter_set_at(base, Duration::from_secs(1), 10));
let mut partial = MetricSet::at(base + Duration::from_secs(1), Duration::from_millis(200));
partial.insert_counter("ops", Labels::default(), 5, base + Duration::from_secs(1));
reporter.scope_close(&labels, partial);
reporter.flush_for_tests();
let ring = reporter.ring(&labels, Duration::from_secs(1));
assert_eq!(
ring.len(),
2,
"the pulse window AND the partial flush are both retained"
);
assert!(!ring[0].is_partial(), "first is the pulse-closed window");
assert_eq!(first_counter_total(&ring[0]), 10);
assert!(
ring[1].is_partial(),
"second is the scope_close partial, flagged"
);
assert_eq!(first_counter_total(&ring[1]), 5);
assert!(
ring[1].interval() < Duration::from_secs(1),
"partial carries its own sub-cadence interval"
);
assert!(
ring[0].captured_at() < ring[1].captured_at(),
"distinct timestamps, not coalesced"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn metric_sink_fans_out_and_cadence_reporter_is_a_sink() {
use std::sync::atomic::{AtomicU64, Ordering};
struct Recorder(Arc<AtomicU64>);
impl MetricSink for Recorder {
fn submit(&self, _labels: &Labels, snapshot: MetricSet) {
self.0
.fetch_add(first_counter_total(&snapshot), Ordering::SeqCst);
}
}
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let rec_a = Arc::new(AtomicU64::new(0));
let rec_b = Arc::new(AtomicU64::new(0));
let fan = FanOutSink::new(vec![
reporter.clone() as Arc<dyn MetricSink>,
Arc::new(Recorder(rec_a.clone())),
Arc::new(Recorder(rec_b.clone())),
]);
let labels = Labels::of("phase", "f");
fan.submit(&labels, counter_set(Duration::from_millis(50), 7));
reporter.flush_for_tests();
assert_eq!(
rec_a.load(Ordering::SeqCst),
7,
"fan-out reached recorder A"
);
assert_eq!(
rec_b.load(Ordering::SeqCst),
7,
"fan-out reached recorder B"
);
let latest = reporter
.latest(&labels, Duration::from_millis(50))
.expect("cadence reporter ingested via MetricSink::submit");
assert_eq!(first_counter_total(&latest), 7);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn force_close_publishes_partial_at_shutdown() {
let cadences = Cadences::new(&[Duration::from_millis(1000)]).unwrap();
let tree = CadenceTree::plan_default(cadences);
let reporter = CadenceReporter::new(tree);
let labels = Labels::of("phase", "trail");
reporter.ingest(&labels, counter_set(Duration::from_millis(200), 3));
reporter.flush_for_tests();
assert!(
reporter
.latest(&labels, Duration::from_millis(1000))
.is_none()
);
reporter.shutdown_flush().await;
let partial = reporter
.latest(&labels, Duration::from_millis(1000))
.expect("shutdown must publish trailing partial");
assert_eq!(first_counter_total(&partial), 3);
assert!(
partial.interval() < Duration::from_millis(1000),
"partial interval must be < cadence: {:?}",
partial.interval()
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn prebuffer_visible_for_in_flight_data() {
let cadences = Cadences::new(&[Duration::from_millis(1000)]).unwrap();
let tree = CadenceTree::plan_default(cadences);
let reporter = CadenceReporter::new(tree);
let labels = Labels::of("phase", "p");
reporter.ingest(&labels, counter_set(Duration::from_millis(300), 7));
reporter.ingest(&labels, counter_set(Duration::from_millis(300), 8));
reporter.flush_for_tests();
let pb = reporter
.prebuffer(&labels, Duration::from_millis(1000))
.expect("prebuffer present");
assert_eq!(
first_counter_total(&pb),
8,
"in-flight prebuffer holds the latest cumulative"
);
assert!(
reporter
.latest(&labels, Duration::from_millis(1000))
.is_none()
);
}
struct CountingReporter {
count: Arc<AtomicU64>,
}
impl crate::scheduler::Reporter for CountingReporter {
fn report(&mut self, _snapshot: &MetricSet) {
self.count.fetch_add(1, Ordering::Relaxed);
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn subscribe_receives_snapshots_on_dispatch_thread() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let labels = Labels::of("phase", "sub");
let count = Arc::new(AtomicU64::new(0));
let _id = reporter
.subscribe(
Duration::from_millis(100),
Box::new(CountingReporter {
count: count.clone(),
}),
SubscriptionOpts::default(),
)
.unwrap();
for _ in 0..5 {
reporter.ingest(&labels, counter_set(Duration::from_millis(100), 1));
}
let deadline = Instant::now() + Duration::from_millis(500);
while count.load(Ordering::Relaxed) < 5 && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(10));
}
assert_eq!(count.load(Ordering::Relaxed), 5);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn quiesce_synchronously_drains_subscriber_and_is_non_terminal() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let labels = Labels::of("phase", "q");
let count = Arc::new(AtomicU64::new(0));
let _id = reporter
.subscribe(
Duration::from_millis(100),
Box::new(CountingReporter {
count: count.clone(),
}),
SubscriptionOpts::default(),
)
.unwrap();
for _ in 0..3 {
reporter.ingest(&labels, counter_set(Duration::from_millis(100), 1));
}
reporter.quiesce(Duration::from_secs(2)).await;
let after_first = count.load(Ordering::Relaxed);
assert!(
after_first >= 1,
"quiesce must synchronously deliver the force-closed window; got {after_first}"
);
reporter.ingest(&labels, counter_set(Duration::from_millis(100), 1));
reporter.quiesce(Duration::from_secs(2)).await;
assert!(
count.load(Ordering::Relaxed) > after_first,
"subscriber must remain alive + receiving after quiesce"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn subscribe_rejects_unknown_cadence() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let err = reporter
.subscribe(
Duration::from_millis(250),
Box::new(CountingReporter {
count: Arc::new(AtomicU64::new(0)),
}),
SubscriptionOpts::default(),
)
.unwrap_err();
assert!(matches!(err, SubscribeError::UnknownCadence(_)));
}
struct SlowReporter {
block: Arc<AtomicBool>,
}
impl crate::scheduler::Reporter for SlowReporter {
fn report(&mut self, _snapshot: &MetricSet) {
while self.block.load(Ordering::Relaxed) {
std::thread::sleep(Duration::from_millis(5));
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn stalled_subscriber_fires_timeout_without_blocking_cascade() {
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let labels = Labels::of("phase", "stall");
let block = Arc::new(AtomicBool::new(true));
let fired = Arc::new(AtomicU64::new(0));
let fired_for_cb = fired.clone();
let opts = SubscriptionOpts {
channel_capacity: 1, timeout: Some(Duration::from_millis(100)),
on_timeout: Some(Arc::new(move |_ev| {
fired_for_cb.fetch_add(1, Ordering::Relaxed);
})),
context_wrap: None,
};
let _id = reporter
.subscribe(
Duration::from_millis(50),
Box::new(SlowReporter {
block: block.clone(),
}),
opts,
)
.unwrap();
let start = Instant::now();
for _ in 0..20 {
reporter.ingest(&labels, counter_set(Duration::from_millis(50), 1));
std::thread::sleep(Duration::from_millis(20));
}
assert!(
start.elapsed() < Duration::from_secs(2),
"cascade took {:?} — subscriber must have blocked it",
start.elapsed()
);
assert!(
fired.load(Ordering::Relaxed) >= 1,
"expected timeout callback to fire; got {}",
fired.load(Ordering::Relaxed)
);
block.store(false, Ordering::Relaxed);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn unsubscribe_stops_delivery() {
let cadences = Cadences::new(&[Duration::from_millis(50)]).unwrap();
let reporter = Arc::new(CadenceReporter::new(CadenceTree::plan_default(cadences)));
let labels = Labels::of("phase", "unsub");
let count = Arc::new(AtomicU64::new(0));
let id = reporter
.subscribe(
Duration::from_millis(50),
Box::new(CountingReporter {
count: count.clone(),
}),
SubscriptionOpts::default(),
)
.unwrap();
reporter.ingest(&labels, counter_set(Duration::from_millis(50), 1));
let deadline = Instant::now() + Duration::from_millis(200);
while count.load(Ordering::Relaxed) < 1 && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(5));
}
assert_eq!(count.load(Ordering::Relaxed), 1);
reporter.unsubscribe(id);
reporter.ingest(&labels, counter_set(Duration::from_millis(50), 1));
reporter.ingest(&labels, counter_set(Duration::from_millis(50), 1));
std::thread::sleep(Duration::from_millis(100));
assert_eq!(count.load(Ordering::Relaxed), 1);
}
fn gauge_set(interval: Duration, value: f64) -> MetricSet {
let mut s = MetricSet::new(interval);
s.insert_gauge("temp", Labels::default(), value, Instant::now());
s
}
fn histogram_set(interval: Duration, samples: &[u64]) -> MetricSet {
use hdrhistogram::Histogram as HdrHistogram;
let mut h = HdrHistogram::<u64>::new(3).unwrap();
for v in samples {
h.record(*v).unwrap();
}
let mut s = MetricSet::new(interval);
s.insert_histogram("latency", Labels::default(), h, Instant::now());
s
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn scope_close_marks_partial_and_publishes_immediately() {
let cadences = Cadences::new(&[Duration::from_secs(1)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "short");
let mut delta = MetricSet::new(Duration::from_millis(200));
delta.insert_counter("ops", Labels::default(), 7, Instant::now());
reporter.scope_close(&labels, delta);
reporter.flush_for_tests();
let latest = reporter
.latest(&labels, Duration::from_secs(1))
.expect("scope_close must publish a partial snapshot at smallest cadence");
assert!(
latest.is_partial(),
"scope-close snapshot must be marked partial"
);
assert_eq!(first_counter_total(&latest), 7);
assert!(
latest.interval() < Duration::from_secs(1),
"partial interval must be < cadence, got {:?}",
latest.interval()
);
assert_eq!(
latest.close_reason(),
Some(CloseReason::ScopeClose),
"scope_close must stamp the typed reason (SRD-93 M4)"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn close_reason_names_the_sealer_not_the_partial_flag() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let nat = Labels::of("phase", "natural");
for _ in 0..2 {
let mut d = MetricSet::new(Duration::from_millis(50));
d.insert_counter("ops", Labels::default(), 1, Instant::now());
reporter.ingest(&nat, d);
}
reporter.flush_for_tests();
let w = reporter.latest(&nat, Duration::from_millis(100)).unwrap();
assert_eq!(
w.close_reason(),
None,
"a cadence-closed window carries no lifecycle reason"
);
let qui = Labels::of("phase", "quiescing");
let mut d = MetricSet::new(Duration::from_millis(10));
d.insert_counter("ops", Labels::default(), 2, Instant::now());
reporter.ingest(&qui, d);
reporter.quiesce(Duration::from_secs(5)).await;
let w = reporter.latest(&qui, Duration::from_millis(100)).unwrap();
assert!(w.is_partial());
assert_eq!(w.close_reason(), Some(CloseReason::Quiesce));
let end = Labels::of("phase", "ending");
let mut d = MetricSet::new(Duration::from_millis(10));
d.insert_counter("ops", Labels::default(), 3, Instant::now());
reporter.ingest(&end, d);
reporter.shutdown_flush().await;
let w = reporter.latest(&end, Duration::from_millis(100)).unwrap();
assert_eq!(w.close_reason(), Some(CloseReason::Shutdown));
assert!(CloseReason::Shutdown > CloseReason::ScopeClose);
assert!(CloseReason::ScopeClose > CloseReason::Quiesce);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn scope_close_counter_partial_carries_through_cascade() {
let cadences =
Cadences::new(&[Duration::from_millis(100), Duration::from_millis(400)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "burst");
let mut partial = MetricSet::new(Duration::from_millis(50));
partial.insert_counter("ops", Labels::default(), 5, Instant::now());
reporter.scope_close(&labels, partial);
reporter.flush_for_tests();
let p100 = reporter
.latest(&labels, Duration::from_millis(100))
.unwrap();
assert_eq!(first_counter_total(&p100), 5);
assert!(
p100.is_partial(),
"smallest-cadence publish must be partial"
);
let p400 = reporter
.latest(&labels, Duration::from_millis(400))
.expect("close_path cascade publishes at every layer");
assert_eq!(
first_counter_total(&p400),
5,
"partial cascades unchanged through the chain"
);
assert!(
p400.is_partial(),
"partial flag is sticky across the cascade fold"
);
for v in [6, 9, 12, 15] {
reporter.ingest(&labels, counter_set(Duration::from_millis(100), v));
}
reporter.flush_for_tests();
let np400 = reporter
.latest(&labels, Duration::from_millis(400))
.unwrap();
assert_eq!(
first_counter_total(&np400),
15,
"natural-pulse window holds the latest cumulative"
);
assert!(
!np400.is_partial(),
"natural-pulse close must NOT be marked partial"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn scope_close_gauge_partials_use_combine_rules() {
let cadences = Cadences::new(&[Duration::from_secs(1)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "g");
reporter.scope_close(&labels, gauge_set(Duration::from_millis(200), 4.0));
reporter.flush_for_tests();
let cadences = Cadences::new(&[Duration::from_secs(1)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
reporter.ingest(&labels, gauge_set(Duration::from_millis(200), 4.0));
reporter.ingest(&labels, gauge_set(Duration::from_millis(200), 8.0));
reporter.scope_close(&labels, MetricSet::new(Duration::ZERO));
reporter.flush_for_tests();
let latest = reporter.latest(&labels, Duration::from_secs(1)).unwrap();
assert!(
latest.is_partial(),
"scope_close must publish prebuffer as partial"
);
let g = latest
.family("temp")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value();
let value = match g {
MetricValue::Gauge(g) => g.value,
_ => panic!("expected gauge, got {:?}", g),
};
assert!(
(value - 8.0).abs() < 1e-9,
"expected last-written 8.0, got {value}"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn scope_close_histogram_partials_hdr_merge() {
use hdrhistogram::Histogram as HdrHistogram;
let cadences = Cadences::new(&[Duration::from_secs(1)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "h");
reporter.ingest(
&labels,
histogram_set(Duration::from_millis(100), &[10, 20, 30]),
);
reporter.ingest(
&labels,
histogram_set(Duration::from_millis(100), &[100, 200, 300]),
);
reporter.scope_close(&labels, MetricSet::new(Duration::ZERO));
reporter.flush_for_tests();
let latest = reporter.latest(&labels, Duration::from_secs(1)).unwrap();
assert!(latest.is_partial());
let v = latest
.family("latency")
.unwrap()
.metrics()
.next()
.unwrap()
.point()
.unwrap()
.value();
let h: &HdrHistogram<u64> = match v {
MetricValue::Histogram(h) => h.reservoir.as_ref(),
_ => panic!("expected histogram"),
};
assert_eq!(h.len(), 6, "all six histogram samples must merge");
assert!(h.value_at_quantile(0.0) <= 10);
assert!(h.value_at_quantile(1.0) >= 300);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn scope_close_only_runs_when_component_is_running() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let reporter = CadenceReporter::new(CadenceTree::plan_default(cadences));
let labels = Labels::of("phase", "never");
reporter.scope_close(&labels, MetricSet::new(Duration::ZERO));
reporter.flush_for_tests();
assert!(
reporter
.latest(&labels, Duration::from_millis(100))
.is_none(),
"empty-delta scope_close on never-seen path must not invent a snapshot"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn separate_components_keyed_independently() {
let cadences = Cadences::new(&[Duration::from_millis(100)]).unwrap();
let tree = CadenceTree::plan_default(cadences);
let reporter = CadenceReporter::new(tree);
reporter.ingest(
&Labels::of("phase", "a"),
counter_set(Duration::from_millis(100), 1),
);
reporter.ingest(
&Labels::of("phase", "b"),
counter_set(Duration::from_millis(100), 99),
);
reporter.flush_for_tests();
assert_eq!(
first_counter_total(
&reporter
.latest(&Labels::of("phase", "a"), Duration::from_millis(100))
.unwrap()
),
1
);
assert_eq!(
first_counter_total(
&reporter
.latest(&Labels::of("phase", "b"), Duration::from_millis(100))
.unwrap()
),
99
);
}
}