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}