corescout_substrate/observation/
interrupts.rs1use crate::observation::source;
42use crate::observation::{
43 BindContext, Perturbation, Sensor, SensorDescriptor, SensorId, SensorOutcome, StateWriter,
44 Uncertainty,
45};
46use corescout_core::error::{Error, Result};
47use corescout_mirror::entity::keys;
48use corescout_mirror::state::{ChannelId, Semantics, Unit};
49use std::path::PathBuf;
50
51pub struct InterruptSensor {
52 path: PathBuf,
53 columns: Vec<(usize, u32)>,
55 channel_total: Option<ChannelId>,
56 channel_sources: Option<ChannelId>,
57}
58
59impl InterruptSensor {
60 pub fn new() -> InterruptSensor {
61 InterruptSensor {
62 path: PathBuf::new(),
63 columns: Vec::new(),
64 channel_total: None,
65 channel_sources: None,
66 }
67 }
68}
69
70impl Default for InterruptSensor {
71 fn default() -> Self {
72 Self::new()
73 }
74}
75
76impl Sensor for InterruptSensor {
77 fn descriptor(&self) -> SensorDescriptor {
78 SensorDescriptor {
79 id: SensorId(6),
80 key: "interrupts",
81 physical_fact: "the number of hardware interrupts each CPU has serviced since boot",
82 source: "/proc/interrupts",
83 max_rate_hz: 50.0,
87 perturbation: Perturbation::Negligible,
88 uncertainty: Uncertainty::unknown(
89 "counts are exact; which CPU a future interrupt lands on depends on IRQ \
90 affinity and on irqbalance, either of which may change between samples",
91 ),
92 requires_privilege: false,
93 }
94 }
95
96 fn bind(&mut self, ctx: &mut BindContext<'_>) -> Result<()> {
97 self.path = ctx.substrate().roots.proc.join("interrupts");
98 let Some(text) = source::string(&self.path) else {
99 return Err(Error::unsupported("/proc/interrupts is not readable"));
100 };
101
102 let cpus = parse_header(&text);
107 if cpus.is_empty() {
108 return Err(Error::unsupported(
109 "/proc/interrupts has no per-CPU columns",
110 ));
111 }
112 for (column, cpu) in cpus.iter().enumerate() {
113 if let Some(row) = ctx.row_of(&keys::logical_cpu(*cpu)) {
114 self.columns.push((column, row));
115 }
116 }
117
118 self.channel_total =
119 Some(ctx.declare_channel("cpu.interrupts.total", Unit::Count, Semantics::Cumulative));
120 self.channel_sources = Some(ctx.declare_channel(
121 "cpu.interrupts.active_sources",
122 Unit::Count,
123 Semantics::Instant,
124 ));
125 Ok(())
126 }
127
128 fn observe(&mut self, out: &mut StateWriter<'_>) -> SensorOutcome {
129 let mut outcome = SensorOutcome::default();
130 let Some(text) = source::string(&self.path) else {
131 outcome.error();
132 return outcome;
133 };
134
135 let width = self.columns.iter().map(|(c, _)| *c + 1).max().unwrap_or(0);
136 let (totals, sources) = sum_columns(&text, width);
137
138 for (column, row) in &self.columns {
139 source::emit(
140 out,
141 &mut outcome,
142 *row,
143 self.channel_total,
144 totals[*column] as f64,
145 );
146 source::emit(
147 out,
148 &mut outcome,
149 *row,
150 self.channel_sources,
151 sources[*column] as f64,
152 );
153 }
154 outcome
155 }
156}
157
158fn parse_header(text: &str) -> Vec<u32> {
160 let Some(header) = text.lines().next() else {
161 return Vec::new();
162 };
163 header
164 .split_whitespace()
165 .filter_map(|token| token.strip_prefix("CPU")?.parse::<u32>().ok())
166 .collect()
167}
168
169fn sum_columns(text: &str, width: usize) -> (Vec<u64>, Vec<u32>) {
173 let mut totals = vec![0u64; width];
174 let mut sources = vec![0u32; width];
175
176 for line in text.lines().skip(1) {
177 let Some((_, rest)) = line.split_once(':') else {
180 continue;
181 };
182 for (column, token) in rest.split_whitespace().take(width).enumerate() {
183 let Ok(count) = token.parse::<u64>() else {
187 break;
188 };
189 totals[column] = totals[column].saturating_add(count);
190 if count > 0 {
191 sources[column] += 1;
192 }
193 }
194 }
195 (totals, sources)
196}
197
198#[cfg(test)]
199mod tests {
200 use super::*;
201
202 const INTERRUPTS: &str = "\
203 CPU0 CPU1 CPU2 CPU3
204 0: 31 0 0 0 IO-APIC 2-edge timer
205 8: 1 0 0 0 IO-APIC 8-edge rtc0
206 9: 0 0 0 0 IO-APIC 9-fasteoi acpi
207 16: 1523 0 0 0 IO-APIC 16-fasteoi ehci_hcd:usb1
208124: 0 88192 0 0 PCI-MSI 524288-edge eth0-rx-0
209NMI: 2 2 2 2 Non-maskable interrupts
210LOC: 9812345 8712345 7612345 6512345 Local timer interrupts
211TLB: 104 233 311 498 TLB shootdowns
212";
213
214 #[test]
215 fn the_header_names_the_cpu_columns() {
216 assert_eq!(parse_header(INTERRUPTS), vec![0, 1, 2, 3]);
217 }
218
219 #[test]
220 fn a_sparse_cpu_set_is_read_from_the_header_not_assumed() {
221 let text = " CPU0 CPU3\n 0: 5 7 IO-APIC timer\n";
225 assert_eq!(parse_header(text), vec![0, 3]);
226 }
227
228 #[test]
229 fn totals_are_summed_across_every_source() {
230 let (totals, _) = sum_columns(INTERRUPTS, 4);
231 assert_eq!(totals[0], 31 + 1 + 1523 + 2 + 9_812_345 + 104);
233 assert_eq!(totals[1], 88_192 + 2 + 8_712_345 + 233);
235 assert!(totals[1] < totals[0]);
236 }
237
238 #[test]
239 fn active_sources_counts_only_sources_that_have_fired() {
240 let (_, sources) = sum_columns(INTERRUPTS, 4);
241 assert_eq!(sources[0], 6);
243 assert_eq!(sources[2], 3);
245 }
246
247 #[test]
248 fn trailing_description_columns_are_not_counted_as_interrupts() {
249 let text = " CPU0\n 0: 31 IO-APIC 2-edge timer\n";
252 let (totals, _) = sum_columns(text, 1);
253 assert_eq!(totals[0], 31);
254 }
255
256 #[test]
257 fn malformed_input_is_survivable() {
258 assert!(parse_header("").is_empty());
259 let (totals, sources) = sum_columns("garbage with no colon\n", 2);
260 assert_eq!(totals, vec![0, 0]);
261 assert_eq!(sources, vec![0, 0]);
262 }
263
264 #[test]
265 fn a_row_shorter_than_the_cpu_count_does_not_panic() {
266 let text = " CPU0 CPU1\n 0: 5 7 IO-APIC timer\nERR: 3\n";
269 let (totals, _) = sum_columns(text, 2);
270 assert_eq!(totals[0], 8);
271 assert_eq!(totals[1], 7);
272 }
273}