Skip to main content

scx_cosmos/
stats.rs

1use std::io::Write;
2use std::sync::atomic::AtomicBool;
3use std::sync::atomic::Ordering;
4use std::sync::Arc;
5use std::time::Duration;
6
7use anyhow::Result;
8use scx_stats::prelude::*;
9use scx_stats_derive::stat_doc;
10use scx_stats_derive::Stats;
11use serde::Deserialize;
12use serde::Serialize;
13
14#[stat_doc]
15#[derive(Clone, Debug, Default, Serialize, Deserialize, Stats)]
16#[stat(top)]
17pub struct Metrics {
18    #[stat(desc = "Direct dispatch due to high perf events (migration)")]
19    pub nr_event_dispatches: u64,
20    #[stat(desc = "Kept on same CPU due to perf sticky threshold")]
21    pub nr_ev_sticky_dispatches: u64,
22    #[stat(desc = "Direct dispatch due to GPU affinity")]
23    pub nr_gpu_dispatches: u64,
24    #[stat(desc = "Primary domain overload events (spill to non-primary CPUs)")]
25    pub nr_overload_events: u64,
26}
27
28impl Metrics {
29    fn format<W: Write>(&self, w: &mut W) -> Result<()> {
30        writeln!(
31            w,
32            "[{}] ev_dispatches={} ev_sticky_dispatches={} gpu_dispatches={} overload_events={}",
33            crate::SCHEDULER_NAME,
34            self.nr_event_dispatches,
35            self.nr_ev_sticky_dispatches,
36            self.nr_gpu_dispatches,
37            self.nr_overload_events,
38        )?;
39        Ok(())
40    }
41
42    fn delta(&self, rhs: &Self) -> Self {
43        Self {
44            nr_event_dispatches: self.nr_event_dispatches - rhs.nr_event_dispatches,
45            nr_ev_sticky_dispatches: self.nr_ev_sticky_dispatches - rhs.nr_ev_sticky_dispatches,
46            nr_gpu_dispatches: self.nr_gpu_dispatches - rhs.nr_gpu_dispatches,
47            nr_overload_events: self.nr_overload_events - rhs.nr_overload_events,
48        }
49    }
50}
51
52pub fn server_data() -> StatsServerData<(), Metrics> {
53    let open: Box<dyn StatsOpener<(), Metrics>> = Box::new(move |(req_ch, res_ch)| {
54        req_ch.send(())?;
55        let mut prev = res_ch.recv()?;
56
57        let read: Box<dyn StatsReader<(), Metrics>> = Box::new(move |_args, (req_ch, res_ch)| {
58            req_ch.send(())?;
59            let cur = res_ch.recv()?;
60            let delta = cur.delta(&prev);
61            prev = cur;
62            delta.to_json()
63        });
64
65        Ok(read)
66    });
67
68    StatsServerData::new()
69        .add_meta(Metrics::meta())
70        .add_ops("top", StatsOps { open, close: None })
71}
72
73pub fn monitor(intv: Duration, shutdown: Arc<AtomicBool>) -> Result<()> {
74    scx_utils::monitor_stats::<Metrics>(
75        &[],
76        intv,
77        || shutdown.load(Ordering::Relaxed),
78        |metrics| metrics.format(&mut std::io::stdout()),
79    )
80}