joule_profiler_source_perf_event/
lib.rs1use std::{
10 collections::HashSet,
11 fs::File,
12 path::{Path, PathBuf},
13};
14
15use joule_profiler_core::{
16 sensor::{Sensor, Sensors},
17 source::MetricReader,
18 types::{Metric, Metrics},
19 unit::{MetricUnit, Unit, UnitPrefix},
20};
21use log::{debug, info, trace};
22
23use crate::{
24 config::PerfConfig,
25 error::PerfEventError,
26 event::{EVENTS, Event},
27 hardware::{PerfEventCounters, PerfEventHardware, Target},
28 snapshot::{Phase, Snapshot},
29};
30
31pub mod config;
32mod error;
33mod event;
34mod hardware;
35mod snapshot;
36
37type Result<T> = std::result::Result<T, PerfEventError>;
38
39const PERF_EVENT_METRIC_UNIT: MetricUnit = MetricUnit {
40 prefix: UnitPrefix::None,
41 unit: Unit::Count,
42};
43
44const CGROUP_ROOT: &str = "/sys/fs/cgroup";
47
48struct CgroupConfig {
49 name: PathBuf,
50 root: PathBuf,
51 cpu_spec: Option<HashSet<u32>>,
52}
53
54pub struct PerfEvent<H: PerfEventHardware = PerfEventCounters> {
63 hardware: H,
64 events: Vec<Event>,
65 cgroup_config: Option<CgroupConfig>,
66 begin_snapshot: Option<Snapshot>,
67 last_snapshot: Option<Snapshot>,
68}
69
70impl<H: PerfEventHardware + 'static> MetricReader for PerfEvent<H> {
71 type Type = Phase;
72 type Error = PerfEventError;
73 type Config = PerfConfig;
74
75 fn from_config(mut config: PerfConfig) -> Result<Self> {
76 let cgroup_config = if let Some(cgroup_name) = config.cgroup_name {
77 Some(CgroupConfig {
78 name: cgroup_name,
79 root: config.cgroup_root.take().unwrap_or(CGROUP_ROOT.into()),
80 cpu_spec: config.cpu_spec,
81 })
82 } else {
83 None
84 };
85
86 Ok(Self {
87 events: config
88 .events
89 .map_or(EVENTS.to_vec(), |e| e.into_iter().collect()),
90 cgroup_config,
91 hardware: H::default(),
92 begin_snapshot: None,
93 last_snapshot: None,
94 })
95 }
96
97 async fn pre_init(&mut self) -> Result<()> {
98 if let Some(cgroup_config) = &mut self.cgroup_config {
99 let path = Path::new(&cgroup_config.root).join(&cgroup_config.name);
100 info!(
101 "Initializing perf_event source for cgroup {}",
102 path.display()
103 );
104 let target = Target::Cgroup(File::open(path)?, cgroup_config.cpu_spec.take());
105 self.hardware.init_counters(&self.events, target).await?;
106 }
107 Ok(())
108 }
109
110 async fn init(&mut self, pid: i32) -> Result<()> {
113 if self.cgroup_config.is_none() {
114 info!("Initializing perf_event source for PID {pid}");
115 self.hardware
116 .init_counters(&self.events, Target::Pid(pid))
117 .await?;
118 }
119 Ok(())
120 }
121
122 async fn measure(&mut self) -> Result<()> {
124 trace!("Reading perf_event counters");
125 let new_snapshot = self.hardware.read_snapshot().await?;
126 if self.begin_snapshot.is_none() {
127 self.begin_snapshot = Some(new_snapshot);
128 } else {
129 self.last_snapshot = Some(new_snapshot);
130 }
131 Ok(())
132 }
133
134 async fn retrieve(&mut self) -> Result<Self::Type> {
136 if let Some(begin) = self.begin_snapshot.take()
137 && let Some(end) = self.last_snapshot.take()
138 {
139 self.begin_snapshot = Some(end.clone());
140 Ok(Phase { begin, end })
141 } else {
142 Err(PerfEventError::NotEnoughSamples)
143 }
144 }
145
146 fn get_sensors(&self) -> Result<Sensors> {
148 trace!("Building perf_event sensor list");
149 let sensors: Sensors = self
150 .events
151 .iter()
152 .map(|event| {
153 trace!("Registering sensor: {event}");
154 Sensor::new(*event, PERF_EVENT_METRIC_UNIT, Self::get_name())
155 })
156 .collect();
157
158 debug!("Registered {} perf_event sensors", sensors.len());
159 Ok(sensors)
160 }
161
162 fn to_metrics(&self, result: Self::Type) -> Result<Metrics> {
167 trace!(
168 "Converting {} counters to metrics",
169 result.begin.metrics.len()
170 );
171 let diff = result.diff();
172 Ok(diff
173 .metrics
174 .into_iter()
175 .map(|(event, per_cpu)| {
176 let value: u64 = per_cpu.values().sum();
177 Metric::new(event, value, PERF_EVENT_METRIC_UNIT, Self::get_name())
178 })
179 .collect())
180 }
181
182 fn get_name() -> &'static str {
183 "perf_event"
184 }
185
186 fn get_id() -> &'static str {
187 "perf"
188 }
189}
190
191#[cfg(test)]
192mod tests {
193 use joule_profiler_core::types::MetricValue;
194
195 use super::*;
196 use crate::{event::Event, hardware::MockPerfEventHardware, snapshot::Snapshot};
197
198 fn snapshot(entries: Vec<(Event, u64)>) -> Snapshot {
200 Snapshot {
201 metrics: entries
202 .into_iter()
203 .map(|(event, value)| (event, std::collections::HashMap::from([(0, value)])))
204 .collect(),
205 }
206 }
207
208 fn total(snapshot: &Snapshot, event: Event) -> u64 {
209 snapshot.metrics[&event].values().sum()
210 }
211
212 fn with_hardware(hardware: MockPerfEventHardware) -> PerfEvent<MockPerfEventHardware> {
213 PerfEvent {
214 hardware,
215 events: EVENTS.to_vec(),
216 cgroup_config: None,
217 begin_snapshot: None,
218 last_snapshot: None,
219 }
220 }
221
222 #[tokio::test]
223 async fn measure_stores_begin_snapshot() {
224 let mut hardware = MockPerfEventHardware::new();
225 hardware
226 .expect_read_snapshot()
227 .returning(|| Box::pin(async { Ok(snapshot(vec![(Event::CpuCycles, 100)])) }));
228
229 let mut source = with_hardware(hardware);
230 source.measure().await.unwrap();
231
232 assert!(source.begin_snapshot.is_some());
233 assert!(source.last_snapshot.is_none());
234 }
235
236 #[tokio::test]
237 async fn measure_twice_stores_last_snapshot() {
238 let mut hardware = MockPerfEventHardware::new();
239 let mut read_snapshot_call_count = 0u64;
240 hardware.expect_read_snapshot().returning(move || {
241 read_snapshot_call_count += 1;
242 Box::pin(async move {
243 Ok(snapshot(vec![(
244 Event::CpuCycles,
245 read_snapshot_call_count * 100,
246 )]))
247 })
248 });
249
250 let mut source = with_hardware(hardware);
251 source.measure().await.unwrap();
252 source.measure().await.unwrap();
253
254 assert!(source.begin_snapshot.is_some());
255 assert!(source.last_snapshot.is_some());
256 }
257
258 #[tokio::test]
259 async fn retrieve_without_enough_snapshots_returns_error() {
260 let mut hardware = MockPerfEventHardware::new();
261 hardware
262 .expect_read_snapshot()
263 .returning(|| Box::pin(async { Ok(snapshot(vec![(Event::CpuCycles, 100)])) }));
264
265 let mut source = with_hardware(hardware);
266 source.measure().await.unwrap();
267
268 assert!(matches!(
269 source.retrieve().await,
270 Err(PerfEventError::NotEnoughSamples)
271 ));
272 }
273
274 #[tokio::test]
275 async fn retrieve_returns_correct_phase() {
276 let mut hardware = MockPerfEventHardware::new();
277 let mut read_snapshot_call_count = 0u64;
278 hardware.expect_read_snapshot().returning(move || {
279 read_snapshot_call_count += 1;
280 Box::pin(async move {
281 Ok(snapshot(vec![(
282 Event::CpuCycles,
283 read_snapshot_call_count * 100,
284 )]))
285 })
286 });
287
288 let mut source = with_hardware(hardware);
289 source.measure().await.unwrap();
290 source.measure().await.unwrap();
291 let phase = source.retrieve().await.unwrap();
292
293 assert_eq!(total(&phase.begin, Event::CpuCycles), 100);
294 assert_eq!(total(&phase.end, Event::CpuCycles), 200);
295 }
296
297 #[tokio::test]
298 async fn retrieve_rolls_begin_snapshot_to_end() {
299 let mut hardware = MockPerfEventHardware::new();
300 let mut read_snapshot_call_count = 0u64;
301 hardware.expect_read_snapshot().returning(move || {
302 read_snapshot_call_count += 1;
303 Box::pin(async move {
304 Ok(snapshot(vec![(
305 Event::CpuCycles,
306 read_snapshot_call_count * 100,
307 )]))
308 })
309 });
310
311 let mut source = with_hardware(hardware);
312 source.measure().await.unwrap();
313 source.measure().await.unwrap();
314 source.retrieve().await.unwrap();
315 assert_eq!(
316 total(source.begin_snapshot.as_ref().unwrap(), Event::CpuCycles),
317 200
318 );
319 assert!(source.last_snapshot.is_none());
320 }
321
322 #[tokio::test]
323 async fn to_metrics_returns_correct_values() {
324 let mut hardware = MockPerfEventHardware::new();
325 let mut read_snapshot_call_count = 0;
326 hardware.expect_read_snapshot().returning(move || {
327 read_snapshot_call_count += 1;
328 Box::pin(async move {
329 Ok(match read_snapshot_call_count {
330 1 => snapshot(vec![(Event::CpuCycles, 0)]),
331 _ => snapshot(vec![(Event::CpuCycles, 500)]),
332 })
333 })
334 });
335
336 let mut source = with_hardware(hardware);
337 source.measure().await.unwrap();
338 source.measure().await.unwrap();
339 let phase = source.retrieve().await.unwrap();
340 let metrics = source.to_metrics(phase).unwrap();
341 let cycles = metrics
342 .iter()
343 .find(|m| m.name == Event::CpuCycles.to_string())
344 .unwrap();
345
346 assert_eq!(cycles.value, MetricValue::UnsignedInteger(500));
347 assert_eq!(cycles.unit, PERF_EVENT_METRIC_UNIT);
348 }
349}