1use std::collections::HashMap;
15use std::sync::Mutex;
16
17#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
19pub enum Backend {
20 Sidecar,
22 Rust,
24}
25
26impl Backend {
27 pub fn parse(raw: &str) -> Option<Self> {
31 match raw {
32 "sidecar" => Some(Backend::Sidecar),
33 "rust" => Some(Backend::Rust),
34 _ => None,
35 }
36 }
37
38 pub fn as_str(self) -> &'static str {
39 match self {
40 Backend::Sidecar => "sidecar",
41 Backend::Rust => "rust",
42 }
43 }
44}
45
46const WINDOW: usize = 1024;
53
54#[derive(Debug, Default, Clone)]
56pub struct Stats {
57 ok_micros: Vec<u64>,
59 err_micros: Vec<u64>,
62 pub ok_count: u64,
64 pub err_count: u64,
65}
66
67impl Stats {
68 fn record(&mut self, micros: u64, ok: bool) {
69 let (samples, count) = if ok {
70 (&mut self.ok_micros, &mut self.ok_count)
71 } else {
72 (&mut self.err_micros, &mut self.err_count)
73 };
74 *count += 1;
75 if samples.len() == WINDOW {
76 samples.remove(0);
77 }
78 samples.push(micros);
79 }
80
81 pub fn percentile(&self, p: f64) -> Option<u64> {
87 percentile_of(&self.ok_micros, p)
88 }
89
90 pub fn error_percentile(&self, p: f64) -> Option<u64> {
92 percentile_of(&self.err_micros, p)
93 }
94}
95
96fn percentile_of(samples: &[u64], p: f64) -> Option<u64> {
98 if samples.is_empty() {
99 return None;
100 }
101 let mut sorted = samples.to_vec();
102 sorted.sort_unstable();
103 let rank = ((p / 100.0) * sorted.len() as f64).ceil().max(1.0) as usize;
105 Some(sorted[rank.min(sorted.len()) - 1])
106}
107
108#[derive(Debug, Clone)]
110pub struct Row {
111 pub backend: Backend,
112 pub op: String,
113 pub stats: Stats,
114}
115
116const MAX_PENDING: usize = 16_384;
122
123#[derive(Debug, Default)]
125pub struct RepoMetrics {
126 stats: Mutex<HashMap<(Backend, &'static str), Stats>>,
127 pending: Mutex<Vec<Sample>>,
134}
135
136#[derive(Debug, Clone, Copy)]
138struct Sample {
139 backend: Backend,
140 op: &'static str,
141 micros: u64,
142 ok: bool,
143}
144
145impl RepoMetrics {
146 pub fn new() -> Self {
147 Self::default()
148 }
149
150 pub fn record(&self, backend: Backend, op: &'static str, micros: u64, ok: bool) {
152 self.stats
157 .lock()
158 .unwrap_or_else(|p| p.into_inner())
159 .entry((backend, op))
160 .or_default()
161 .record(micros, ok);
162
163 let mut pending = self.pending.lock().unwrap_or_else(|p| p.into_inner());
164 if pending.len() >= MAX_PENDING {
165 pending.remove(0);
166 }
167 pending.push(Sample {
168 backend,
169 op,
170 micros,
171 ok,
172 });
173 }
174
175 pub fn snapshot(&self) -> Vec<Row> {
180 let stats = self.stats.lock().unwrap_or_else(|p| p.into_inner());
181 let mut rows: Vec<Row> = stats
182 .iter()
183 .map(|((backend, op), stats)| Row {
184 backend: *backend,
185 op: (*op).to_string(),
186 stats: stats.clone(),
187 })
188 .collect();
189 rows.sort_by(|a, b| (a.backend, &a.op).cmp(&(b.backend, &b.op)));
190 rows
191 }
192
193 fn drain(&self) -> Vec<Sample> {
195 std::mem::take(&mut *self.pending.lock().unwrap_or_else(|p| p.into_inner()))
196 }
197}
198
199pub async fn timed<T, E, F>(
205 metrics: &RepoMetrics,
206 backend: Backend,
207 op: &'static str,
208 call: F,
209) -> Result<T, E>
210where
211 F: std::future::Future<Output = Result<T, E>>,
212{
213 let started = std::time::Instant::now();
214 let result = call.await;
215 metrics.record(
216 backend,
217 op,
218 started.elapsed().as_micros() as u64,
219 result.is_ok(),
220 );
221 result
222}
223
224pub async fn flush(metrics: &RepoMetrics, pool: &sqlx::SqlitePool, now: i64) -> anyhow::Result<()> {
230 let samples = metrics.drain();
231 if samples.is_empty() {
232 return Ok(());
233 }
234
235 let mut tx = pool.begin().await?;
236 for sample in &samples {
237 sqlx::query(
238 "INSERT INTO repo_timing (backend, op, micros, ok, at) VALUES (?1, ?2, ?3, ?4, ?5)",
239 )
240 .bind(sample.backend.as_str())
241 .bind(sample.op)
242 .bind(sample.micros as i64)
243 .bind(i64::from(sample.ok))
244 .bind(now)
245 .execute(&mut *tx)
246 .await?;
247
248 sqlx::query(
253 "INSERT INTO repo_timing_total (backend, op, ok_count, err_count) \
254 VALUES (?1, ?2, ?3, ?4) \
255 ON CONFLICT(backend, op) DO UPDATE SET \
256 ok_count = ok_count + excluded.ok_count, \
257 err_count = err_count + excluded.err_count",
258 )
259 .bind(sample.backend.as_str())
260 .bind(sample.op)
261 .bind(i64::from(sample.ok))
262 .bind(i64::from(!sample.ok))
263 .execute(&mut *tx)
264 .await?;
265 }
266
267 for (backend, op, ok) in distinct_keys(&samples) {
271 sqlx::query(
272 "DELETE FROM repo_timing WHERE backend = ?1 AND op = ?2 AND ok = ?3 AND id NOT IN \
273 (SELECT id FROM repo_timing WHERE backend = ?1 AND op = ?2 AND ok = ?3 \
274 ORDER BY id DESC LIMIT ?4)",
275 )
276 .bind(backend.as_str())
277 .bind(op)
278 .bind(i64::from(ok))
279 .bind(WINDOW as i64)
280 .execute(&mut *tx)
281 .await?;
282 }
283
284 tx.commit().await?;
285 Ok(())
286}
287
288fn distinct_keys(samples: &[Sample]) -> Vec<(Backend, &'static str, bool)> {
290 let mut keys: Vec<(Backend, &'static str, bool)> =
291 samples.iter().map(|s| (s.backend, s.op, s.ok)).collect();
292 keys.sort();
293 keys.dedup();
294 keys
295}
296
297pub async fn persisted_rows(pool: &sqlx::SqlitePool) -> anyhow::Result<Vec<Row>> {
302 let totals: Vec<(String, String, i64, i64)> = sqlx::query_as(
303 "SELECT backend, op, ok_count, err_count FROM repo_timing_total ORDER BY backend, op",
304 )
305 .fetch_all(pool)
306 .await?;
307
308 let mut rows = Vec::with_capacity(totals.len());
309 for (backend, op, ok_count, err_count) in totals {
310 let Some(backend) = Backend::parse(&backend) else {
311 continue;
314 };
315 let samples: Vec<(i64, i64)> =
316 sqlx::query_as("SELECT micros, ok FROM repo_timing WHERE backend = ?1 AND op = ?2")
317 .bind(backend.as_str())
318 .bind(&op)
319 .fetch_all(pool)
320 .await?;
321
322 let mut stats = Stats {
323 ok_count: ok_count as u64,
324 err_count: err_count as u64,
325 ..Stats::default()
326 };
327 for (micros, ok) in samples {
328 if ok == 1 {
329 stats.ok_micros.push(micros as u64);
330 } else {
331 stats.err_micros.push(micros as u64);
332 }
333 }
334 rows.push(Row { backend, op, stats });
335 }
336 Ok(rows)
337}
338
339pub fn render(rows: &[Row]) -> String {
349 let mut out = String::from(
350 "backend operation ok p50ms p95ms err errp50ms\n",
351 );
352 if rows.is_empty() {
353 out.push_str("(no repo operations recorded yet)\n");
354 return out;
355 }
356 for row in rows {
357 out.push_str(&format!(
358 "{:<8} {:<28} {:>4} {:>7} {:>7} {:>6} {:>9}\n",
359 row.backend.as_str(),
360 row.op,
361 row.stats.ok_count,
362 render_micros(row.stats.percentile(50.0)),
363 render_micros(row.stats.percentile(95.0)),
364 row.stats.err_count,
365 render_micros(row.stats.error_percentile(50.0)),
366 ));
367 }
368 out
369}
370
371fn render_micros(micros: Option<u64>) -> String {
374 match micros {
375 Some(micros) => format!("{:.1}", micros as f64 / 1000.0),
376 None => "-".to_string(),
377 }
378}
379
380#[cfg(test)]
381mod tests {
382 use super::*;
383
384 fn stats_of(ok: &[u64], err: &[u64]) -> Stats {
385 let mut stats = Stats::default();
386 for micros in ok {
387 stats.record(*micros, true);
388 }
389 for micros in err {
390 stats.record(*micros, false);
391 }
392 stats
393 }
394
395 #[test]
398 fn percentiles_come_from_the_recorded_samples() {
399 let stats = stats_of(&[10, 20, 30, 40, 50, 60, 70, 80, 90, 100], &[]);
400 assert_eq!(stats.percentile(50.0), Some(50));
401 assert_eq!(stats.percentile(95.0), Some(100));
402 assert_eq!(stats.percentile(100.0), Some(100));
403 assert_eq!(stats.percentile(0.0), Some(10));
405 }
406
407 #[test]
414 fn failures_are_counted_but_kept_out_of_the_success_percentiles() {
415 let stats = stats_of(&[1000; 10], &[1; 90]);
417
418 assert_eq!(
419 stats.percentile(50.0),
420 Some(1000),
421 "fast failures dragged the success percentile down, which is how a \
422 broken backend passes for a fast one"
423 );
424 assert_eq!(stats.ok_count, 10);
425 assert_eq!(stats.err_count, 90);
426 assert_eq!(stats.error_percentile(50.0), Some(1));
427 }
428
429 #[test]
432 fn an_unexercised_operation_reports_nothing_rather_than_zero() {
433 let stats = Stats::default();
434 assert_eq!(stats.percentile(50.0), None);
435 assert_eq!(stats.error_percentile(50.0), None);
436 }
437
438 #[test]
441 fn an_operation_that_only_ever_fails_still_reports_its_failures() {
442 let stats = stats_of(&[], &[5, 7, 9]);
443 assert_eq!(stats.percentile(50.0), None);
444 assert_eq!(stats.err_count, 3);
445 assert_eq!(stats.error_percentile(50.0), Some(7));
446 }
447
448 #[test]
452 fn the_window_keeps_the_most_recent_samples() {
453 let mut stats = Stats::default();
454 for i in 0..(WINDOW as u64 + 10) {
455 stats.record(i, true);
456 }
457 assert_eq!(stats.ok_count, WINDOW as u64 + 10, "the total counts all");
458 assert_eq!(
459 stats.percentile(0.0),
460 Some(10),
461 "the oldest samples should have aged out of the window"
462 );
463 assert_eq!(stats.percentile(100.0), Some(WINDOW as u64 + 9));
464 }
465
466 #[test]
469 fn the_two_backends_are_recorded_side_by_side() {
470 let metrics = RepoMetrics::new();
471 metrics.record(Backend::Sidecar, "list_subscriptions", 5_000, true);
472 metrics.record(Backend::Rust, "list_subscriptions", 2_000, true);
473
474 let snapshot = metrics.snapshot();
475 assert_eq!(snapshot.len(), 2);
476 let sidecar = snapshot
477 .iter()
478 .find(|r| r.backend == Backend::Sidecar)
479 .expect("the sidecar row");
480 let rust = snapshot
481 .iter()
482 .find(|r| r.backend == Backend::Rust)
483 .expect("the rust row");
484 assert_eq!(
485 sidecar.op, rust.op,
486 "the same operation name, or the rows cannot be compared"
487 );
488 assert_eq!(sidecar.stats.percentile(50.0), Some(5_000));
489 assert_eq!(rust.stats.percentile(50.0), Some(2_000));
490 }
491
492 #[test]
496 fn an_absent_measurement_renders_as_a_dash_rather_than_zero() {
497 let metrics = RepoMetrics::new();
498 metrics.record(Backend::Rust, "list_folders", 3_000, false);
499 let table = render(&metrics.snapshot());
500
501 assert!(
502 !table.contains("0.0"),
503 "a never-succeeded operation rendered as 0.0ms, which reads as instant:\n{table}"
504 );
505 assert!(
506 table.contains('-'),
507 "expected a dash for the absent p50:\n{table}"
508 );
509 assert!(table.contains("list_folders"), "the row vanished:\n{table}");
511 assert!(
512 table.contains("3.0"),
513 "the failure latency is missing:\n{table}"
514 );
515 }
516
517 #[test]
520 fn an_empty_snapshot_says_so() {
521 assert!(render(&[]).contains("no repo operations recorded"));
522 }
523
524 #[tokio::test]
528 async fn timed_records_success_and_failure_distinctly() {
529 let metrics = RepoMetrics::new();
530 let _: Result<(), ()> = timed(&metrics, Backend::Rust, "op", async { Ok(()) }).await;
531 let _: Result<(), ()> = timed(&metrics, Backend::Rust, "op", async { Err(()) }).await;
532
533 let snapshot = metrics.snapshot();
534 assert_eq!(snapshot[0].stats.ok_count, 1);
535 assert_eq!(snapshot[0].stats.err_count, 1);
536 }
537
538 async fn pool() -> sqlx::SqlitePool {
541 crate::store::init_url("sqlite::memory:").await.unwrap()
542 }
543
544 #[tokio::test]
552 async fn a_backend_flip_keeps_the_earlier_backends_rows() {
553 let pool = pool().await;
554
555 let before = RepoMetrics::new();
557 for _ in 0..5 {
558 before.record(Backend::Sidecar, "list_subscriptions_sorted", 9_000, true);
559 }
560 flush(&before, &pool, 1_700_000_000).await.unwrap();
561
562 let after = RepoMetrics::new();
564 for _ in 0..5 {
565 after.record(Backend::Rust, "list_subscriptions_sorted", 3_000, true);
566 }
567 flush(&after, &pool, 1_700_000_100).await.unwrap();
568
569 assert_eq!(
570 after.snapshot().len(),
571 1,
572 "in-process memory only ever holds the running backend"
573 );
574
575 let rows = persisted_rows(&pool).await.unwrap();
576 assert_eq!(
577 rows.len(),
578 2,
579 "both backends must survive the flip: {rows:?}"
580 );
581 let sidecar = rows.iter().find(|r| r.backend == Backend::Sidecar).unwrap();
582 let rust = rows.iter().find(|r| r.backend == Backend::Rust).unwrap();
583 assert_eq!(sidecar.stats.percentile(50.0), Some(9_000));
584 assert_eq!(rust.stats.percentile(50.0), Some(3_000));
585 assert_eq!(sidecar.op, rust.op, "rows must be comparable by operation");
586 }
587
588 #[tokio::test]
592 async fn persisted_failures_stay_out_of_the_success_percentiles() {
593 let pool = pool().await;
594 let metrics = RepoMetrics::new();
595 for _ in 0..10 {
596 metrics.record(Backend::Rust, "add_subscription", 1_000, true);
597 }
598 for _ in 0..90 {
599 metrics.record(Backend::Rust, "add_subscription", 1, false);
600 }
601 flush(&metrics, &pool, 1_700_000_000).await.unwrap();
602
603 let rows = persisted_rows(&pool).await.unwrap();
604 let stats = &rows[0].stats;
605 assert_eq!(
606 stats.percentile(50.0),
607 Some(1_000),
608 "fast failures dragged the persisted success percentile down"
609 );
610 assert_eq!(stats.ok_count, 10);
611 assert_eq!(stats.err_count, 90);
612 }
613
614 #[tokio::test]
620 async fn pruning_bounds_the_window_without_losing_the_totals() {
621 let pool = pool().await;
622 let metrics = RepoMetrics::new();
623 let total = WINDOW + 250;
624 for i in 0..total {
625 metrics.record(Backend::Rust, "list_folders_sorted", i as u64 + 1, true);
626 }
627 flush(&metrics, &pool, 1_700_000_000).await.unwrap();
628
629 let kept: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM repo_timing")
630 .fetch_one(&pool)
631 .await
632 .unwrap();
633 assert_eq!(kept, WINDOW as i64, "the window is not bounded");
634
635 let rows = persisted_rows(&pool).await.unwrap();
636 assert_eq!(
637 rows[0].stats.ok_count, total as u64,
638 "pruning ate the all-time count"
639 );
640 assert_eq!(rows[0].stats.percentile(0.0), Some(251));
643 }
644
645 #[tokio::test]
648 async fn a_burst_of_failures_does_not_evict_the_successes() {
649 let pool = pool().await;
650 let metrics = RepoMetrics::new();
651 metrics.record(Backend::Rust, "put_read_state", 5_000, true);
652 for _ in 0..(WINDOW + 100) {
653 metrics.record(Backend::Rust, "put_read_state", 2, false);
654 }
655 flush(&metrics, &pool, 1_700_000_000).await.unwrap();
656
657 let rows = persisted_rows(&pool).await.unwrap();
658 assert_eq!(
659 rows[0].stats.percentile(50.0),
660 Some(5_000),
661 "the only success was evicted by a flood of failures"
662 );
663 }
664
665 #[tokio::test]
667 async fn flushing_twice_does_not_double_count() {
668 let pool = pool().await;
669 let metrics = RepoMetrics::new();
670 metrics.record(Backend::Rust, "remove_saved", 1_000, true);
671 flush(&metrics, &pool, 1_700_000_000).await.unwrap();
672 flush(&metrics, &pool, 1_700_000_001).await.unwrap();
673
674 let rows = persisted_rows(&pool).await.unwrap();
675 assert_eq!(rows[0].stats.ok_count, 1, "the sample was counted twice");
676 }
677
678 #[test]
681 fn the_pending_buffer_is_bounded() {
682 let metrics = RepoMetrics::new();
683 for _ in 0..(MAX_PENDING + 100) {
684 metrics.record(Backend::Rust, "op", 1, true);
685 }
686 assert!(
687 metrics.pending.lock().unwrap().len() <= MAX_PENDING,
688 "the write buffer grew past its bound"
689 );
690 }
691
692 #[test]
696 fn an_unknown_backend_name_is_not_guessed() {
697 assert_eq!(Backend::parse("sidecar"), Some(Backend::Sidecar));
698 assert_eq!(Backend::parse("rust"), Some(Backend::Rust));
699 for unknown in ["", "RUST", "postgres", "rust "] {
700 assert_eq!(Backend::parse(unknown), None, "guessed at {unknown:?}");
701 }
702 }
703
704 #[test]
707 fn the_rendered_table_distinguishes_p50_from_p95() {
708 let metrics = RepoMetrics::new();
709 for _ in 0..90 {
711 metrics.record(Backend::Rust, "op", 1_000, true);
712 }
713 for _ in 0..10 {
714 metrics.record(Backend::Rust, "op", 900_000, true);
715 }
716 let table = render(&metrics.snapshot());
717 assert!(table.contains("1.0"), "p50 missing from:\n{table}");
718 assert!(
719 table.contains("900.0"),
720 "the p95 column is not showing p95:\n{table}"
721 );
722 }
723}