Skip to main content

nmbrs_metrics/
scheduler.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Metrics snapshot scheduler with hierarchical frame coalescing.
5//!
6//! A dedicated thread captures frames at the base interval from the
7//! component tree. Each reporter is registered at its own interval
8//! (must be an exact multiple of the base). Schedule nodes accumulate
9//! and coalesce frames for slower reporters.
10//!
11//! At every tick the scheduler also feeds the installed
12//! [`CadenceReporter`] (SRD-42), which owns the windowed snapshot
13//! store read by every consumer through
14//! [`crate::metrics_query::MetricsQuery`].
15
16use std::sync::{Arc, Condvar, Mutex};
17use std::time::{Duration, Instant};
18
19use crate::cadence_reporter::CadenceReporter;
20use crate::labels::Labels;
21use crate::snapshot::MetricSet;
22
23/// Trait for metrics reporters (external consumers: SQLite, CSV, etc.).
24pub trait Reporter: Send + 'static {
25    fn report(&mut self, snapshot: &MetricSet);
26    fn flush(&mut self) {}
27
28    /// Self-termination signal. After a [`report`](Reporter::report)
29    /// that leaves this `true`, the subscriber's cadence-feed dispatch
30    /// worker exits its loop (calling [`flush`](Reporter::flush) on the
31    /// way out) — the subscriber receives no further pulses. A one-shot
32    /// subscriber — e.g. a settle / stop evaluator that has set a
33    /// terminal phase disposition — uses this to **unregister itself**
34    /// without a self-join deadlock (it runs on the worker thread, so it
35    /// cannot call `unsubscribe` on itself directly). Default `false`
36    /// (a long-lived subscriber that never self-terminates).
37    fn finished(&self) -> bool {
38        false
39    }
40}
41
42/// Capture function that produces per-component delta snapshots
43/// from the component tree.
44///
45/// Returns one `(effective_labels, delta_snapshot)` per RUNNING
46/// component that has instruments with data.
47pub type CaptureFunc = Box<dyn Fn() -> Vec<(Labels, MetricSet)> + Send>;
48
49/// A node in the schedule tree that accumulates and coalesces snapshots.
50struct ScheduleNode {
51    interval: Duration,
52    accumulated: Vec<MetricSet>,
53    accumulated_duration: Duration,
54    reporters: Vec<Box<dyn Reporter>>,
55    children: Vec<ScheduleNode>,
56}
57
58impl ScheduleNode {
59    fn new(interval: Duration) -> Self {
60        Self {
61            interval,
62            accumulated: Vec::new(),
63            accumulated_duration: Duration::ZERO,
64            reporters: Vec::new(),
65            children: Vec::new(),
66        }
67    }
68
69    /// Ingest a combined snapshot. Accumulate, and when the
70    /// interval is satisfied, coalesce and emit.
71    fn ingest(&mut self, snapshot: MetricSet) {
72        self.accumulated_duration += snapshot.interval();
73        self.accumulated.push(snapshot);
74
75        if self.accumulated_duration >= self.interval {
76            let coalesced = MetricSet::coalesce(&self.accumulated);
77            self.accumulated.clear();
78            self.accumulated_duration = Duration::ZERO;
79
80            for reporter in &mut self.reporters {
81                reporter.report(&coalesced);
82            }
83            for child in &mut self.children {
84                child.ingest(coalesced.clone());
85            }
86        }
87    }
88}
89
90/// Configuration for the snapshot scheduler.
91pub struct SchedulerConfig {
92    pub base_interval: Duration,
93}
94
95impl Default for SchedulerConfig {
96    fn default() -> Self {
97        Self {
98            base_interval: Duration::from_secs(1),
99        }
100    }
101}
102
103/// Builder for constructing a scheduler with reporters.
104pub struct SchedulerBuilder {
105    config: SchedulerConfig,
106    reporters: Vec<(Duration, Box<dyn Reporter>)>,
107    cadence_reporter: Option<Arc<CadenceReporter>>,
108    cadence_tree: Option<crate::cadence::CadenceTree>,
109}
110
111impl Default for SchedulerBuilder {
112    fn default() -> Self {
113        Self::new()
114    }
115}
116
117impl SchedulerBuilder {
118    pub fn new() -> Self {
119        Self {
120            config: SchedulerConfig::default(),
121            reporters: Vec::new(),
122            cadence_reporter: None,
123            cadence_tree: None,
124        }
125    }
126
127    pub fn base_interval(mut self, interval: Duration) -> Self {
128        self.config.base_interval = interval;
129        self
130    }
131
132    pub fn add_reporter(mut self, interval: Duration, reporter: impl Reporter) -> Self {
133        self.reporters.push((interval, Box::new(reporter)));
134        self
135    }
136
137    /// Install the cadence reporter that owns the windowed snapshot
138    /// store. On every scheduler tick, captured per-component
139    /// snapshots are fed into this reporter, which cascades them
140    /// up the cadence tree and publishes closed windows to
141    /// [`crate::metrics_query::MetricsQuery`] readers.
142    pub fn with_cadence_reporter(mut self, reporter: Arc<CadenceReporter>) -> Self {
143        self.cadence_reporter = Some(reporter);
144        self
145    }
146
147    /// Install a cadence tree (SRD-42 §"Tree Construction"). When set,
148    /// `build()` constructs a chained schedule where each layer feeds
149    /// the next via [`ScheduleNode::ingest`] rather than coalescing
150    /// from base frames independently. Hidden layers participate in
151    /// accumulation but have no reporters of their own.
152    ///
153    /// Reporters at intervals matching a tree layer attach at that
154    /// layer; reporters at intervals outside the tree continue to
155    /// attach as flat children of root (backward-compatible).
156    pub fn with_cadence_tree(mut self, tree: crate::cadence::CadenceTree) -> Self {
157        self.cadence_tree = Some(tree);
158        self
159    }
160
161    /// Build the schedule tree and return a handle.
162    ///
163    /// The scheduler is not yet running — call `start()` on the handle.
164    pub fn build(self, capture: CaptureFunc) -> SchedulerHandle {
165        let base = self.config.base_interval;
166        let mut root = ScheduleNode::new(base);
167
168        let mut by_interval: std::collections::BTreeMap<Duration, Vec<Box<dyn Reporter>>> =
169            std::collections::BTreeMap::new();
170        for (interval, reporter) in self.reporters {
171            by_interval.entry(interval).or_default().push(reporter);
172        }
173
174        // Reporters that match the base interval always live on root.
175        if let Some(reps) = by_interval.remove(&base) {
176            root.reporters.extend(reps);
177        }
178
179        // If a cadence tree was provided, build the chained sub-tree.
180        // Walking layers largest → smallest builds the chain from the
181        // leaf inward, so each node owns its single child.
182        if let Some(tree) = self.cadence_tree {
183            let mut chain: Option<ScheduleNode> = None;
184            for layer in tree.layers().iter().rev() {
185                if layer.interval == base {
186                    // Base-interval "layer" is just the root itself —
187                    // any reporters at that interval are already on
188                    // root. Skip without nesting.
189                    continue;
190                }
191                assert!(
192                    layer.interval.as_millis() % base.as_millis() == 0,
193                    "cadence layer {:?} must be an exact multiple of base {:?}",
194                    layer.interval,
195                    base,
196                );
197                let mut node = ScheduleNode::new(layer.interval);
198                if !layer.hidden
199                    && let Some(reps) = by_interval.remove(&layer.interval)
200                {
201                    node.reporters = reps;
202                }
203                if let Some(child) = chain.take() {
204                    node.children.push(child);
205                }
206                chain = Some(node);
207            }
208            if let Some(top) = chain {
209                root.children.push(top);
210            }
211        }
212
213        // Reporters not consumed by the tree (intervals outside it,
214        // or no tree at all) attach as flat children of root — same
215        // behavior as before this layering existed.
216        for (interval, reporters) in by_interval {
217            assert!(
218                interval.as_millis() % base.as_millis() == 0,
219                "reporter interval {:?} must be an exact multiple of base {:?}",
220                interval,
221                base
222            );
223            let mut node = ScheduleNode::new(interval);
224            node.reporters = reporters;
225            root.children.push(node);
226        }
227
228        SchedulerHandle {
229            root: Arc::new(Mutex::new(root)),
230            capture,
231            base_interval: base,
232            running: Arc::new(Mutex::new(false)),
233            cadence_reporter: self.cadence_reporter,
234        }
235    }
236}
237
238/// Handle to a running (or startable) scheduler.
239pub struct SchedulerHandle {
240    root: Arc<Mutex<ScheduleNode>>,
241    capture: CaptureFunc,
242    base_interval: Duration,
243    running: Arc<Mutex<bool>>,
244    cadence_reporter: Option<Arc<CadenceReporter>>,
245}
246
247impl SchedulerHandle {
248    /// Reference to the installed cadence reporter, if any.
249    pub fn cadence_reporter(&self) -> Option<&Arc<CadenceReporter>> {
250        self.cadence_reporter.as_ref()
251    }
252
253    /// Flush a retiring component's final delta through the
254    /// cadence reporter (if present). Called from the executor
255    /// thread when a phase completes, outside the scheduler tick
256    /// loop.
257    pub fn flush_component(&self, labels: &Labels, final_delta: MetricSet) {
258        if let Some(reporter) = &self.cadence_reporter {
259            reporter.ingest(labels, final_delta);
260        }
261    }
262
263    /// Start the scheduler on a dedicated thread.
264    ///
265    /// Returns a `StopHandle` that can be used to shut down.
266    pub fn start(self) -> StopHandle {
267        let root = self.root.clone();
268        let root_for_stop = self.root;
269        let capture = self.capture;
270        let interval = self.base_interval;
271        let running = self.running.clone();
272        let cadence_reporter = self.cadence_reporter.clone();
273        let cadence_reporter_for_stop = self.cadence_reporter.clone();
274
275        let (frame_tx, frame_rx) = std::sync::mpsc::channel::<MetricSet>();
276        // Hot-path split (SRD-102 §6): the `timing` thread only *captures*
277        // deltas and enqueues here; a single ordered `io`-pool worker drains
278        // this channel and does the potentially-slow reporter delivery
279        // (report()/ingest — CSV/SQLite/HTTP), keeping the timing thread's
280        // critical section to capture + enqueue.
281        let (io_tx, io_rx) = std::sync::mpsc::channel::<MetricSet>();
282        // Stop signal: `running` (a Mutex<bool>) gates the loop and the
283        // Condvar wakes the timing thread out of its inter-tick wait
284        // immediately on stop instead of dwelling a full base interval.
285        // Condvar + Arc<Mutex> are `Sync`, so a session host can still stop
286        // through a shared `Arc<StopHandle>` now that the wait is a std
287        // primitive rather than a tokio `Notify`.
288        let stop_cv = Arc::new(Condvar::new());
289        let stop_cv_thread = stop_cv.clone();
290        // The timing thread fires `done` AFTER its final flush, so the sync
291        // `stop()` can wait for the trailing window to land (summary reports
292        // read complete data).
293        let (done_tx, done_rx) = std::sync::mpsc::channel::<()>();
294
295        *running.lock().unwrap_or_else(|e| e.into_inner()) = true;
296
297        let stop_running = running.clone();
298        // Capture a runtime handle (start() is called from the async runner)
299        // so the std timing thread can `block_on` the one async shutdown call
300        // (`cadence_reporter.shutdown_flush`). The timing thread is never a
301        // runtime worker, so blocking on it is safe.
302        let rt_handle = tokio::runtime::Handle::try_current().ok();
303
304        // The ordered `io`-pool reporter worker. Exits when the timing thread
305        // drops `io_tx` on shutdown, after it has delivered every enqueued
306        // snapshot. A single consumer preserves reporter delivery order (the
307        // `io` pool's thread *count* is capacity for other future consumers).
308        let root_io = root.clone();
309        let io_handle = crate::thread_pools::global()
310            .spawn("io", "reporter", move || {
311                while let Ok(snapshot) = io_rx.recv() {
312                    let mut node = root_io.lock().unwrap_or_else(|e| e.into_inner());
313                    for reporter in &mut node.reporters {
314                        reporter.report(&snapshot);
315                    }
316                    for child in &mut node.children {
317                        child.ingest(snapshot.clone());
318                    }
319                }
320            })
321            .expect("spawn io reporter thread");
322
323        // The cadence tick loop runs on a dedicated `timing`-pool OS thread
324        // (SRD-102): realtime scheduling policy + affinity applied at spawn,
325        // never sharing duty with the async worker runtime, so timer wake-ups
326        // are not queued behind workload fibers.
327        let sched_thread = crate::thread_pools::global()
328            .spawn_timing("cadence", move || {
329                // Divergence surveillance (SRD-102 §6): compare the nominal
330                // deadline (`scheduled_ts`) to the actual fire instant. If
331                // they diverge by more than 250 ms the timing thread is being
332                // delayed (CPU starvation / oversleep). Warn — rate-limited —
333                // so the operator can correlate anomalies with scheduler
334                // health. Recorded snapshot intervals stay at the nominal
335                // cadence (canonical cadence as a matter of record); the
336                // divergence is reported out-of-band and stamped on the
337                // snapshot as scheduled_ts vs actual_ts.
338                let divergence_threshold = Duration::from_millis(250);
339                let divergence_warn_min_interval = Duration::from_secs(60);
340                let mut last_divergence_warn: Option<Instant> = None;
341                let mut next_tick = Instant::now() + interval;
342                loop {
343                    // Interruptible wait to the absolute `next_tick`. Holds the
344                    // `running` guard across `wait_timeout` (which atomically
345                    // releases + reacquires), so a stop set by `stop()` is seen
346                    // the instant the Condvar wakes us.
347                    let fire = {
348                        let mut guard = stop_running.lock().unwrap_or_else(|e| e.into_inner());
349                        loop {
350                            if !*guard {
351                                break false;
352                            }
353                            let now = Instant::now();
354                            if now >= next_tick {
355                                break true;
356                            }
357                            let (g, _) = stop_cv_thread
358                                .wait_timeout(guard, next_tick - now)
359                                .unwrap_or_else(|e| e.into_inner());
360                            guard = g;
361                        }
362                    };
363                    if !fire {
364                        break;
365                    }
366
367                    let scheduled = next_tick;
368                    // Fixed-rate: advance by the nominal interval regardless of
369                    // when we actually woke, so cadence does not drift.
370                    next_tick += interval;
371                    let actual = Instant::now();
372
373                    if let Some(divergence) = divergence_warning(
374                        scheduled,
375                        actual,
376                        divergence_threshold,
377                        divergence_warn_min_interval,
378                        last_divergence_warn,
379                    ) {
380                        last_divergence_warn = Some(actual);
381                        crate::diag::warn(&format!(
382                            "scheduler cadence divergence: scheduled vs actual off \
383                             by {:?} (>250ms) — snapshots still recorded at nominal \
384                             cadence; the `timing` pool thread is being delayed \
385                             (CPU starvation / oversleep)",
386                            divergence,
387                        ));
388                    }
389
390                    // Drain async snapshot channel (lifecycle flushes from
391                    // executor) → offload delivery to the io worker.
392                    while let Ok(snapshot) = frame_rx.try_recv() {
393                        let _ = io_tx.send(snapshot);
394                    }
395
396                    // Capture per-component deltas from the tree.
397                    let component_snapshots = (capture)();
398
399                    // Feed each per-component delta into the cadence reporter
400                    // (single writer of windowed snapshots — a non-blocking
401                    // crossbeam send, kept on the timing thread as part of
402                    // capture).
403                    if let Some(ref cr) = cadence_reporter {
404                        for (labels, snapshot) in &component_snapshots {
405                            cr.ingest(labels, snapshot.clone());
406                        }
407                    }
408
409                    // Merge component snapshots into one combined snapshot for
410                    // the scheduler-tree reporters (CSV / SQLite / etc.), stamp
411                    // the scheduled/actual timestamp pair, and hand off to io.
412                    let all_snapshots: Vec<MetricSet> = component_snapshots
413                        .into_iter()
414                        .map(|(_, snapshot)| snapshot)
415                        .collect();
416                    let mut combined = if all_snapshots.is_empty() {
417                        MetricSet::new(interval)
418                    } else {
419                        let mut merged = MetricSet::coalesce(&all_snapshots);
420                        // Interval reflects the scheduler interval, not the sum
421                        // from coalesce (which sums intervals).
422                        merged.set_interval(interval);
423                        merged
424                    };
425                    combined.set_scheduled_ts(scheduled);
426                    let _ = io_tx.send(combined);
427                }
428
429                // Final capture before shutdown: ensures short-lived phases
430                // that completed between ticks get their data to reporters.
431                {
432                    let component_snapshots = (capture)();
433                    if let Some(ref cr) = cadence_reporter {
434                        for (labels, snapshot) in &component_snapshots {
435                            cr.ingest(labels, snapshot.clone());
436                        }
437                    }
438                    let all_snapshots: Vec<MetricSet> = component_snapshots
439                        .into_iter()
440                        .map(|(_, snapshot)| snapshot)
441                        .collect();
442                    if !all_snapshots.is_empty() {
443                        let mut merged = MetricSet::coalesce(&all_snapshots);
444                        merged.set_interval(interval);
445                        let _ = io_tx.send(merged);
446                    }
447                }
448
449                // No more steady-state deliveries: close the io channel and
450                // join the reporter worker so every enqueued snapshot has
451                // landed before the final flush. After this the timing thread
452                // is the sole toucher of `root`.
453                drop(io_tx);
454                let _ = io_handle.join();
455
456                // Force-close any unpromoted cadence partials so the trailing
457                // window is not lost. The only async call — `block_on` on this
458                // dedicated (non-runtime) thread.
459                if let (Some(h), Some(cr)) = (rt_handle.as_ref(), cadence_reporter.as_ref()) {
460                    h.block_on(cr.shutdown_flush());
461                }
462
463                // Drain any remaining async frames directly (io worker joined).
464                while let Ok(snapshot) = frame_rx.try_recv() {
465                    let mut r = root.lock().unwrap_or_else(|e| e.into_inner());
466                    for reporter in &mut r.reporters {
467                        reporter.report(&snapshot);
468                    }
469                    for child in &mut r.children {
470                        child.ingest(snapshot.clone());
471                    }
472                }
473                // Flush all reporters on shutdown.
474                flush_tree(&mut root.lock().unwrap_or_else(|e| e.into_inner()));
475                // Trailing window has landed — release a waiting `stop()`.
476                let _ = done_tx.send(());
477            })
478            .expect("spawn timing scheduler thread");
479
480        StopHandle {
481            running: self.running,
482            cadence_reporter: cadence_reporter_for_stop,
483            root: root_for_stop,
484            task: Mutex::new(Some(sched_thread)),
485            frame_tx,
486            stop_cv,
487            done_rx: Mutex::new(Some(done_rx)),
488        }
489    }
490}
491
492/// SRD-102 §6 divergence-warning decision (extracted for testability). The
493/// nominal `scheduled` deadline vs the `actual` fire instant must diverge by
494/// more than `threshold`, AND `min_interval` must have elapsed since the last
495/// warning (rate-limit so sustained drift warns once per window, not every
496/// tick). Returns the divergence magnitude when a warning is due.
497fn divergence_warning(
498    scheduled: Instant,
499    actual: Instant,
500    threshold: Duration,
501    min_interval: Duration,
502    last_warn: Option<Instant>,
503) -> Option<Duration> {
504    let divergence = if actual >= scheduled {
505        actual - scheduled
506    } else {
507        scheduled - actual
508    };
509    let due = last_warn
510        .map(|t| actual.duration_since(t) >= min_interval)
511        .unwrap_or(true);
512    (divergence > threshold && due).then_some(divergence)
513}
514
515fn flush_tree(node: &mut ScheduleNode) {
516    for reporter in &mut node.reporters {
517        reporter.flush();
518    }
519    for child in &mut node.children {
520        flush_tree(child);
521    }
522}
523
524/// Handle to stop a running scheduler.
525pub struct StopHandle {
526    running: Arc<Mutex<bool>>,
527    cadence_reporter: Option<Arc<CadenceReporter>>,
528    #[allow(dead_code)] // retained for future direct-flush access
529    root: Arc<Mutex<ScheduleNode>>,
530    /// The scheduler thread handle — a dedicated `timing`-pool OS thread
531    /// (SRD-102), not a runtime task. Interior-mutable so the session host
532    /// can stop the scheduler through a shared `Arc<StopHandle>` (SRD-88 —
533    /// host owns the session-tier scheduler; executions only `report_frame`).
534    /// `take`n by whichever of `stop` / `drop` runs first; the other sees
535    /// `None` and is a no-op (idempotent).
536    task: Mutex<Option<std::thread::JoinHandle<()>>>,
537    /// Channel for async frame delivery — the executor sends frames here
538    /// instead of writing to reporters inline. The scheduler thread drains
539    /// this channel on each tick.
540    frame_tx: std::sync::mpsc::Sender<MetricSet>,
541    /// Wakes the timing thread out of its inter-tick Condvar wait so shutdown
542    /// is prompt instead of waiting out a base interval. Paired with the
543    /// `running` Mutex the thread waits on.
544    stop_cv: Arc<Condvar>,
545    /// Signalled by the task after its final flush; `stop()` waits on it
546    /// so the trailing window is committed before it returns.
547    done_rx: Mutex<Option<std::sync::mpsc::Receiver<()>>>,
548}
549
550impl StopHandle {
551    /// Stop the scheduler and join the capture thread. `&self` +
552    /// interior-mutable `thread` so a session host holding a shared
553    /// `Arc<StopHandle>` can stop it without sole ownership.
554    /// Idempotent — a second call (or `drop` after) sees `thread`
555    /// already taken and no-ops.
556    ///
557    /// ASYNC on purpose — the wait must YIELD to the runtime, never
558    /// block it. The timing thread's final flush `block_on`s
559    /// `CadenceReporter::shutdown_flush`, whose ack comes from the
560    /// reporter's OWNER — a tokio task on the caller's runtime. A
561    /// blocking wait here deadlocked a CURRENT-THREAD runtime three
562    /// ways: this (the only runtime) thread parked in `recv()`, the
563    /// timing thread parked awaiting the owner's ack, and the owner
564    /// task unable to run on the parked runtime. Multi-thread runtimes
565    /// escaped via `block_in_place` + spare workers, which is why only
566    /// single-threaded harnesses hung. The yield-poll below lets
567    /// same-runtime tasks progress on every flavor; this is a
568    /// session-end one-shot, so the 2 ms poll cadence touches no hot
569    /// path.
570    pub async fn stop(&self) {
571        *self.running.lock().unwrap_or_else(|e| e.into_inner()) = false;
572        self.stop_cv.notify_all(); // wake the inter-tick wait
573        // Wait for the task's final flush to land (the trailing window).
574        let done = self
575            .done_rx
576            .lock()
577            .unwrap_or_else(|e| e.into_inner())
578            .take();
579        if let Some(done) = done {
580            loop {
581                match done.try_recv() {
582                    Ok(()) | Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
583                    Err(std::sync::mpsc::TryRecvError::Empty) => {
584                        tokio::time::sleep(std::time::Duration::from_millis(2)).await;
585                    }
586                }
587            }
588        }
589        // The task has finished; drop its handle (no abort needed).
590        let _ = self.task.lock().unwrap_or_else(|e| e.into_inner()).take();
591    }
592
593    /// Reference to the cadence reporter, if any.
594    ///
595    /// Remains valid and queryable after the scheduler is stopped.
596    pub fn cadence_reporter(&self) -> Option<&Arc<CadenceReporter>> {
597        self.cadence_reporter.as_ref()
598    }
599
600    /// Deliver a frame to reporters asynchronously.
601    ///
602    /// The frame is enqueued on a channel and processed by the
603    /// scheduler thread on its next tick. This never blocks the
604    /// caller — safe to call from tokio worker threads.
605    pub fn report_frame(&self, snapshot: &MetricSet) {
606        let _ = self.frame_tx.send(snapshot.clone());
607    }
608}
609
610impl Drop for StopHandle {
611    fn drop(&mut self) {
612        *self.running.lock().unwrap_or_else(|e| e.into_inner()) = false;
613        self.stop_cv.notify_all(); // wake the inter-tick wait
614        let done = self
615            .done_rx
616            .lock()
617            .unwrap_or_else(|e| e.into_inner())
618            .take();
619        if let Some(done) = done {
620            // Drop is sync, so it can only WAIT where blocking is safe:
621            // OUTSIDE any tokio runtime, a plain blocking recv (the timing
622            // thread's `done` needs no progress from this thread — except
623            // through the CadenceReporter owner task, which lives on a
624            // runtime this thread is not part of). INSIDE a runtime,
625            // blocking would starve exactly the task the timing thread's
626            // final flush awaits (see `stop`'s doc), so DETACH instead:
627            // the timing thread completes its flush on its own; only the
628            // ordering guarantee ("flush landed before return") is
629            // forfeited, and a caller that needs that guarantee uses the
630            // async `stop()`.
631            if tokio::runtime::Handle::try_current().is_err() {
632                let _ = done.recv();
633            }
634        }
635        // The timing thread signals `done` as its last act, then returns —
636        // drop the join handle (detach); no abort exists for an OS thread and
637        // none is needed.
638        let _ = self.task.lock().unwrap_or_else(|e| e.into_inner()).take();
639    }
640}
641
642#[cfg(test)]
643mod tests {
644    use super::*;
645    use crate::snapshot::MetricValue;
646    use std::sync::atomic::{AtomicU64, Ordering};
647
648    struct CountingReporter {
649        count: Arc<AtomicU64>,
650    }
651
652    impl Reporter for CountingReporter {
653        fn report(&mut self, _snapshot: &MetricSet) {
654            self.count.fetch_add(1, Ordering::Relaxed);
655        }
656    }
657
658    fn mock_capture() -> Vec<(Labels, MetricSet)> {
659        let mut s = MetricSet::new(Duration::from_millis(100));
660        s.insert_counter("ops", Labels::default(), 10, Instant::now());
661        vec![(Labels::of("phase", "test"), s)]
662    }
663
664    fn empty_snapshot(interval: Duration) -> MetricSet {
665        MetricSet::new(interval)
666    }
667
668    #[test]
669    fn divergence_under_threshold_does_not_warn() {
670        let scheduled = Instant::now();
671        let actual = scheduled + Duration::from_millis(100); // < 250ms
672        assert!(
673            divergence_warning(
674                scheduled,
675                actual,
676                Duration::from_millis(250),
677                Duration::from_secs(60),
678                None,
679            )
680            .is_none()
681        );
682    }
683
684    #[test]
685    fn divergence_over_threshold_warns_first_time() {
686        let scheduled = Instant::now();
687        let actual = scheduled + Duration::from_millis(300); // > 250ms
688        let d = divergence_warning(
689            scheduled,
690            actual,
691            Duration::from_millis(250),
692            Duration::from_secs(60),
693            None,
694        );
695        assert_eq!(d, Some(Duration::from_millis(300)));
696    }
697
698    #[test]
699    fn divergence_is_rate_limited_within_window() {
700        let scheduled = Instant::now();
701        let actual = scheduled + Duration::from_millis(400);
702        // A warning fired 10s ago; the 60s window has not elapsed → suppressed.
703        let last_warn = Some(actual - Duration::from_secs(10));
704        assert!(
705            divergence_warning(
706                scheduled,
707                actual,
708                Duration::from_millis(250),
709                Duration::from_secs(60),
710                last_warn,
711            )
712            .is_none()
713        );
714        // Once the window elapses, it warns again.
715        let last_warn = Some(actual - Duration::from_secs(61));
716        assert!(
717            divergence_warning(
718                scheduled,
719                actual,
720                Duration::from_millis(250),
721                Duration::from_secs(60),
722                last_warn,
723            )
724            .is_some()
725        );
726    }
727
728    #[test]
729    fn scheduled_ts_is_stamped_on_tick_snapshots() {
730        // The scheduler stamps `scheduled_ts` on the combined snapshot each
731        // tick (captured_at is the actual fire instant). Verify the MetricSet
732        // carries the pair.
733        let mut s = MetricSet::new(Duration::from_millis(100));
734        assert!(s.scheduled_ts().is_none());
735        let sched = Instant::now();
736        s.set_scheduled_ts(sched);
737        assert_eq!(s.scheduled_ts(), Some(sched));
738        // actual_ts aliases captured_at.
739        assert_eq!(s.actual_ts(), s.captured_at());
740    }
741
742    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
743    async fn scheduler_builds_and_reports() {
744        let count = Arc::new(AtomicU64::new(0));
745        let c = count.clone();
746        let handle = SchedulerBuilder::new()
747            .base_interval(Duration::from_millis(100))
748            .add_reporter(Duration::from_millis(100), CountingReporter { count: c })
749            .build(Box::new(mock_capture));
750
751        let stop = handle.start();
752        tokio::time::sleep(Duration::from_millis(350)).await;
753        stop.stop().await;
754
755        let c = count.load(Ordering::Relaxed);
756        assert!((2..=5).contains(&c), "expected ~3 reports, got {c}");
757    }
758
759    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
760    async fn scheduler_feeds_cadence_reporter() {
761        use crate::cadence::{CadenceTree, Cadences};
762
763        let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_millis(100)]).unwrap());
764        let reporter = Arc::new(CadenceReporter::new(tree));
765        let handle = SchedulerBuilder::new()
766            .base_interval(Duration::from_millis(100))
767            .with_cadence_reporter(reporter.clone())
768            .build(Box::new(mock_capture));
769
770        let stop = handle.start();
771        tokio::time::sleep(Duration::from_millis(350)).await;
772        stop.stop().await;
773
774        // Reporter received ingests — has the component tracked.
775        let components = reporter.component_labels();
776        assert_eq!(components.len(), 1);
777        // The 100ms cadence should have at least one closed snapshot.
778        let component = &components[0];
779        let latest = reporter
780            .latest(component, Duration::from_millis(100))
781            .expect("cadence reporter should have a closed 100ms snapshot");
782        let ops_total = match latest
783            .family("ops")
784            .unwrap()
785            .metrics()
786            .next()
787            .unwrap()
788            .point()
789            .unwrap()
790            .value()
791        {
792            MetricValue::Counter(c) => c.cumulative,
793            _ => panic!("expected counter"),
794        };
795        assert_eq!(ops_total, 10, "one tick = 10");
796    }
797
798    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
799    async fn scheduler_coalesces_for_slow_reporter() {
800        let fast_count = Arc::new(AtomicU64::new(0));
801        let slow_count = Arc::new(AtomicU64::new(0));
802        let fc = fast_count.clone();
803        let sc = slow_count.clone();
804
805        let handle = SchedulerBuilder::new()
806            .base_interval(Duration::from_millis(50))
807            .add_reporter(Duration::from_millis(50), CountingReporter { count: fc })
808            .add_reporter(Duration::from_millis(200), CountingReporter { count: sc })
809            .build(Box::new(|| {
810                vec![(
811                    Labels::of("phase", "test"),
812                    empty_snapshot(Duration::from_millis(50)),
813                )]
814            }));
815
816        let stop = handle.start();
817        tokio::time::sleep(Duration::from_millis(450)).await;
818        stop.stop().await;
819
820        let fast = fast_count.load(Ordering::Relaxed);
821        let slow = slow_count.load(Ordering::Relaxed);
822        assert!(fast >= 6, "fast should get many reports, got {fast}");
823        assert!((1..=3).contains(&slow), "slow should get ~2, got {slow}");
824    }
825
826    /// With a CadenceTree installed, a slow reporter at the largest
827    /// declared cadence is fed *through* the chain (root → smallest
828    /// → … → largest). Functionally indistinguishable from the flat
829    /// arrangement at the consumer level — same number of reports,
830    /// same coalesced data — but internally the largest layer's
831    /// accumulation is bounded by the next-smaller cadence, not by
832    /// every base frame.
833    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
834    async fn scheduler_chained_tree_delivers_to_largest_cadence() {
835        use crate::cadence::{CadenceTree, Cadences};
836
837        let small_count = Arc::new(AtomicU64::new(0));
838        let large_count = Arc::new(AtomicU64::new(0));
839        let sc = small_count.clone();
840        let lc = large_count.clone();
841
842        // Cadences: 100ms (smallest declared) and 400ms (largest).
843        // Ratio 4 — well under default fan-in, no hidden inserts.
844        let tree = CadenceTree::plan_default(
845            Cadences::new(&[Duration::from_millis(100), Duration::from_millis(400)]).unwrap(),
846        );
847
848        let handle = SchedulerBuilder::new()
849            .base_interval(Duration::from_millis(100))
850            .with_cadence_tree(tree)
851            .add_reporter(Duration::from_millis(100), CountingReporter { count: sc })
852            .add_reporter(Duration::from_millis(400), CountingReporter { count: lc })
853            .build(Box::new(|| {
854                vec![(
855                    Labels::of("phase", "test"),
856                    empty_snapshot(Duration::from_millis(100)),
857                )]
858            }));
859
860        let stop = handle.start();
861        tokio::time::sleep(Duration::from_millis(900)).await;
862        stop.stop().await;
863
864        let small = small_count.load(Ordering::Relaxed);
865        let large = large_count.load(Ordering::Relaxed);
866        // ~9 base ticks → smallest fires every tick (≥6) and
867        // largest fires every 4 (≥1, ≤3).
868        assert!(small >= 6, "smallest cadence reports = {small}");
869        assert!(
870            (1..=3).contains(&large),
871            "largest cadence reports = {large}"
872        );
873    }
874
875    /// Hidden intermediate layers (auto-inserted by the planner)
876    /// participate in accumulation but never deliver to a reporter.
877    /// Verify that a reporter only at the *largest* declared cadence
878    /// still gets its expected report count even when a hidden
879    /// layer sits between it and the smallest cadence.
880    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
881    async fn scheduler_hidden_layers_pass_through_to_visible_reporters() {
882        use crate::cadence::{CadenceTree, Cadences};
883
884        let large_count = Arc::new(AtomicU64::new(0));
885        let lc = large_count.clone();
886
887        // 50ms → 1500ms is ratio 30 — exceeds default K=20, so the
888        // planner inserts a hidden intermediate. Ensures the chain
889        // flows through it correctly.
890        let tree = CadenceTree::plan_default(
891            Cadences::new(&[Duration::from_millis(50), Duration::from_millis(1500)]).unwrap(),
892        );
893        // Sanity check the planner actually inserted one.
894        let inserted: Vec<_> = tree.hidden().collect();
895        assert!(!inserted.is_empty(), "test relies on hidden insertion");
896
897        let handle = SchedulerBuilder::new()
898            .base_interval(Duration::from_millis(50))
899            .with_cadence_tree(tree)
900            .add_reporter(Duration::from_millis(1500), CountingReporter { count: lc })
901            .build(Box::new(|| {
902                vec![(
903                    Labels::of("phase", "test"),
904                    empty_snapshot(Duration::from_millis(50)),
905                )]
906            }));
907
908        let stop = handle.start();
909        tokio::time::sleep(Duration::from_millis(3300)).await;
910        stop.stop().await;
911
912        let large = large_count.load(Ordering::Relaxed);
913        assert!(large >= 1, "largest reporter saw 0 frames — chain broken");
914    }
915
916    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
917    async fn flush_component_routes_to_cadence_reporter() {
918        use crate::cadence::{CadenceTree, Cadences};
919
920        let tree = CadenceTree::plan_default(Cadences::new(&[Duration::from_secs(1)]).unwrap());
921        let reporter = Arc::new(CadenceReporter::new(tree));
922        let handle = SchedulerBuilder::new()
923            .with_cadence_reporter(reporter.clone())
924            .build(Box::new(Vec::new));
925
926        // Flush without starting — simulates lifecycle retirement
927        let labels = Labels::of("phase", "done");
928        let mut snapshot = MetricSet::new(Duration::from_secs(1));
929        snapshot.insert_counter("final_ops", Labels::default(), 42, Instant::now());
930        handle.flush_component(&labels, snapshot);
931        reporter.flush_for_tests();
932
933        // The flush went straight into the reporter's smallest
934        // cadence accumulator and promoted (interval matched).
935        let latest = reporter
936            .latest(&labels, Duration::from_secs(1))
937            .expect("flush should produce a closed snapshot");
938        assert!(latest.family("final_ops").is_some());
939    }
940}