1use crate::extract::ExtractSchedUtilOpts;
7use crate::process::PerfSchedScriptRecord;
8use anyhow::{Context as _, Result};
9use serde::{Deserialize, Serialize};
10use serde_json::{json, Value};
11use std::collections::{BTreeSet, HashMap, HashSet};
12use std::fs::File;
13use std::io::{BufRead, BufReader};
14
15const DEFAULT_SYSTEM_CATEGORY_NAMES: [&str; 5] = [
16 "hardirq",
17 "softirq-rx",
18 "softirq-tx",
19 "softirq-other",
20 "nmi",
21];
22const HARDIRQ_CATEGORY: &str = "hardirq";
23const SOFTIRQ_RX_CATEGORY: &str = "softirq-rx";
24const SOFTIRQ_TX_CATEGORY: &str = "softirq-tx";
25const SOFTIRQ_OTHER_CATEGORY: &str = "softirq-other";
26const NMI_CATEGORY: &str = "nmi";
27
28#[derive(Debug, Clone, Serialize, Deserialize)]
29struct BusyUtilRecord {
30 time_ms: u64,
31 window_ms: u64,
32 cpu_count: usize,
33 total: f64,
34 uncategorized: f64,
35 categories: HashMap<String, f64>,
36}
37
38#[derive(Debug, Clone)]
39struct BusyInterval {
40 start_ns: u64,
41 end_ns: u64,
42 comm: String,
43 hint: u64,
44}
45
46#[derive(Debug, Clone)]
47enum CategoryCommMatcher {
48 Exact(String),
49 Glob(String),
50}
51
52#[derive(Debug, Clone)]
53struct CategorySpec {
54 name: String,
55 comm_matcher: CategoryCommMatcher,
56 hint: Option<u64>,
57}
58
59#[derive(Debug, Default)]
60struct CompiledCategoryMatcher {
61 exact_any_hint: HashMap<String, Vec<usize>>,
62 exact_by_hint: HashMap<String, HashMap<u64, Vec<usize>>>,
63 glob_specs: Vec<GlobCategory>,
64}
65
66#[derive(Debug)]
67struct GlobCategory {
68 index: usize,
69 pattern: String,
70 hint: Option<u64>,
71}
72
73#[derive(Debug, Clone)]
74struct RunningTask {
75 tid: i32,
76 comm: String,
77 hint: u64,
78}
79
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81enum ActiveSystem {
82 HardIrq,
83 SoftIrqRx,
84 SoftIrqTx,
85 SoftIrqOther,
86}
87
88impl ActiveSystem {
89 fn category_name(self) -> &'static str {
90 match self {
91 Self::HardIrq => HARDIRQ_CATEGORY,
92 Self::SoftIrqRx => SOFTIRQ_RX_CATEGORY,
93 Self::SoftIrqTx => SOFTIRQ_TX_CATEGORY,
94 Self::SoftIrqOther => SOFTIRQ_OTHER_CATEGORY,
95 }
96 }
97
98 fn is_softirq(self) -> bool {
99 matches!(self, Self::SoftIrqRx | Self::SoftIrqTx | Self::SoftIrqOther)
100 }
101}
102
103#[derive(Debug, Clone, PartialEq, Eq)]
104enum ActiveExecution {
105 Task { comm: String, hint: u64 },
106 System(ActiveSystem),
107}
108
109impl ActiveExecution {
110 fn to_interval(&self, start_ns: u64, end_ns: u64) -> Option<BusyInterval> {
111 if end_ns <= start_ns {
112 return None;
113 }
114
115 let (comm, hint) = match self {
116 Self::Task { comm, hint } => (comm.clone(), *hint),
117 Self::System(system) => (system.category_name().to_string(), 0),
118 };
119
120 Some(BusyInterval {
121 start_ns,
122 end_ns,
123 comm,
124 hint,
125 })
126 }
127}
128
129#[derive(Debug, Default)]
130struct CpuState {
131 segment_start_ns: u64,
132 last_seen_ns: u64,
133 running_task: Option<RunningTask>,
134 system_stack: Vec<ActiveSystem>,
135 active_execution: Option<ActiveExecution>,
136}
137
138#[derive(Debug)]
139enum ClassifiedEvent {
140 SchedSwitch {
141 next_comm: String,
142 next_hint: u64,
143 next_tid: i32,
144 next_is_idle: bool,
145 },
146 HardIrqEntry,
147 HardIrqExit,
148 SoftIrqEntry(ActiveSystem),
149 SoftIrqExit,
150 Nmi {
151 delta_ns: u64,
152 },
153 Other,
154}
155
156trait BusyIntervalSink {
157 fn on_interval(&mut self, interval: &BusyInterval);
158}
159
160#[derive(Debug)]
161struct BucketAggregator {
162 trace_start_ns: u64,
163 window_ns: u64,
164 total_busy_ns: Vec<u64>,
165 uncategorized_busy_ns: Vec<u64>,
166 category_busy_ns: Vec<Vec<u64>>,
167 categories: Vec<CategorySpec>,
168 matcher: CompiledCategoryMatcher,
169 interval_count: usize,
170}
171
172impl BucketAggregator {
173 fn new(trace_start_ns: u64, window_ns: u64, categories: Vec<CategorySpec>) -> Self {
174 let category_busy_ns = vec![Vec::new(); categories.len()];
175 let matcher = compile_category_matcher(&categories);
176 Self {
177 trace_start_ns,
178 window_ns,
179 total_busy_ns: Vec::new(),
180 uncategorized_busy_ns: Vec::new(),
181 category_busy_ns,
182 categories,
183 matcher,
184 interval_count: 0,
185 }
186 }
187
188 fn ensure_bucket(&mut self, idx: usize) {
189 let required_len = idx + 1;
190 if self.total_busy_ns.len() >= required_len {
191 return;
192 }
193
194 self.total_busy_ns.resize(required_len, 0);
195 self.uncategorized_busy_ns.resize(required_len, 0);
196 for values in &mut self.category_busy_ns {
197 values.resize(required_len, 0);
198 }
199 }
200
201 fn finalize(
202 self,
203 trace_end_ns: u64,
204 cpu_count: usize,
205 window_ms: u64,
206 ) -> Result<(Vec<BusyUtilRecord>, usize)> {
207 if trace_end_ns <= self.trace_start_ns {
208 anyhow::bail!("non-positive sched trace duration");
209 }
210 if cpu_count == 0 {
211 anyhow::bail!("no CPUs observed in sched trace");
212 }
213
214 let total_duration_ns = trace_end_ns - self.trace_start_ns;
215 let bucket_count = total_duration_ns.div_ceil(self.window_ns) as usize;
216
217 let mut output = Vec::with_capacity(bucket_count);
218 for idx in 0..bucket_count {
219 let bucket_start_ns = self.trace_start_ns + idx as u64 * self.window_ns;
220 let bucket_end_ns = (bucket_start_ns + self.window_ns).min(trace_end_ns);
221 let bucket_width_ns = bucket_end_ns - bucket_start_ns;
222 let capacity_ns = bucket_width_ns * cpu_count as u64;
223 let scale = if capacity_ns == 0 {
224 0.0
225 } else {
226 100.0 / capacity_ns as f64
227 };
228
229 let total_busy_ns = self.total_busy_ns.get(idx).copied().unwrap_or(0);
230 let uncategorized_busy_ns = self.uncategorized_busy_ns.get(idx).copied().unwrap_or(0);
231 let categorized_busy_ns: u64 = self
232 .category_busy_ns
233 .iter()
234 .map(|values| values.get(idx).copied().unwrap_or(0))
235 .sum();
236
237 if categorized_busy_ns + uncategorized_busy_ns != total_busy_ns {
238 let category_totals: Vec<_> = self
239 .categories
240 .iter()
241 .enumerate()
242 .filter_map(|(idx_cat, category)| {
243 let value = self
244 .category_busy_ns
245 .get(idx_cat)
246 .and_then(|values| values.get(idx))
247 .copied()
248 .unwrap_or(0);
249 (value > 0).then_some(format!("{}={}ns", category.name, value))
250 })
251 .collect();
252
253 anyhow::bail!(
254 "sched util categories are not mutually exclusive in bucket starting at {}ms: total={}ns uncategorized={}ns categorized_sum={}ns [{}]. Fix --categories so each busy interval matches at most one category",
255 (bucket_start_ns - self.trace_start_ns) / 1_000_000,
256 total_busy_ns,
257 uncategorized_busy_ns,
258 categorized_busy_ns,
259 category_totals.join(", "),
260 );
261 }
262
263 let categories_json = self
264 .categories
265 .iter()
266 .enumerate()
267 .map(|(idx_cat, category)| {
268 let value = self
269 .category_busy_ns
270 .get(idx_cat)
271 .and_then(|values| values.get(idx))
272 .copied()
273 .unwrap_or(0);
274 (category.name.clone(), value as f64 * scale)
275 })
276 .collect();
277
278 output.push(BusyUtilRecord {
279 time_ms: (bucket_start_ns - self.trace_start_ns) / 1_000_000,
280 window_ms,
281 cpu_count,
282 total: total_busy_ns as f64 * scale,
283 uncategorized: uncategorized_busy_ns as f64 * scale,
284 categories: categories_json,
285 });
286 }
287
288 Ok((output, self.interval_count))
289 }
290}
291
292impl BusyIntervalSink for BucketAggregator {
293 fn on_interval(&mut self, interval: &BusyInterval) {
294 let interval_start = interval.start_ns.max(self.trace_start_ns);
295 let interval_end = interval.end_ns;
296 if interval_end <= interval_start {
297 return;
298 }
299
300 let start_idx = ((interval_start - self.trace_start_ns) / self.window_ns) as usize;
301 let end_idx = ((interval_end - self.trace_start_ns - 1) / self.window_ns) as usize;
302 self.ensure_bucket(end_idx);
303
304 let matched_categories = self.matcher.match_indices(interval);
305 for idx in start_idx..=end_idx {
306 let bucket_start = self.trace_start_ns + idx as u64 * self.window_ns;
307 let bucket_end = bucket_start + self.window_ns;
308 let overlap_ns = interval_end.min(bucket_end) - interval_start.max(bucket_start);
309 if overlap_ns == 0 {
310 continue;
311 }
312
313 self.total_busy_ns[idx] += overlap_ns;
314 if matched_categories.is_empty() {
315 self.uncategorized_busy_ns[idx] += overlap_ns;
316 } else {
317 for cat_idx in &matched_categories {
318 self.category_busy_ns[*cat_idx][idx] += overlap_ns;
319 }
320 }
321 }
322
323 self.interval_count += 1;
324 }
325}
326
327#[derive(Debug, Default)]
328struct TraceStats {
329 trace_start_ns: Option<u64>,
330 trace_end_ns: Option<u64>,
331 observed_cpus: BTreeSet<u32>,
332 saw_switch: bool,
333}
334
335impl TraceStats {
336 fn observe(&mut self, record: &PerfSchedScriptRecord, time_ns: u64) {
337 self.observed_cpus.insert(record.cpu);
338 self.trace_start_ns.get_or_insert(time_ns);
339 self.trace_end_ns = Some(time_ns);
340 if record.event == "sched:sched_switch" {
341 self.saw_switch = true;
342 }
343 }
344
345 fn trace_start_ns(&self) -> u64 {
346 self.trace_start_ns.unwrap_or(0)
347 }
348
349 fn finish(&self) -> Result<(u64, usize)> {
350 if !self.saw_switch {
351 anyhow::bail!("no sched:sched_switch events found in perf.sched.jsonl");
352 }
353 let end_ns = self.trace_end_ns.unwrap_or(0);
354 let start_ns = self.trace_start_ns();
355 if end_ns <= start_ns {
356 anyhow::bail!("non-positive sched trace duration");
357 }
358 if self.observed_cpus.is_empty() {
359 anyhow::bail!("no CPUs observed in sched trace");
360 }
361 Ok((end_ns, self.observed_cpus.len()))
362 }
363}
364
365#[derive(Debug)]
366struct SchedBusyTracker<S> {
367 cpu_states: HashMap<u32, CpuState>,
368 sink: S,
369}
370
371impl<S: BusyIntervalSink> SchedBusyTracker<S> {
372 fn new(sink: S) -> Self {
373 Self {
374 cpu_states: HashMap::new(),
375 sink,
376 }
377 }
378
379 fn process_record(&mut self, record: PerfSchedScriptRecord) {
380 let Some(time_ns) = record
381 .sample_time_ns()
382 .or_else(|| sched_time_to_ns(record.time))
383 else {
384 return;
385 };
386
387 let mut emitted = Vec::new();
388 {
389 let state = self
390 .cpu_states
391 .entry(record.cpu)
392 .or_insert_with(|| CpuState {
393 segment_start_ns: time_ns,
394 last_seen_ns: time_ns,
395 ..CpuState::default()
396 });
397 state.last_seen_ns = time_ns;
398
399 if let Some(task) = state.running_task.as_mut() {
400 if record.tid == task.tid && record.hint != task.hint {
401 task.hint = record.hint;
402 if let Some(interval) = sync_cpu_state(state, time_ns) {
403 emitted.push(interval);
404 }
405 }
406 }
407
408 match classify_event(&record) {
409 ClassifiedEvent::SchedSwitch {
410 next_comm,
411 next_hint,
412 next_tid,
413 next_is_idle,
414 } => {
415 if next_is_idle {
416 state.running_task = None;
417 } else {
418 state.running_task = Some(RunningTask {
419 tid: next_tid,
420 comm: next_comm,
421 hint: next_hint,
422 });
423 }
424 if let Some(interval) = sync_cpu_state(state, time_ns) {
425 emitted.push(interval);
426 }
427 }
428 ClassifiedEvent::HardIrqEntry => {
429 state.system_stack.push(ActiveSystem::HardIrq);
430 if let Some(interval) = sync_cpu_state(state, time_ns) {
431 emitted.push(interval);
432 }
433 }
434 ClassifiedEvent::HardIrqExit => {
435 pop_system_override(state, ActiveSystem::HardIrq);
436 if let Some(interval) = sync_cpu_state(state, time_ns) {
437 emitted.push(interval);
438 }
439 }
440 ClassifiedEvent::SoftIrqEntry(system) => {
441 state.system_stack.push(system);
442 if let Some(interval) = sync_cpu_state(state, time_ns) {
443 emitted.push(interval);
444 }
445 }
446 ClassifiedEvent::SoftIrqExit => {
447 pop_softirq_override(state);
448 if let Some(interval) = sync_cpu_state(state, time_ns) {
449 emitted.push(interval);
450 }
451 }
452 ClassifiedEvent::Nmi { delta_ns } => {
453 let nmi_start_ns = time_ns.saturating_sub(delta_ns);
454 if let Some(interval) = flush_current_interval(state, nmi_start_ns) {
455 emitted.push(interval);
456 }
457 emitted.push(BusyInterval {
458 start_ns: nmi_start_ns,
459 end_ns: time_ns,
460 comm: NMI_CATEGORY.to_string(),
461 hint: 0,
462 });
463 state.segment_start_ns = time_ns;
464 }
465 ClassifiedEvent::Other => {}
466 }
467 }
468
469 for interval in emitted {
470 self.sink.on_interval(&interval);
471 }
472 }
473
474 fn finish(mut self, _trace_end_ns: u64) -> S {
475 for state in self.cpu_states.values_mut() {
476 let cpu_end_ns = state.last_seen_ns.max(state.segment_start_ns);
477 if let Some(interval) = flush_current_interval(state, cpu_end_ns) {
478 self.sink.on_interval(&interval);
479 }
480 }
481 self.sink
482 }
483}
484
485#[cfg(test)]
486#[derive(Debug, Default)]
487struct VecIntervalSink {
488 intervals: Vec<BusyInterval>,
489}
490
491#[cfg(test)]
492impl BusyIntervalSink for VecIntervalSink {
493 fn on_interval(&mut self, interval: &BusyInterval) {
494 self.intervals.push(interval.clone());
495 }
496}
497
498fn current_active_execution(state: &CpuState) -> Option<ActiveExecution> {
499 if let Some(system) = state.system_stack.last().copied() {
500 return Some(ActiveExecution::System(system));
501 }
502
503 state.running_task.as_ref().map(|task| {
504 if is_ksoftirqd_comm(&task.comm) {
505 ActiveExecution::System(ActiveSystem::SoftIrqOther)
506 } else {
507 ActiveExecution::Task {
508 comm: task.comm.clone(),
509 hint: task.hint,
510 }
511 }
512 })
513}
514
515fn is_ksoftirqd_comm(comm: &str) -> bool {
516 comm.starts_with("ksoftirqd/")
517}
518
519fn flush_current_interval(state: &mut CpuState, end_ns: u64) -> Option<BusyInterval> {
520 let interval = state
521 .active_execution
522 .as_ref()
523 .and_then(|active| active.to_interval(state.segment_start_ns, end_ns));
524 state.segment_start_ns = end_ns;
525 interval
526}
527
528fn sync_cpu_state(state: &mut CpuState, time_ns: u64) -> Option<BusyInterval> {
529 let next_active = current_active_execution(state);
530 if state.active_execution == next_active {
531 return None;
532 }
533
534 let interval = flush_current_interval(state, time_ns);
535 state.active_execution = next_active;
536 interval
537}
538
539fn pop_system_override(state: &mut CpuState, system: ActiveSystem) {
540 if let Some(pos) = state
541 .system_stack
542 .iter()
543 .rposition(|value| *value == system)
544 {
545 state.system_stack.remove(pos);
546 }
547}
548
549fn pop_softirq_override(state: &mut CpuState) {
550 if let Some(pos) = state
551 .system_stack
552 .iter()
553 .rposition(|value| value.is_softirq())
554 {
555 state.system_stack.remove(pos);
556 }
557}
558
559fn parse_sched_categories(spec: &str) -> Result<Vec<CategorySpec>> {
560 let mut categories = Vec::new();
561 let mut seen_names = HashSet::new();
562 let mut hinted_comm_parts = BTreeSet::new();
563 let mut exact_zero_hint_comms = HashSet::new();
564 let mut glob_zero_hint_comms = HashSet::new();
565
566 for raw_item in spec.split(',') {
567 let item = raw_item.trim();
568 if item.is_empty() {
569 continue;
570 }
571
572 let (comm_part, hint) = match item.split_once("@hint=") {
573 Some((comm, hint_str)) => (
574 comm.trim(),
575 Some(
576 hint_str
577 .trim()
578 .parse::<u64>()
579 .with_context(|| format!("invalid hint selector in category '{item}'"))?,
580 ),
581 ),
582 None => (item, None),
583 };
584
585 let comm_matcher = if comm_part.contains('*') || comm_part.contains('?') {
586 CategoryCommMatcher::Glob(comm_part.to_string())
587 } else {
588 CategoryCommMatcher::Exact(comm_part.to_string())
589 };
590
591 if hint.is_some() {
592 hinted_comm_parts.insert(comm_part.to_string());
593 }
594 if hint == Some(0) {
595 match &comm_matcher {
596 CategoryCommMatcher::Exact(comm) => {
597 exact_zero_hint_comms.insert(comm.clone());
598 }
599 CategoryCommMatcher::Glob(pattern) => {
600 glob_zero_hint_comms.insert(pattern.clone());
601 }
602 }
603 }
604
605 if !seen_names.insert(item.to_string()) {
606 continue;
607 }
608
609 categories.push(CategorySpec {
610 name: item.to_string(),
611 comm_matcher,
612 hint,
613 });
614 }
615
616 for comm_part in hinted_comm_parts {
617 let needs_zero_hint = if comm_part.contains('*') || comm_part.contains('?') {
618 !glob_zero_hint_comms.contains(&comm_part)
619 } else {
620 !exact_zero_hint_comms.contains(&comm_part)
621 };
622
623 if !needs_zero_hint {
624 continue;
625 }
626
627 let zero_hint_name = format!("{comm_part}@hint=0");
628 if !seen_names.insert(zero_hint_name.clone()) {
629 continue;
630 }
631
632 let comm_matcher = if comm_part.contains('*') || comm_part.contains('?') {
633 CategoryCommMatcher::Glob(comm_part.clone())
634 } else {
635 CategoryCommMatcher::Exact(comm_part.clone())
636 };
637
638 categories.push(CategorySpec {
639 name: zero_hint_name,
640 comm_matcher,
641 hint: Some(0),
642 });
643 }
644
645 for name in DEFAULT_SYSTEM_CATEGORY_NAMES {
646 if seen_names.insert(name.to_string()) {
647 categories.push(CategorySpec {
648 name: name.to_string(),
649 comm_matcher: CategoryCommMatcher::Exact(name.to_string()),
650 hint: None,
651 });
652 }
653 }
654
655 Ok(categories)
656}
657
658fn compile_category_matcher(categories: &[CategorySpec]) -> CompiledCategoryMatcher {
659 let mut matcher = CompiledCategoryMatcher::default();
660
661 for (index, category) in categories.iter().enumerate() {
662 match &category.comm_matcher {
663 CategoryCommMatcher::Exact(comm) => {
664 if let Some(hint) = category.hint {
665 matcher
666 .exact_by_hint
667 .entry(comm.clone())
668 .or_default()
669 .entry(hint)
670 .or_default()
671 .push(index);
672 } else {
673 matcher
674 .exact_any_hint
675 .entry(comm.clone())
676 .or_default()
677 .push(index);
678 }
679 }
680 CategoryCommMatcher::Glob(pattern) => matcher.glob_specs.push(GlobCategory {
681 index,
682 pattern: pattern.clone(),
683 hint: category.hint,
684 }),
685 }
686 }
687
688 matcher
689}
690
691impl CompiledCategoryMatcher {
692 fn match_indices(&self, interval: &BusyInterval) -> Vec<usize> {
693 let mut matched = Vec::new();
694
695 if let Some(indices) = self.exact_any_hint.get(&interval.comm) {
696 matched.extend(indices.iter().copied());
697 }
698 if let Some(by_hint) = self.exact_by_hint.get(&interval.comm) {
699 if let Some(indices) = by_hint.get(&interval.hint) {
700 matched.extend(indices.iter().copied());
701 }
702 }
703
704 for glob in &self.glob_specs {
705 if glob.hint.map(|hint| hint == interval.hint).unwrap_or(true)
706 && glob_match(&glob.pattern, &interval.comm)
707 {
708 matched.push(glob.index);
709 }
710 }
711
712 matched.sort_unstable();
713 matched.dedup();
714 matched
715 }
716}
717
718fn classify_event(record: &PerfSchedScriptRecord) -> ClassifiedEvent {
719 match record.event.as_str() {
720 "sched:sched_switch" => {
721 let Some(fields) = record.fields.as_ref() else {
722 return ClassifiedEvent::Other;
723 };
724
725 let next_comm = fields
726 .get("next_comm")
727 .and_then(Value::as_str)
728 .unwrap_or_default()
729 .to_string();
730 let next_hint = fields.get("next_hint").and_then(Value::as_u64).unwrap_or(0);
731 let next_tid = fields
732 .get("next_pid")
733 .and_then(Value::as_i64)
734 .and_then(|tid| i32::try_from(tid).ok())
735 .unwrap_or(0);
736 ClassifiedEvent::SchedSwitch {
737 next_is_idle: sched_switch_next_is_idle(fields),
738 next_comm,
739 next_hint,
740 next_tid,
741 }
742 }
743 "irq:irq_handler_entry" => ClassifiedEvent::HardIrqEntry,
744 "irq:irq_handler_exit" => ClassifiedEvent::HardIrqExit,
745 "irq:softirq_entry" => ClassifiedEvent::SoftIrqEntry(classify_softirq_system(record)),
746 "irq:softirq_exit" => ClassifiedEvent::SoftIrqExit,
747 _ => parse_nmi_duration_ns(record)
748 .map(|delta_ns| ClassifiedEvent::Nmi { delta_ns })
749 .unwrap_or(ClassifiedEvent::Other),
750 }
751}
752
753fn classify_softirq_system(record: &PerfSchedScriptRecord) -> ActiveSystem {
754 let action = record
755 .fields
756 .as_ref()
757 .and_then(|fields| fields.get("action"))
758 .and_then(Value::as_str)
759 .unwrap_or_default();
760
761 match action {
762 "NET_RX" => ActiveSystem::SoftIrqRx,
763 "NET_TX" => ActiveSystem::SoftIrqTx,
764 _ => ActiveSystem::SoftIrqOther,
765 }
766}
767
768fn sched_switch_next_is_idle(fields: &serde_json::Map<String, Value>) -> bool {
769 let next_pid = fields.get("next_pid").and_then(Value::as_i64).unwrap_or(-1);
770 let next_comm = fields
771 .get("next_comm")
772 .and_then(Value::as_str)
773 .unwrap_or_default();
774
775 next_pid == 0 || next_comm.starts_with("swapper")
776}
777
778fn parse_nmi_duration_ns(record: &PerfSchedScriptRecord) -> Option<u64> {
779 if !record.event.starts_with("nmi:") {
780 return None;
781 }
782
783 let needle = "delta_ns:";
784 let start = record.trace.find(needle)? + needle.len();
785 record.trace[start..]
786 .split_whitespace()
787 .next()?
788 .trim_end_matches(',')
789 .parse::<u64>()
790 .ok()
791}
792
793fn sched_time_to_ns(time: f64) -> Option<u64> {
794 if !time.is_finite() || time < 0.0 {
795 return None;
796 }
797 let ns = time * 1_000_000_000.0;
798 if ns < 0.0 || ns > u64::MAX as f64 {
799 return None;
800 }
801 Some(ns as u64)
802}
803
804fn glob_match(pattern: &str, value: &str) -> bool {
805 let p = pattern.as_bytes();
806 let v = value.as_bytes();
807 let (mut pi, mut vi) = (0usize, 0usize);
808 let mut star_pi = None;
809 let mut star_vi = 0usize;
810
811 while vi < v.len() {
812 if pi < p.len() && (p[pi] == b'?' || p[pi] == v[vi]) {
813 pi += 1;
814 vi += 1;
815 } else if pi < p.len() && p[pi] == b'*' {
816 star_pi = Some(pi);
817 pi += 1;
818 star_vi = vi;
819 } else if let Some(saved_pi) = star_pi {
820 pi = saved_pi + 1;
821 star_vi += 1;
822 vi = star_vi;
823 } else {
824 return false;
825 }
826 }
827
828 while pi < p.len() && p[pi] == b'*' {
829 pi += 1;
830 }
831
832 pi == p.len()
833}
834
835pub fn cmd_extract_sched_util(opts: ExtractSchedUtilOpts) -> Result<()> {
836 if opts.window_ms == 0 {
837 anyhow::bail!("--window-ms must be greater than zero");
838 }
839
840 let categories = parse_sched_categories(&opts.categories)?;
841 let file = File::open(&opts.file).context("failed to open perf.sched.jsonl")?;
842 let reader = BufReader::new(file);
843
844 let mut stats = TraceStats::default();
845 let window_ns = opts.window_ms * 1_000_000;
846 let mut aggregator = BucketAggregator::new(0, window_ns, categories);
847 let mut tracker = SchedBusyTracker::new(aggregator);
848
849 for line in reader.lines() {
850 let line = line.context("failed to read line")?;
851 if line.trim().is_empty() {
852 continue;
853 }
854
855 let record: PerfSchedScriptRecord =
856 serde_json::from_str(&line).context("failed to parse perf.sched.jsonl record")?;
857 let Some(time_ns) = record
858 .sample_time_ns()
859 .or_else(|| sched_time_to_ns(record.time))
860 else {
861 continue;
862 };
863
864 if stats.trace_start_ns.is_none() {
865 tracker.sink.trace_start_ns = time_ns;
866 }
867 stats.observe(&record, time_ns);
868 tracker.process_record(record);
869 }
870
871 let (trace_end_ns, cpu_count) = stats.finish()?;
872 aggregator = tracker.finish(trace_end_ns);
873 let category_count = aggregator.categories.len();
874 let (output, interval_count) = aggregator.finalize(trace_end_ns, cpu_count, opts.window_ms)?;
875
876 if opts.verbose {
877 eprintln!(
878 "sched util: {} intervals, {} buckets, {} cpus, {} categories",
879 interval_count,
880 output.len(),
881 cpu_count,
882 category_count
883 );
884 }
885
886 for record in output {
887 println!(
888 "{}",
889 serde_json::to_string(&json!({
890 "time_ms": record.time_ms,
891 "window_ms": record.window_ms,
892 "cpu_count": record.cpu_count,
893 "total": record.total,
894 "uncategorized": record.uncategorized,
895 "categories": record.categories,
896 }))?
897 );
898 }
899
900 Ok(())
901}
902
903#[cfg(test)]
904fn build_busy_intervals(
905 records: impl IntoIterator<Item = PerfSchedScriptRecord>,
906) -> Result<(Vec<BusyInterval>, u64, u64, usize)> {
907 let mut stats = TraceStats::default();
908 let mut tracker = SchedBusyTracker::new(VecIntervalSink::default());
909
910 for record in records {
911 let Some(time_ns) = record
912 .sample_time_ns()
913 .or_else(|| sched_time_to_ns(record.time))
914 else {
915 continue;
916 };
917 stats.observe(&record, time_ns);
918 tracker.process_record(record);
919 }
920
921 let (trace_end_ns, cpu_count) = stats.finish()?;
922 let sink = tracker.finish(trace_end_ns);
923 Ok((
924 sink.intervals,
925 stats.trace_start_ns(),
926 trace_end_ns,
927 cpu_count,
928 ))
929}
930
931#[cfg(test)]
932mod tests {
933 use super::*;
934 use std::fs;
935 use std::path::PathBuf;
936 use tempfile::TempDir;
937
938 fn write_sched_jsonl(tempdir: &TempDir, lines: &[&str]) -> PathBuf {
939 let path = tempdir.path().join("perf.sched.jsonl");
940 fs::write(&path, lines.join("\n") + "\n").expect("failed to write perf.sched.jsonl");
941 path
942 }
943
944 #[test]
945 fn sched_util_includes_total_uncategorized_and_hint_categories() -> Result<()> {
946 let tempdir = TempDir::new()?;
947 let path = write_sched_jsonl(
948 &tempdir,
949 &[
950 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"worker-a","next_pid":10,"next_prio":120},"hint":0}"#,
951 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0005,"event":"sched:sched_stat_runtime","trace":"comm=worker-a runtime=500000 [ns] vruntime=0 [ns]","fields":{"comm":"worker-a","runtime":500000,"vruntime":0},"hint":640}"#,
952 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120},"hint":640}"#,
953 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> svc-4:20 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"svc-4","next_pid":20,"next_prio":120},"hint":0}"#,
954 r#"{"comm":"svc-4","pid":20,"tid":20,"cpu":0,"time":0.002,"event":"sched:sched_switch","trace":"svc-4:20 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"svc-4","prev_pid":20,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120},"hint":0}"#,
955 ],
956 );
957
958 let categories = parse_sched_categories("worker-a@hint=640,svc-*")?;
959 let file = File::open(&path)?;
960 let reader = BufReader::new(file);
961 let records: Vec<PerfSchedScriptRecord> = reader
962 .lines()
963 .map(|line| -> Result<_> {
964 let line = line?;
965 Ok(serde_json::from_str(&line)?)
966 })
967 .collect::<Result<_>>()?;
968
969 let (intervals, start_ns, end_ns, cpu_count) = build_busy_intervals(records)?;
970 assert_eq!(cpu_count, 1);
971 assert_eq!(start_ns, 0);
972 assert_eq!(end_ns, 2_000_000);
973 assert_eq!(intervals.len(), 3);
974 assert_eq!(intervals[0].comm, "worker-a");
975 assert_eq!(intervals[0].hint, 0);
976 assert_eq!(intervals[1].comm, "worker-a");
977 assert_eq!(intervals[1].hint, 640);
978 assert_eq!(intervals[2].comm, "svc-4");
979 assert_eq!(intervals[2].hint, 0);
980
981 let matcher = compile_category_matcher(&categories);
982 let category_names: Vec<_> = categories.iter().map(|cat| cat.name.as_str()).collect();
983 assert_eq!(
984 category_names[matcher.match_indices(&intervals[0])[0]],
985 "worker-a@hint=0"
986 );
987 assert_eq!(
988 category_names[matcher.match_indices(&intervals[1])[0]],
989 "worker-a@hint=640"
990 );
991 assert_eq!(
992 category_names[matcher.match_indices(&intervals[2])[0]],
993 "svc-*"
994 );
995
996 Ok(())
997 }
998
999 #[test]
1000 fn sched_util_auto_adds_zero_hint_category_for_hinted_specs() -> Result<()> {
1001 let categories =
1002 parse_sched_categories("hhvmworker@hint=256,hhvmworker@hint=640,mcrpxy-*")?;
1003 let names: BTreeSet<_> = categories.iter().map(|cat| cat.name.as_str()).collect();
1004
1005 assert!(names.contains("hhvmworker@hint=0"));
1006 assert!(names.contains("hhvmworker@hint=256"));
1007 assert!(names.contains("hhvmworker@hint=640"));
1008 assert!(!names.contains("mcrpxy-*@hint=0"));
1009
1010 Ok(())
1011 }
1012
1013 #[test]
1014 fn sched_util_does_not_duplicate_explicit_zero_hint_category() -> Result<()> {
1015 let categories = parse_sched_categories("hhvmworker@hint=0,hhvmworker@hint=384")?;
1016 let names: Vec<_> = categories
1017 .iter()
1018 .filter(|cat| cat.name.starts_with("hhvmworker@hint="))
1019 .map(|cat| cat.name.as_str())
1020 .collect();
1021
1022 assert_eq!(names, vec!["hhvmworker@hint=0", "hhvmworker@hint=384"]);
1023
1024 Ok(())
1025 }
1026
1027 #[test]
1028 fn sched_util_outputs_total_and_uncategorized_records() -> Result<()> {
1029 let tempdir = TempDir::new()?;
1030 let path = write_sched_jsonl(
1031 &tempdir,
1032 &[
1033 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"worker-a","next_pid":10,"next_prio":120},"hint":0}"#,
1034 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0005,"event":"sched:sched_stat_runtime","trace":"comm=worker-a runtime=500000 [ns] vruntime=0 [ns]","fields":{"comm":"worker-a","runtime":500000,"vruntime":0},"hint":640}"#,
1035 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120},"hint":640}"#,
1036 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> uncat:20 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"uncat","next_pid":20,"next_prio":120},"hint":0}"#,
1037 r#"{"comm":"uncat","pid":20,"tid":20,"cpu":0,"time":0.002,"event":"sched:sched_switch","trace":"uncat:20 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"uncat","prev_pid":20,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120},"hint":0}"#,
1038 ],
1039 );
1040
1041 let opts = ExtractSchedUtilOpts {
1042 file: path,
1043 window_ms: 1,
1044 categories: "worker-a@hint=640".to_string(),
1045 verbose: false,
1046 };
1047
1048 let categories = parse_sched_categories(&opts.categories)?;
1049 let window_ns = opts.window_ms * 1_000_000;
1050 let mut agg = BucketAggregator::new(0, window_ns, categories);
1051 let file = File::open(&opts.file)?;
1052 let reader = BufReader::new(file);
1053 let mut stats = TraceStats::default();
1054 let mut tracker = SchedBusyTracker::new(agg);
1055
1056 for line in reader.lines() {
1057 let line = line?;
1058 let record: PerfSchedScriptRecord = serde_json::from_str(&line)?;
1059 let Some(time_ns) = record
1060 .sample_time_ns()
1061 .or_else(|| sched_time_to_ns(record.time))
1062 else {
1063 continue;
1064 };
1065 if stats.trace_start_ns.is_none() {
1066 tracker.sink.trace_start_ns = time_ns;
1067 }
1068 stats.observe(&record, time_ns);
1069 tracker.process_record(record);
1070 }
1071 let (trace_end_ns, cpu_count) = stats.finish()?;
1072 agg = tracker.finish(trace_end_ns);
1073
1074 assert_eq!(cpu_count, 1);
1075 let (output, _interval_count) = agg.finalize(trace_end_ns, cpu_count, opts.window_ms)?;
1076 assert_eq!(output.len(), 2);
1077 assert_eq!(output[0].total, 100.0);
1078 assert_eq!(output[0].uncategorized, 0.0);
1079 assert_eq!(output[0].categories["worker-a@hint=0"], 50.0);
1080 assert_eq!(output[0].categories["worker-a@hint=640"], 50.0);
1081 assert_eq!(output[1].uncategorized, 100.0);
1082
1083 Ok(())
1084 }
1085
1086 #[test]
1087 fn sched_util_rejects_overlapping_categories() -> Result<()> {
1088 let tempdir = TempDir::new()?;
1089 let path = write_sched_jsonl(
1090 &tempdir,
1091 &[
1092 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"worker-a","next_pid":10,"next_prio":120,"next_hint":640},"hint":0}"#,
1093 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120,"next_hint":0},"hint":640}"#,
1094 ],
1095 );
1096
1097 let opts = ExtractSchedUtilOpts {
1098 file: path,
1099 window_ms: 1,
1100 categories: "worker-a,worker-a@hint=640".to_string(),
1101 verbose: false,
1102 };
1103
1104 let categories = parse_sched_categories(&opts.categories)?;
1105 let window_ns = opts.window_ms * 1_000_000;
1106 let mut agg = BucketAggregator::new(0, window_ns, categories);
1107 let file = File::open(&opts.file)?;
1108 let reader = BufReader::new(file);
1109 let mut stats = TraceStats::default();
1110 let mut tracker = SchedBusyTracker::new(agg);
1111
1112 for line in reader.lines() {
1113 let line = line?;
1114 let record: PerfSchedScriptRecord = serde_json::from_str(&line)?;
1115 let Some(time_ns) = record
1116 .sample_time_ns()
1117 .or_else(|| sched_time_to_ns(record.time))
1118 else {
1119 continue;
1120 };
1121 if stats.trace_start_ns.is_none() {
1122 tracker.sink.trace_start_ns = time_ns;
1123 }
1124 stats.observe(&record, time_ns);
1125 tracker.process_record(record);
1126 }
1127 let (trace_end_ns, cpu_count) = stats.finish()?;
1128 agg = tracker.finish(trace_end_ns);
1129
1130 let err = agg
1131 .finalize(trace_end_ns, cpu_count, opts.window_ms)
1132 .expect_err("expected overlapping categories to be rejected");
1133 assert!(err
1134 .to_string()
1135 .contains("sched util categories are not mutually exclusive"));
1136
1137 Ok(())
1138 }
1139
1140 #[test]
1141 fn sched_util_splits_busy_time_on_on_cpu_hint_changes_and_tracks_irqs() -> Result<()> {
1142 let tempdir = TempDir::new()?;
1143 let path = write_sched_jsonl(
1144 &tempdir,
1145 &[
1146 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"worker-a","next_pid":10,"next_prio":120},"hint":0}"#,
1147 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0002,"event":"sched:sched_stat_runtime","trace":"comm=worker-a runtime=200000 [ns] vruntime=0 [ns]","fields":{"comm":"worker-a","runtime":200000,"vruntime":0},"hint":384}"#,
1148 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0006,"event":"irq:softirq_entry","trace":"vec=3 [action=NET_RX]","fields":{"action":"NET_RX","vec":3},"hint":384}"#,
1149 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0007,"event":"irq:softirq_exit","trace":"vec=3 [action=NET_RX]","fields":{"action":"NET_RX","vec":3},"hint":384}"#,
1150 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0008,"event":"sched:sched_stat_runtime","trace":"comm=worker-a runtime=800000 [ns] vruntime=0 [ns]","fields":{"comm":"worker-a","runtime":800000,"vruntime":0},"hint":640}"#,
1151 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.0010,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","next_comm":"swapper/0","next_pid":0,"next_prio":120},"hint":640}"#,
1152 ],
1153 );
1154
1155 let file = File::open(&path)?;
1156 let reader = BufReader::new(file);
1157 let records: Vec<PerfSchedScriptRecord> = reader
1158 .lines()
1159 .map(|line| -> Result<_> {
1160 let line = line?;
1161 Ok(serde_json::from_str(&line)?)
1162 })
1163 .collect::<Result<_>>()?;
1164
1165 let (intervals, _start_ns, _end_ns, cpu_count) = build_busy_intervals(records)?;
1166 assert_eq!(cpu_count, 1);
1167 assert_eq!(intervals.len(), 5);
1168
1169 assert_eq!(intervals[0].comm, "worker-a");
1170 assert_eq!(intervals[0].hint, 0);
1171 assert_eq!((intervals[0].start_ns, intervals[0].end_ns), (0, 200_000));
1172
1173 assert_eq!(intervals[1].comm, "worker-a");
1174 assert_eq!(intervals[1].hint, 384);
1175 assert_eq!(
1176 (intervals[1].start_ns, intervals[1].end_ns),
1177 (200_000, 600_000)
1178 );
1179
1180 assert_eq!(intervals[2].comm, "softirq-rx");
1181 assert_eq!(intervals[2].hint, 0);
1182 assert_eq!(
1183 (intervals[2].start_ns, intervals[2].end_ns),
1184 (600_000, 700_000)
1185 );
1186
1187 assert_eq!(intervals[3].comm, "worker-a");
1188 assert_eq!(intervals[3].hint, 384);
1189 assert_eq!(
1190 (intervals[3].start_ns, intervals[3].end_ns),
1191 (700_000, 800_000)
1192 );
1193
1194 assert_eq!(intervals[4].comm, "worker-a");
1195 assert_eq!(intervals[4].hint, 640);
1196 assert_eq!(
1197 (intervals[4].start_ns, intervals[4].end_ns),
1198 (800_000, 1_000_000)
1199 );
1200
1201 Ok(())
1202 }
1203
1204 #[test]
1205 fn sched_util_uses_next_hint_from_switch_for_immediate_on_cpu_time() -> Result<()> {
1206 let tempdir = TempDir::new()?;
1207 let path = write_sched_jsonl(
1208 &tempdir,
1209 &[
1210 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","prev_hint":0,"next_comm":"worker-a","next_pid":10,"next_prio":120,"next_hint":640},"hint":0}"#,
1211 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":0,"time":0.001,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/0:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","prev_hint":640,"next_comm":"swapper/0","next_pid":0,"next_prio":120,"next_hint":0},"hint":640}"#,
1212 ],
1213 );
1214
1215 let file = File::open(&path)?;
1216 let reader = BufReader::new(file);
1217 let records: Vec<PerfSchedScriptRecord> = reader
1218 .lines()
1219 .map(|line| -> Result<_> {
1220 let line = line?;
1221 Ok(serde_json::from_str(&line)?)
1222 })
1223 .collect::<Result<_>>()?;
1224
1225 let (intervals, _start_ns, _end_ns, cpu_count) = build_busy_intervals(records)?;
1226 assert_eq!(cpu_count, 1);
1227 assert_eq!(intervals.len(), 1);
1228 assert_eq!(intervals[0].comm, "worker-a");
1229 assert_eq!(intervals[0].hint, 640);
1230 assert_eq!((intervals[0].start_ns, intervals[0].end_ns), (0, 1_000_000));
1231
1232 Ok(())
1233 }
1234
1235 #[test]
1236 fn sched_util_does_not_extrapolate_last_running_task_to_global_trace_end() -> Result<()> {
1237 let tempdir = TempDir::new()?;
1238 let path = write_sched_jsonl(
1239 &tempdir,
1240 &[
1241 r#"{"comm":"idle","pid":0,"tid":0,"cpu":0,"time":0.0,"event":"sched:sched_switch","trace":"swapper/0:0 [120] R ==> perf:900 [120]","fields":{"prev_comm":"swapper/0","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"perf","next_pid":900,"next_prio":120},"hint":0}"#,
1242 r#"{"comm":"perf","pid":900,"tid":900,"cpu":0,"time":0.001,"event":"sched:sched_stat_runtime","trace":"comm=perf runtime=1000000 [ns] vruntime=0 [ns]","fields":{"comm":"perf","runtime":1000000,"vruntime":0},"hint":0}"#,
1243 r#"{"comm":"idle","pid":0,"tid":0,"cpu":1,"time":0.0,"event":"sched:sched_switch","trace":"swapper/1:0 [120] R ==> worker-a:10 [120]","fields":{"prev_comm":"swapper/1","prev_pid":0,"prev_prio":120,"prev_state":"R","next_comm":"worker-a","next_pid":10,"next_prio":120},"hint":0}"#,
1244 r#"{"comm":"worker-a","pid":10,"tid":10,"cpu":1,"time":0.002,"event":"sched:sched_switch","trace":"worker-a:10 [120] R ==> swapper/1:0 [120]","fields":{"prev_comm":"worker-a","prev_pid":10,"prev_prio":120,"prev_state":"R","next_comm":"swapper/1","next_pid":0,"next_prio":120},"hint":0}"#,
1245 ],
1246 );
1247
1248 let file = File::open(&path)?;
1249 let reader = BufReader::new(file);
1250 let records: Vec<PerfSchedScriptRecord> = reader
1251 .lines()
1252 .map(|line| -> Result<_> {
1253 let line = line?;
1254 Ok(serde_json::from_str(&line)?)
1255 })
1256 .collect::<Result<_>>()?;
1257
1258 let (intervals, _start_ns, end_ns, cpu_count) = build_busy_intervals(records)?;
1259 assert_eq!(cpu_count, 2);
1260 assert_eq!(end_ns, 2_000_000);
1261
1262 let perf_interval = intervals
1263 .iter()
1264 .find(|interval| interval.comm == "perf")
1265 .expect("expected perf interval");
1266 assert_eq!(
1267 (perf_interval.start_ns, perf_interval.end_ns),
1268 (0, 1_000_000)
1269 );
1270
1271 let worker_interval = intervals
1272 .iter()
1273 .find(|interval| interval.comm == "worker-a")
1274 .expect("expected worker interval");
1275 assert_eq!(
1276 (worker_interval.start_ns, worker_interval.end_ns),
1277 (0, 2_000_000)
1278 );
1279
1280 Ok(())
1281 }
1282}