1use std::collections::HashMap;
23use std::path::Path;
24use std::thread::available_parallelism;
25
26use serde::{Deserialize, Serialize};
27
28use crate::helpers::result_archive_path;
29
30pub const PHASES: [&str; 17] = [
35 "bootstrap",
36 "digest",
37 "build",
38 "rewrite",
39 "math_parse",
40 "post_xml_parse",
41 "post_scan",
42 "bibliography",
43 "crossref",
44 "graphics",
45 "math_images",
46 "mathml_pres",
47 "mathml_cont",
48 "split",
49 "xslt",
50 "html5_fixups",
51 "serialize",
52];
53
54#[derive(Debug, Default, Deserialize)]
58pub struct TelemetryRecord {
59 #[serde(default)]
61 pub paper_id: String,
62 #[serde(default)]
64 pub git_sha: String,
65 #[serde(default)]
67 pub host: String,
68 #[serde(default)]
70 pub category: String,
71 #[serde(default)]
73 pub exit_code: i64,
74 #[serde(default)]
76 pub wall_us: u64,
77 #[serde(default)]
79 pub max_rss_kb: u64,
80 #[serde(default)]
83 pub phase_us: Vec<u64>,
84 #[serde(default)]
86 pub warnings: u64,
87 #[serde(default)]
89 pub errors: u64,
90 #[serde(default)]
92 pub fatal_errors: u64,
93 #[serde(default)]
95 pub formulae: u64,
96 #[serde(default)]
98 pub math_parse_attempts: u64,
99 #[serde(default)]
102 pub math_parse_count: u64,
103 #[serde(default)]
105 pub graphics_assets: u64,
106 #[serde(default)]
108 pub output_bytes: u64,
109}
110
111pub fn read_telemetry_json(result: &Path) -> Result<TelemetryRecord, String> {
118 let file = std::fs::File::open(result).map_err(|e| format!("cannot open result archive: {e}"))?;
119 let mut archive = zip::ZipArchive::new(file).map_err(|e| format!("not a readable zip: {e}"))?;
120 let mut entry = archive
121 .by_name("telemetry.json")
122 .map_err(|e| format!("no telemetry.json entry: {e}"))?;
123 let mut raw = Vec::new();
124 {
125 use std::io::Read;
126 entry
127 .read_to_end(&mut raw)
128 .map_err(|e| format!("reading telemetry.json failed: {e}"))?;
129 }
130 serde_json::from_slice(&raw).map_err(|e| format!("malformed telemetry.json: {e}"))
131}
132
133#[derive(Debug, Default, Clone, Serialize)]
135pub struct Percentiles {
136 pub p50: u64,
138 pub p90: u64,
140 pub p99: u64,
142 pub max: u64,
144}
145
146pub fn percentiles(values: &[u64]) -> Percentiles {
150 let n = values.len();
151 if n == 0 {
152 return Percentiles::default();
153 }
154 let mut sorted = values.to_vec();
155 sorted.sort_unstable();
156 let at = |p: usize| -> u64 {
158 let rank = (p * n).div_ceil(100).clamp(1, n);
159 sorted[rank - 1]
160 };
161 Percentiles {
162 p50: at(50),
163 p90: at(90),
164 p99: at(99),
165 max: sorted[n - 1],
166 }
167}
168
169#[derive(Debug, Default, Clone, Serialize)]
173pub struct TailStats {
174 pub top1pct_wall_share: f64,
176 pub top5pct_wall_share: f64,
178 pub over_30s: u64,
180 pub over_60s: u64,
182 pub over_120s: u64,
184 pub over_180s: u64,
186}
187
188#[derive(Debug, Default, Clone, Serialize)]
191pub struct RssBuckets {
192 pub over_2gib: u64,
194 pub over_3gib: u64,
196 pub over_4gib: u64,
198}
199
200#[derive(Debug, Default, Clone, Serialize)]
203pub struct MathStats {
204 pub formulae: u64,
206 pub parse_invocations: u64,
208 pub parse_count: u64,
210 pub parses_per_formula: f64,
212}
213
214#[derive(Debug, Default, Clone, Serialize)]
217pub struct OutcomeWall {
218 pub n: usize,
220 pub median_ms: u64,
222 pub mean_ms: u64,
224 pub p99_ms: u64,
226}
227
228#[derive(Debug, Clone, Serialize)]
231pub struct TelemetrySummary {
232 pub corpus: String,
234 pub service: String,
236 pub sample_count: usize,
238 pub skipped: usize,
240 pub outcome_counts: Vec<(String, u64)>,
242 pub wall_ms: Percentiles,
244 pub rss_mib: Percentiles,
246 pub phase_p99_ms: Vec<(String, u64)>,
248 pub phase_wall_pct: Vec<(String, f64)>,
251 pub tail: TailStats,
253 pub rss_buckets: RssBuckets,
255 pub math: MathStats,
257 pub slow_tail_dominant: Vec<(String, u64)>,
260 pub fatal_profile: OutcomeWall,
262 pub no_problem_profile: OutcomeWall,
264 pub slowest: Option<(String, u64)>,
266 pub highest_rss: Option<(String, u64)>,
268 pub total_formulae: u64,
270 pub total_graphics_assets: u64,
272 pub total_output_bytes: u64,
274 pub generated_unix: i64,
276}
277
278pub fn summarize(corpus: &str, service: &str, records: &[TelemetryRecord]) -> TelemetrySummary {
283 let bucket = |r: &TelemetryRecord| -> &'static str {
286 if r.fatal_errors > 0 || r.category.contains("fatal") || r.exit_code >= 3 {
287 "fatal"
288 } else if r.errors > 0 || r.category == "conversion_error" || r.exit_code == 2 {
289 "error"
290 } else if r.warnings > 0 {
291 "warning"
292 } else {
293 "no_problem"
294 }
295 };
296 let (mut no_problem, mut warning, mut error, mut fatal) = (0u64, 0u64, 0u64, 0u64);
297 for record in records {
298 match bucket(record) {
299 "fatal" => fatal += 1,
300 "error" => error += 1,
301 "warning" => warning += 1,
302 _ => no_problem += 1,
303 }
304 }
305 let outcome_counts = vec![
306 ("no_problem".to_string(), no_problem),
307 ("warning".to_string(), warning),
308 ("error".to_string(), error),
309 ("fatal".to_string(), fatal),
310 ];
311
312 let wall_ms_values: Vec<u64> = records.iter().map(|record| record.wall_us / 1000).collect();
315 let rss_mib_values: Vec<u64> = records
316 .iter()
317 .map(|record| record.max_rss_kb / 1024)
318 .collect();
319 let wall_ms = percentiles(&wall_ms_values);
320 let rss_mib = percentiles(&rss_mib_values);
321
322 let phase_p99_ms = PHASES
325 .iter()
326 .enumerate()
327 .map(|(index, &name)| {
328 let phase_ms: Vec<u64> = records
329 .iter()
330 .filter_map(|record| record.phase_us.get(index).map(|us| us / 1000))
331 .collect();
332 (name.to_string(), percentiles(&phase_ms).p99)
333 })
334 .collect();
335
336 let slowest = records
338 .iter()
339 .max_by_key(|record| record.wall_us)
340 .map(|record| (record.paper_id.clone(), record.wall_us / 1000));
341 let highest_rss = records
342 .iter()
343 .max_by_key(|record| record.max_rss_kb)
344 .map(|record| (record.paper_id.clone(), record.max_rss_kb / 1024));
345
346 let total_formulae = records
348 .iter()
349 .fold(0u64, |acc, record| acc.saturating_add(record.formulae));
350 let total_graphics_assets = records.iter().fold(0u64, |acc, record| {
351 acc.saturating_add(record.graphics_assets)
352 });
353 let total_output_bytes = records
354 .iter()
355 .fold(0u64, |acc, record| acc.saturating_add(record.output_bytes));
356
357 let total_wall: u64 = records
359 .iter()
360 .map(|r| r.wall_us)
361 .fold(0, u64::saturating_add);
362 let mut phase_tot = [0u128; 17];
363 for r in records {
364 for (i, &us) in r.phase_us.iter().take(17).enumerate() {
365 phase_tot[i] += us as u128;
366 }
367 }
368 let mut phase_wall_pct: Vec<(String, f64)> = PHASES
369 .iter()
370 .enumerate()
371 .map(|(i, &name)| {
372 let pct = if total_wall > 0 {
373 100.0 * phase_tot[i] as f64 / total_wall as f64
374 } else {
375 0.0
376 };
377 (name.to_string(), pct)
378 })
379 .collect();
380 phase_wall_pct.sort_by(|a, b| b.1.total_cmp(&a.1));
381
382 let mut walls_us: Vec<u64> = records.iter().map(|r| r.wall_us).collect();
384 walls_us.sort_unstable();
385 let tot_wall_f = total_wall as f64;
386 let top_share = |frac: f64| -> f64 {
387 let k = ((walls_us.len() as f64) * frac).floor() as usize;
388 if k == 0 || tot_wall_f == 0.0 {
389 return 0.0;
390 }
391 let s: u128 = walls_us[walls_us.len() - k..]
392 .iter()
393 .map(|&w| w as u128)
394 .sum();
395 100.0 * s as f64 / tot_wall_f
396 };
397 let count_over = |secs: u64| walls_us.iter().filter(|&&w| w >= secs * 1_000_000).count() as u64;
398 let tail = TailStats {
399 top1pct_wall_share: top_share(0.01),
400 top5pct_wall_share: top_share(0.05),
401 over_30s: count_over(30),
402 over_60s: count_over(60),
403 over_120s: count_over(120),
404 over_180s: count_over(180),
405 };
406
407 let rss_over = |gib: u64| {
409 records
410 .iter()
411 .filter(|r| r.max_rss_kb >= gib * 1024 * 1024)
412 .count() as u64
413 };
414 let rss_buckets = RssBuckets {
415 over_2gib: rss_over(2),
416 over_3gib: rss_over(3),
417 over_4gib: rss_over(4),
418 };
419
420 let m_inv = records
422 .iter()
423 .fold(0u64, |a, r| a.saturating_add(r.math_parse_attempts));
424 let m_cnt = records
425 .iter()
426 .fold(0u64, |a, r| a.saturating_add(r.math_parse_count));
427 let math = MathStats {
428 formulae: total_formulae,
429 parse_invocations: m_inv,
430 parse_count: m_cnt,
431 parses_per_formula: if m_inv > 0 {
432 m_cnt as f64 / m_inv as f64
433 } else {
434 0.0
435 },
436 };
437
438 let mut by_wall: Vec<&TelemetryRecord> = records.iter().collect();
440 by_wall.sort_by_key(|r| std::cmp::Reverse(r.wall_us));
441 let mut dom: HashMap<&'static str, u64> = HashMap::new();
442 for r in by_wall.iter().take(50) {
443 if let Some((i, _)) = r
444 .phase_us
445 .iter()
446 .take(17)
447 .enumerate()
448 .max_by_key(|&(_, us)| *us)
449 {
450 *dom.entry(PHASES[i]).or_insert(0) += 1;
451 }
452 }
453 let mut slow_tail_dominant: Vec<(String, u64)> =
454 dom.into_iter().map(|(k, v)| (k.to_string(), v)).collect();
455 slow_tail_dominant.sort_by_key(|d| std::cmp::Reverse(d.1));
456
457 let profile = |name: &str| -> OutcomeWall {
459 let mut ms: Vec<u64> = records
460 .iter()
461 .filter(|r| bucket(r) == name)
462 .map(|r| r.wall_us / 1000)
463 .collect();
464 if ms.is_empty() {
465 return OutcomeWall::default();
466 }
467 ms.sort_unstable();
468 let sum: u128 = ms.iter().map(|&x| x as u128).sum();
469 let p = percentiles(&ms);
470 OutcomeWall {
471 n: ms.len(),
472 median_ms: p.p50,
473 mean_ms: (sum / ms.len() as u128) as u64,
474 p99_ms: p.p99,
475 }
476 };
477 let fatal_profile = profile("fatal");
478 let no_problem_profile = profile("no_problem");
479
480 let generated_unix = std::time::SystemTime::now()
481 .duration_since(std::time::UNIX_EPOCH)
482 .map(|elapsed| elapsed.as_secs() as i64)
483 .unwrap_or(0);
484
485 TelemetrySummary {
486 corpus: corpus.to_string(),
487 service: service.to_string(),
488 sample_count: records.len(),
489 skipped: 0,
490 outcome_counts,
491 wall_ms,
492 rss_mib,
493 phase_p99_ms,
494 phase_wall_pct,
495 tail,
496 rss_buckets,
497 math,
498 slow_tail_dominant,
499 fatal_profile,
500 no_problem_profile,
501 slowest,
502 highest_rss,
503 total_formulae,
504 total_graphics_assets,
505 total_output_bytes,
506 generated_unix,
507 }
508}
509
510fn read_chunk(
514 chunk: &[String],
515 service_name: &str,
516 sandbox_id: Option<i32>,
517) -> (Vec<TelemetryRecord>, usize) {
518 let mut records = Vec::new();
519 let mut skipped = 0usize;
520 for entry in chunk {
521 match result_archive_path(entry, service_name, sandbox_id) {
522 Some(path) => match read_telemetry_json(&path) {
523 Ok(mut record) => {
524 record.paper_id = crate::helpers::entry_document_name(entry);
529 records.push(record);
530 },
531 Err(_) => skipped += 1,
532 },
533 None => skipped += 1,
534 }
535 }
536 (records, skipped)
537}
538
539pub fn aggregate(
546 corpus_name: &str,
547 service_name: &str,
548 sandbox_id: Option<i32>,
549 entries: Vec<String>,
550) -> TelemetrySummary {
551 let total = entries.len();
552 if total == 0 {
553 return summarize(corpus_name, service_name, &[]);
554 }
555 let workers = available_parallelism()
557 .map(|n| n.get())
558 .unwrap_or(1)
559 .clamp(1, 16)
560 .min(total);
561 let chunk_size = total.div_ceil(workers);
562
563 let mut records: Vec<TelemetryRecord> = Vec::new();
564 let mut skipped = 0usize;
565 std::thread::scope(|scope| {
566 let handles: Vec<_> = entries
569 .chunks(chunk_size)
570 .map(|chunk| {
571 (
572 chunk.len(),
573 scope.spawn(move || read_chunk(chunk, service_name, sandbox_id)),
574 )
575 })
576 .collect();
577 for (chunk_len, handle) in handles {
578 match handle.join() {
579 Ok((chunk_records, chunk_skipped)) => {
580 records.extend(chunk_records);
581 skipped += chunk_skipped;
582 },
583 Err(_) => skipped += chunk_len,
584 }
585 }
586 });
587
588 let mut summary = summarize(corpus_name, service_name, &records);
589 summary.skipped = skipped;
590 summary
591}
592
593#[cfg(test)]
594mod tests {
595 use super::{PHASES, TelemetryRecord, percentiles, summarize};
596
597 fn record(paper_id: &str) -> TelemetryRecord {
599 TelemetryRecord {
600 paper_id: paper_id.to_string(),
601 phase_us: vec![0; PHASES.len()],
602 ..Default::default()
603 }
604 }
605
606 #[test]
607 fn nearest_rank_percentiles_on_a_known_vector() {
608 let values = [50, 20, 100, 40, 10, 70, 30, 90, 60, 80];
610 let p = percentiles(&values);
611 assert_eq!(p.p50, 50);
613 assert_eq!(p.p90, 90);
614 assert_eq!(p.p99, 100);
615 assert_eq!(p.max, 100);
616 let empty = percentiles(&[]);
618 assert_eq!((empty.p50, empty.p90, empty.p99, empty.max), (0, 0, 0, 0));
619 }
620
621 #[test]
622 fn summarize_buckets_percentiles_and_phase_length() {
623 let mut clean = record("clean"); clean.wall_us = 1_000_000; clean.max_rss_kb = 1_048_576; let mut warned = record("warned");
628 warned.warnings = 3;
629 warned.wall_us = 2_000_000; let mut errored = record("errored");
631 errored.errors = 1;
632 errored.category = "conversion_error".to_string();
633 errored.exit_code = 2;
634 errored.wall_us = 3_000_000; let mut fataled = record("fataled");
636 fataled.fatal_errors = 1;
637 fataled.category = "conversion_fatal".to_string();
638 fataled.exit_code = 3;
639 fataled.wall_us = 4_000_000; let records = [clean, warned, errored, fataled];
642 let summary = summarize("c", "s", &records);
643
644 assert_eq!(
646 summary.outcome_counts,
647 vec![
648 ("no_problem".to_string(), 1),
649 ("warning".to_string(), 1),
650 ("error".to_string(), 1),
651 ("fatal".to_string(), 1),
652 ]
653 );
654 assert_eq!(summary.wall_ms.p50, 2000);
656 assert_eq!(summary.wall_ms.p90, 4000);
657 assert_eq!(summary.wall_ms.p99, 4000);
658 assert_eq!(summary.wall_ms.max, 4000);
659 assert_eq!(summary.phase_p99_ms.len(), 17);
661 assert_eq!(summary.phase_p99_ms[0].0, "bootstrap");
662 assert_eq!(summary.slowest, Some(("fataled".to_string(), 4000)));
664 assert_eq!(summary.sample_count, 4);
665 assert_eq!(summary.skipped, 0);
666 }
667
668 #[test]
669 fn summarize_budget_tail_math_and_profiles() {
670 let mut records = Vec::new();
671 for i in 0..8 {
673 let mut r = record(&format!("ok{i}"));
674 r.wall_us = 1_000_000;
675 r.phase_us[1] = 700_000; r.phase_us[4] = 200_000; r.formulae = 10;
678 r.math_parse_attempts = 10;
679 r.math_parse_count = 13;
680 records.push(r);
681 }
682 let mut slow_fatal = record("slowfatal");
684 slow_fatal.wall_us = 40_000_000;
685 slow_fatal.phase_us[1] = 39_000_000;
686 slow_fatal.fatal_errors = 1;
687 records.push(slow_fatal);
688 let mut slow_warn = record("slowwarn");
690 slow_warn.wall_us = 65_000_000;
691 slow_warn.phase_us[4] = 60_000_000;
692 slow_warn.warnings = 1;
693 records.push(slow_warn);
694
695 let s = summarize("c", "s", &records);
696
697 assert_eq!(s.phase_wall_pct.len(), 17);
699 assert_eq!(s.phase_wall_pct[0].0, "math_parse");
700 assert!(s.phase_wall_pct.windows(2).all(|w| w[0].1 >= w[1].1));
701 let budget_sum: f64 = s.phase_wall_pct.iter().map(|(_, p)| p).sum();
702 assert!(
703 (90.0..=100.0).contains(&budget_sum),
704 "instrumented phases ≈ total wall"
705 );
706
707 assert_eq!(s.tail.over_30s, 2);
709 assert_eq!(s.tail.over_60s, 1);
710 assert_eq!(s.tail.over_180s, 0);
711
712 assert_eq!(s.math.formulae, 80);
714 assert!((s.math.parses_per_formula - 1.3).abs() < 1e-9);
715
716 assert_eq!(s.fatal_profile.n, 1);
718 assert_eq!(s.fatal_profile.median_ms, 40_000);
719 assert_eq!(s.no_problem_profile.n, 8);
720 assert_eq!(s.no_problem_profile.median_ms, 1_000);
721
722 assert!(!s.slow_tail_dominant.is_empty());
724 }
725
726 #[test]
727 fn minimal_failure_record_parses_and_summarizes() {
728 let json = r#"{"paper_id":"x","category":"conversion_fatal","exit_code":3}"#;
731 let record: TelemetryRecord = serde_json::from_str(json).expect("minimal record parses");
732 assert_eq!(record.paper_id, "x");
733 assert_eq!(record.exit_code, 3);
734 assert!(
735 record.phase_us.is_empty(),
736 "absent phase_us decodes to empty"
737 );
738
739 let summary = summarize("c", "s", std::slice::from_ref(&record));
742 assert_eq!(summary.phase_p99_ms.len(), 17);
743 assert!(summary.phase_p99_ms.iter().all(|(_, p99)| *p99 == 0));
744 assert_eq!(summary.outcome_counts[3], ("fatal".to_string(), 1));
745 }
746}