1use std::time::Instant;
9
10#[repr(u8)]
39#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
40pub enum Phase {
41 BoundaryStart,
42 SnapshotReady,
43 RequestBuilt,
44 FirstByte,
45 LastBlockEnd,
46 ToolsSpawned,
47 ToolsJoined,
48 BatchProposed,
49 WatermarkDurable,
50 BoundaryEnd,
51}
52
53pub const PHASE_COUNT: usize = 10;
55
56const CAPACITY: usize = 64;
59
60const DELTA_PAIRS: [(Phase, Phase); PHASE_COUNT - 1] = [
69 (Phase::BoundaryStart, Phase::SnapshotReady),
70 (Phase::SnapshotReady, Phase::RequestBuilt),
71 (Phase::RequestBuilt, Phase::FirstByte),
72 (Phase::FirstByte, Phase::LastBlockEnd),
73 (Phase::LastBlockEnd, Phase::ToolsSpawned),
74 (Phase::ToolsSpawned, Phase::ToolsJoined),
75 (Phase::ToolsJoined, Phase::BatchProposed),
76 (Phase::BatchProposed, Phase::WatermarkDurable),
77 (Phase::WatermarkDurable, Phase::BoundaryEnd),
78];
79
80fn now_nanos() -> u64 {
85 static START: std::sync::OnceLock<Instant> = std::sync::OnceLock::new();
86 let start = START.get_or_init(Instant::now);
87 start.elapsed().as_nanos() as u64
88}
89
90fn width(from: u64, to: u64) -> u64 {
95 if from == 0 || to == 0 {
96 0
97 } else {
98 to.saturating_sub(from)
99 }
100}
101
102fn sample_overhead(raw: &[u64; PHASE_COUNT]) -> u64 {
106 let total = width(
107 raw[Phase::BoundaryStart as usize],
108 raw[Phase::BoundaryEnd as usize],
109 );
110 let stream = width(
111 raw[Phase::FirstByte as usize],
112 raw[Phase::LastBlockEnd as usize],
113 );
114 let tools = width(
115 raw[Phase::ToolsSpawned as usize],
116 raw[Phase::ToolsJoined as usize],
117 );
118 total.saturating_sub(stream).saturating_sub(tools)
119}
120
121fn percentile(sorted: &[u64], p: u64) -> u64 {
125 if sorted.is_empty() {
126 return 0;
127 }
128 let n = sorted.len() as u64;
129 let rank = (p * n).div_ceil(100);
132 let idx = rank.clamp(1, n) - 1;
133 sorted[idx as usize]
134}
135
136#[derive(Debug, Clone, serde::Serialize)]
138pub struct PhaseDeltaSummary {
139 pub from: Phase,
140 pub to: Phase,
141 pub p50_ns: u64,
142 pub p99_ns: u64,
143}
144
145#[derive(Debug, Clone, serde::Serialize)]
148pub struct LedgerSummary {
149 pub sample_count: usize,
152 pub overhead_p50_ns: u64,
153 pub overhead_p99_ns: u64,
154 pub phase_deltas: Vec<PhaseDeltaSummary>,
155 pub max_rss_bytes: u64,
156 pub samples: Vec<[u64; PHASE_COUNT]>,
161}
162
163pub struct LoopLedger {
166 samples: [[u64; PHASE_COUNT]; CAPACITY],
167 len: usize,
169 current: usize,
172}
173
174impl Default for LoopLedger {
175 fn default() -> Self {
176 Self::new()
177 }
178}
179
180impl LoopLedger {
181 pub fn new() -> Self {
182 Self {
183 samples: [[0; PHASE_COUNT]; CAPACITY],
184 len: 0,
185 current: CAPACITY,
186 }
187 }
188
189 pub fn start_sample(&mut self) {
193 if self.len < CAPACITY {
194 self.current = self.len;
195 self.len += 1;
196 } else {
197 self.current = CAPACITY;
198 }
199 }
200
201 pub fn stamp(&mut self, phase: Phase) {
205 if self.current >= CAPACITY {
206 return;
207 }
208 let slot = &mut self.samples[self.current][phase as usize];
209 if *slot == 0 {
210 *slot = now_nanos();
211 }
212 }
213
214 pub fn restamp(&mut self, phase: Phase) {
226 if self.current >= CAPACITY {
227 return;
228 }
229 self.samples[self.current][phase as usize] = now_nanos();
230 }
231
232 pub fn summary(&self, max_rss_bytes: u64) -> LedgerSummary {
236 let used = &self.samples[..self.len];
237 let mut overheads: Vec<u64> = used.iter().map(sample_overhead).collect();
238 overheads.sort_unstable();
239 let phase_deltas = DELTA_PAIRS
240 .iter()
241 .map(|&(from, to)| {
242 let mut deltas: Vec<u64> = used
243 .iter()
244 .map(|s| width(s[from as usize], s[to as usize]))
245 .collect();
246 deltas.sort_unstable();
247 PhaseDeltaSummary {
248 from,
249 to,
250 p50_ns: percentile(&deltas, 50),
251 p99_ns: percentile(&deltas, 99),
252 }
253 })
254 .collect();
255 LedgerSummary {
256 sample_count: self.len,
257 overhead_p50_ns: percentile(&overheads, 50),
258 overhead_p99_ns: percentile(&overheads, 99),
259 phase_deltas,
260 max_rss_bytes,
261 samples: used.to_vec(),
262 }
263 }
264}
265
266pub fn max_rss_bytes() -> u64 {
270 let mut usage: libc::rusage = unsafe { std::mem::zeroed() };
271 if unsafe { libc::getrusage(libc::RUSAGE_SELF, &mut usage) } != 0 {
272 return 0;
273 }
274 normalize_rss(usage.ru_maxrss.max(0) as u64)
275}
276
277#[cfg(target_os = "macos")]
278fn normalize_rss(raw: u64) -> u64 {
279 raw
280}
281
282#[cfg(not(target_os = "macos"))]
283fn normalize_rss(raw: u64) -> u64 {
284 raw * 1024
285}
286
287#[cfg(test)]
288mod tests {
289 use super::*;
290
291 fn raw_with(pairs: &[(Phase, u64)]) -> [u64; PHASE_COUNT] {
292 let mut raw = [0u64; PHASE_COUNT];
293 for &(phase, ns) in pairs {
294 raw[phase as usize] = ns;
295 }
296 raw
297 }
298
299 #[test]
300 fn sample_overhead_subtracts_stream_and_tools() {
301 let raw = raw_with(&[
302 (Phase::BoundaryStart, 1_000),
303 (Phase::SnapshotReady, 1_100),
304 (Phase::RequestBuilt, 1_200),
305 (Phase::FirstByte, 1_300),
306 (Phase::LastBlockEnd, 1_800), (Phase::ToolsSpawned, 1_850),
308 (Phase::ToolsJoined, 2_050), (Phase::BatchProposed, 2_100),
310 (Phase::WatermarkDurable, 2_150),
311 (Phase::BoundaryEnd, 2_200), ]);
313 assert_eq!(sample_overhead(&raw), 500);
315 }
316
317 #[test]
318 fn sample_overhead_treats_absent_phases_as_zero_width() {
319 let raw = raw_with(&[(Phase::BoundaryStart, 1_000), (Phase::BoundaryEnd, 1_500)]);
323 assert_eq!(sample_overhead(&raw), 500);
324 }
325
326 #[test]
327 fn width_never_underflows_on_an_out_of_order_pair() {
328 assert_eq!(width(500, 100), 0);
329 }
330
331 #[test]
332 fn percentile_is_nearest_rank_on_sorted_input() {
333 let sorted = [10u64, 20, 30, 40, 50];
334 assert_eq!(percentile(&sorted, 50), 30);
335 assert_eq!(percentile(&sorted, 99), 50);
336 assert_eq!(percentile(&sorted, 0), 10);
337 }
338
339 #[test]
340 fn percentile_of_empty_input_is_zero() {
341 assert_eq!(percentile(&[], 50), 0);
342 }
343
344 #[test]
345 fn start_sample_then_stamp_targets_a_fresh_slot_each_time() {
346 let mut ledger = LoopLedger::new();
347 ledger.start_sample();
348 ledger.stamp(Phase::BoundaryStart);
349 ledger.start_sample();
350 ledger.stamp(Phase::BoundaryStart);
351 let report = ledger.summary(0);
352 assert_eq!(report.sample_count, 2);
353 assert_ne!(
354 report.samples[0][Phase::BoundaryStart as usize],
355 0,
356 "sample 0 must have its own stamp"
357 );
358 assert_ne!(
359 report.samples[1][Phase::BoundaryStart as usize],
360 0,
361 "sample 1 must have its own stamp"
362 );
363 }
364
365 #[test]
366 fn stamp_before_any_start_sample_is_a_noop() {
367 let mut ledger = LoopLedger::new();
368 ledger.stamp(Phase::BoundaryStart); let report = ledger.summary(0);
370 assert_eq!(report.sample_count, 0, "no sample was ever started");
371 }
372
373 #[test]
374 fn first_stamp_wins() {
375 let mut ledger = LoopLedger::new();
376 ledger.start_sample();
377 ledger.stamp(Phase::BoundaryStart);
378 let first = ledger.summary(0).samples[0][Phase::BoundaryStart as usize];
379 std::thread::sleep(std::time::Duration::from_micros(50));
382 ledger.stamp(Phase::BoundaryStart);
383 let second = ledger.summary(0).samples[0][Phase::BoundaryStart as usize];
384 assert_eq!(
385 first, second,
386 "the second stamp must not overwrite the first"
387 );
388 }
389
390 #[test]
391 fn restamp_overwrites_with_the_latest_value() {
392 let mut ledger = LoopLedger::new();
398 ledger.start_sample();
399 ledger.stamp(Phase::BatchProposed);
400 let first = ledger.summary(0).samples[0][Phase::BatchProposed as usize];
401 std::thread::sleep(std::time::Duration::from_micros(50));
402 ledger.restamp(Phase::BatchProposed);
403 let second = ledger.summary(0).samples[0][Phase::BatchProposed as usize];
404 assert!(
405 second > first,
406 "restamp must overwrite with a later value, got first={first} second={second}"
407 );
408 }
409
410 #[test]
411 fn restamp_before_any_start_sample_is_a_noop() {
412 let mut ledger = LoopLedger::new();
413 ledger.restamp(Phase::BatchProposed); let report = ledger.summary(0);
415 assert_eq!(report.sample_count, 0, "no sample was ever started");
416 }
417
418 #[test]
419 fn a_phase_may_be_legitimately_absent() {
420 let mut ledger = LoopLedger::new();
421 ledger.start_sample();
422 ledger.stamp(Phase::BoundaryStart);
423 ledger.stamp(Phase::BoundaryEnd);
424 let report = ledger.summary(0);
426 assert_eq!(report.samples[0][Phase::ToolsSpawned as usize], 0);
427 assert_eq!(report.samples[0][Phase::ToolsJoined as usize], 0);
428 }
429
430 #[test]
431 fn overflow_past_capacity_is_safe_and_drops() {
432 let mut ledger = LoopLedger::new();
433 for _ in 0..(CAPACITY + 36) {
434 ledger.start_sample();
435 ledger.stamp(Phase::BoundaryStart);
436 ledger.stamp(Phase::BoundaryEnd);
437 }
438 let report = ledger.summary(0);
439 assert_eq!(
440 report.sample_count, CAPACITY,
441 "overflow must cap, not grow, the stored sample count"
442 );
443 assert_eq!(report.samples.len(), CAPACITY);
444 }
445
446 #[test]
447 fn summary_reports_sample_count_and_passes_through_max_rss() {
448 let mut ledger = LoopLedger::new();
449 ledger.start_sample();
450 ledger.stamp(Phase::BoundaryStart);
451 ledger.stamp(Phase::BoundaryEnd);
452 let report = ledger.summary(123_456);
453 assert_eq!(report.sample_count, 1);
454 assert_eq!(report.max_rss_bytes, 123_456);
455 }
456
457 #[test]
458 fn summary_overhead_percentiles_span_every_sample() {
459 let mut samples = [[0u64; PHASE_COUNT]; CAPACITY];
466 samples[0] = raw_with(&[(Phase::BoundaryStart, 1_000), (Phase::BoundaryEnd, 1_100)]);
467 samples[1] = raw_with(&[(Phase::BoundaryStart, 1_000), (Phase::BoundaryEnd, 1_500)]);
468 let ledger = LoopLedger {
469 samples,
470 len: 2,
471 current: 1,
472 };
473 let report = ledger.summary(0);
474 assert_eq!(report.overhead_p50_ns, 100);
475 assert_eq!(report.overhead_p99_ns, 500);
476 }
477
478 #[test]
479 fn phase_deltas_cover_every_consecutive_pair_in_declaration_order() {
480 let ledger = LoopLedger::new();
481 let report = ledger.summary(0);
482 assert_eq!(report.phase_deltas.len(), PHASE_COUNT - 1);
483 assert_eq!(report.phase_deltas[0].from, Phase::BoundaryStart);
484 assert_eq!(report.phase_deltas[0].to, Phase::SnapshotReady);
485 assert_eq!(report.phase_deltas[PHASE_COUNT - 2].to, Phase::BoundaryEnd);
486 }
487
488 #[test]
489 fn now_nanos_is_monotonically_nondecreasing() {
490 let a = now_nanos();
491 let b = now_nanos();
492 assert!(b >= a);
493 }
494
495 #[test]
496 fn max_rss_bytes_returns_a_plausible_reading() {
497 assert!(max_rss_bytes() > 0);
500 }
501}