1use crate::observation::source;
38use crate::observation::{
39 BindContext, Perturbation, Sensor, SensorDescriptor, SensorId, SensorOutcome, StateWriter,
40 Uncertainty,
41};
42use corescout_core::error::{Error, Result};
43use corescout_mirror::entity::keys;
44use corescout_mirror::state::{ChannelId, Semantics, Unit};
45use std::path::PathBuf;
46
47const TIME_FIELDS: [&str; 8] = [
49 "user", "nice", "system", "idle", "iowait", "irq", "softirq", "steal",
50];
51
52pub struct SchedulerSensor {
53 stat_path: PathBuf,
54 schedstat_path: PathBuf,
55 loadavg_path: PathBuf,
56 cpu_rows: Vec<(u32, u32)>,
58 machine_row: Option<u32>,
59 ns_per_tick: f64,
61 time_channels: Vec<ChannelId>,
62 channel_run: Option<ChannelId>,
63 channel_wait: Option<ChannelId>,
64 channel_slices: Option<ChannelId>,
65 channel_runnable: Option<ChannelId>,
66 channel_tasks: Option<ChannelId>,
67}
68
69impl SchedulerSensor {
70 pub fn new() -> SchedulerSensor {
71 SchedulerSensor {
72 stat_path: PathBuf::new(),
73 schedstat_path: PathBuf::new(),
74 loadavg_path: PathBuf::new(),
75 cpu_rows: Vec::new(),
76 machine_row: None,
77 ns_per_tick: 0.0,
78 time_channels: Vec::new(),
79 channel_run: None,
80 channel_wait: None,
81 channel_slices: None,
82 channel_runnable: None,
83 channel_tasks: None,
84 }
85 }
86
87 fn row_of_cpu(&self, cpu: u32) -> Option<u32> {
88 self.cpu_rows
89 .iter()
90 .find(|(c, _)| *c == cpu)
91 .map(|(_, r)| *r)
92 }
93}
94
95impl Default for SchedulerSensor {
96 fn default() -> Self {
97 Self::new()
98 }
99}
100
101impl Sensor for SchedulerSensor {
102 fn descriptor(&self) -> SensorDescriptor {
103 SensorDescriptor {
104 id: SensorId(5),
105 key: "scheduler",
106 physical_fact: "cumulative CPU time by class, runqueue wait time, and the number \
107 of tasks runnable right now",
108 source: "/proc/stat, /proc/schedstat, /proc/loadavg",
109 max_rate_hz: 100.0,
110 perturbation: Perturbation::Negligible,
111 uncertainty: Uncertainty::absolute(
112 1.0,
113 "/proc/stat is accounted per timer tick, so it attributes a whole tick to \
114 whatever was running when the tick fired; sub-tick work is misattributed",
115 ),
116 requires_privilege: false,
117 }
118 }
119
120 fn bind(&mut self, ctx: &mut BindContext<'_>) -> Result<()> {
121 let proc_root = ctx.substrate().roots.proc.clone();
122 self.stat_path = proc_root.join("stat");
123 self.schedstat_path = proc_root.join("schedstat");
124 self.loadavg_path = proc_root.join("loadavg");
125
126 if source::string(&self.stat_path).is_none() {
127 return Err(Error::unsupported("/proc/stat is not readable"));
128 }
129
130 self.ns_per_tick = 1_000_000_000.0 / source::clock_ticks_per_second() as f64;
131 self.machine_row = ctx.row_of("machine");
132 for cpu in ctx
133 .substrate()
134 .topology
135 .logical_cpus
136 .iter()
137 .filter(|c| c.online)
138 {
139 if let Some(row) = ctx.row_of(&keys::logical_cpu(cpu.id)) {
140 self.cpu_rows.push((cpu.id, row));
141 }
142 }
143
144 for field in TIME_FIELDS {
145 self.time_channels.push(ctx.declare_channel(
146 format!("cpu.time.{field}"),
147 Unit::Nanosecond,
148 Semantics::Cumulative,
149 ));
150 }
151 self.channel_run = Some(ctx.declare_channel(
152 "cpu.sched.run_time",
153 Unit::Nanosecond,
154 Semantics::Cumulative,
155 ));
156 self.channel_wait = Some(ctx.declare_channel(
157 "cpu.sched.wait_time",
158 Unit::Nanosecond,
159 Semantics::Cumulative,
160 ));
161 self.channel_slices =
162 Some(ctx.declare_channel("cpu.sched.timeslices", Unit::Count, Semantics::Cumulative));
163 self.channel_runnable =
164 Some(ctx.declare_channel("machine.tasks.runnable", Unit::Count, Semantics::Instant));
165 self.channel_tasks =
166 Some(ctx.declare_channel("machine.tasks.total", Unit::Count, Semantics::Instant));
167 Ok(())
168 }
169
170 fn observe(&mut self, out: &mut StateWriter<'_>) -> SensorOutcome {
171 let mut outcome = SensorOutcome::default();
172
173 match source::string(&self.stat_path) {
174 Some(text) => {
175 for (cpu, times) in parse_proc_stat(&text) {
176 let Some(row) = self.row_of_cpu(cpu) else {
177 continue;
178 };
179 for (index, ticks) in times.iter().enumerate() {
180 source::emit(
184 out,
185 &mut outcome,
186 row,
187 self.time_channels.get(index).copied(),
188 *ticks as f64 * self.ns_per_tick,
189 );
190 }
191 }
192 }
193 None => outcome.error(),
194 }
195
196 if let Some(text) = source::string(&self.schedstat_path) {
198 for (cpu, run_ns, wait_ns, slices) in parse_schedstat(&text) {
199 let Some(row) = self.row_of_cpu(cpu) else {
200 continue;
201 };
202 source::emit(out, &mut outcome, row, self.channel_run, run_ns as f64);
203 source::emit(out, &mut outcome, row, self.channel_wait, wait_ns as f64);
204 source::emit(out, &mut outcome, row, self.channel_slices, slices as f64);
205 }
206 }
207
208 if let (Some(machine), Some(text)) = (self.machine_row, source::string(&self.loadavg_path))
209 {
210 if let Some((runnable, total)) = parse_loadavg_tasks(&text) {
211 source::emit(
212 out,
213 &mut outcome,
214 machine,
215 self.channel_runnable,
216 runnable as f64,
217 );
218 source::emit(out, &mut outcome, machine, self.channel_tasks, total as f64);
219 }
220 }
221
222 outcome
223 }
224}
225
226fn parse_proc_stat(text: &str) -> Vec<(u32, [u64; 8])> {
231 let mut out = Vec::new();
232 for line in text.lines() {
233 let Some(rest) = line.strip_prefix("cpu") else {
234 continue;
235 };
236 if !rest.starts_with(|c: char| c.is_ascii_digit()) {
241 continue;
242 }
243 let mut fields = rest.split_whitespace();
244 let Some(index) = fields.next() else { continue };
245 let Ok(cpu) = index.parse::<u32>() else {
246 continue;
247 };
248 let mut times = [0u64; 8];
249 for slot in times.iter_mut() {
250 *slot = fields.next().and_then(|v| v.parse().ok()).unwrap_or(0);
251 }
252 out.push((cpu, times));
253 }
254 out
255}
256
257fn parse_schedstat(text: &str) -> Vec<(u32, u64, u64, u64)> {
264 let mut out = Vec::new();
265 for line in text.lines() {
266 let Some(rest) = line.strip_prefix("cpu") else {
267 continue;
268 };
269 if !rest.starts_with(|c: char| c.is_ascii_digit()) {
270 continue;
271 }
272 let mut fields = rest.split_whitespace();
273 let Some(index) = fields.next() else { continue };
274 let Ok(cpu) = index.parse::<u32>() else {
275 continue;
276 };
277 let values: Vec<u64> = fields.filter_map(|v| v.parse().ok()).collect();
278 if values.len() < 3 {
279 continue;
280 }
281 let tail = &values[values.len() - 3..];
282 out.push((cpu, tail[0], tail[1], tail[2]));
283 }
284 out
285}
286
287fn parse_loadavg_tasks(text: &str) -> Option<(u64, u64)> {
292 let field = text.split_whitespace().nth(3)?;
293 let (runnable, total) = field.split_once('/')?;
294 Some((runnable.parse().ok()?, total.parse().ok()?))
295}
296
297#[cfg(test)]
298mod tests {
299 use super::*;
300
301 const PROC_STAT: &str = "\
302cpu 100 200 300 400 500 600 700 800 0 0
303cpu0 1 2 3 4 5 6 7 8 0 0
304cpu1 10 20 30 40 50 60 70 80 0 0
305intr 12345 1 2 3
306ctxt 987654
307btime 1700000000
308processes 4242
309procs_running 3
310procs_blocked 0
311";
312
313 const SCHEDSTAT: &str = "\
314version 15
315timestamp 4294900000
316cpu0 0 0 0 0 0 0 123456789 987654321 4242
317domain0 00000000,00000003 0 0 0
318cpu1 0 0 0 0 0 0 111 222 333
319";
320
321 #[test]
322 fn proc_stat_yields_per_cpu_times_and_skips_the_aggregate() {
323 let parsed = parse_proc_stat(PROC_STAT);
324 assert_eq!(parsed.len(), 2, "the aggregate `cpu` line must be skipped");
325 assert!(
326 !parsed.iter().any(|(cpu, _)| *cpu == 100),
327 "the aggregate first time field must not be read as a CPU number"
328 );
329 assert_eq!(parsed[0].0, 0);
330 assert_eq!(parsed[0].1, [1, 2, 3, 4, 5, 6, 7, 8]);
331 assert_eq!(parsed[1].0, 1);
332 assert_eq!(parsed[1].1[3], 40, "idle time");
333 }
334
335 #[test]
336 fn proc_stat_tolerates_a_short_line() {
337 let parsed = parse_proc_stat("cpu0 1 2 3 4\n");
340 assert_eq!(parsed[0].1, [1, 2, 3, 4, 0, 0, 0, 0]);
341 }
342
343 #[test]
344 fn schedstat_takes_the_last_three_fields() {
345 let parsed = parse_schedstat(SCHEDSTAT);
346 assert_eq!(parsed.len(), 2, "domain lines must be ignored");
347 assert_eq!(parsed[0], (0, 123_456_789, 987_654_321, 4242));
348 assert_eq!(parsed[1], (1, 111, 222, 333));
349 }
350
351 #[test]
352 fn schedstat_survives_a_version_with_extra_leading_fields() {
353 let text = "cpu0 9 9 9 9 9 9 9 9 9 7 8 9\n";
355 assert_eq!(parse_schedstat(text), vec![(0, 7, 8, 9)]);
356 }
357
358 #[test]
359 fn loadavg_gives_the_instantaneous_counts_only() {
360 let parsed = parse_loadavg_tasks("0.52 0.61 0.70 3/1234 5678\n");
361 assert_eq!(parsed, Some((3, 1234)));
362 }
363
364 #[test]
365 fn loadavg_averages_are_not_extracted() {
366 let text = "0.52 0.61 0.70 3/1234 5678\n";
369 let parsed = parse_loadavg_tasks(text).unwrap();
370 assert_ne!(parsed.0, 0, "runnable count is present");
371 assert!(
372 !TIME_FIELDS.iter().any(|f| f.contains("load")),
373 "no channel should carry a load average"
374 );
375 }
376
377 #[test]
378 fn malformed_input_yields_nothing_rather_than_panicking() {
379 assert!(parse_proc_stat("").is_empty());
380 assert!(parse_schedstat("cpu\n").is_empty());
381 assert_eq!(parse_loadavg_tasks("garbage"), None);
382 assert_eq!(parse_loadavg_tasks("1 2 3 notafraction 5"), None);
383 }
384}