Skip to main content

feather_reader/
metrics.rs

1//! Side-by-side latency for the two repo backends.
2//!
3//! The cutover runs one backend at a time behind a flag, so the comparison is
4//! across a flip rather than within a request. That makes the *shape* of the
5//! measurement the thing to get right: both paths are timed at the same
6//! boundary — the call site in `web.rs`, which is what a user's request
7//! actually waits on — by one wrapper rather than two hand-placed timers, so a
8//! difference in the numbers is a difference in the backends.
9//!
10//! **Failures are timed separately from successes.** A backend that is fast
11//! because it is erroring out early would otherwise look like a win, and that
12//! is precisely the regression a cutover needs to catch.
13
14use std::collections::HashMap;
15use std::sync::Mutex;
16
17/// Which repo implementation served a call.
18#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
19pub enum Backend {
20    /// The Node `@atproto/oauth-client` sidecar.
21    Sidecar,
22    /// The Rust-native OAuth client.
23    Rust,
24}
25
26impl Backend {
27    /// Read back a persisted row's backend name. `None` for anything this
28    /// version does not know, so a row from a newer build is skipped rather
29    /// than failing the whole table.
30    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
46/// How many recent samples are kept per (backend, operation).
47///
48/// Exact percentiles over a bounded recent window, rather than approximate ones
49/// over all time: a cutover comparison cares about how the backend is behaving
50/// *now*, and a window that includes the first cold-start minute forever would
51/// hide a later improvement.
52const WINDOW: usize = 1024;
53
54/// Timings for one (backend, operation) pair.
55#[derive(Debug, Default, Clone)]
56pub struct Stats {
57    /// Durations of SUCCESSFUL calls, most recent `WINDOW` kept.
58    ok_micros: Vec<u64>,
59    /// Durations of FAILED calls, kept apart so they cannot flatter a
60    /// percentile.
61    err_micros: Vec<u64>,
62    /// Totals over all time, unaffected by the window.
63    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    /// The `p`th percentile of SUCCESSFUL calls, in microseconds.
82    ///
83    /// `None` when nothing has succeeded — reporting `0` for an operation that
84    /// has never completed would read as "instant" on exactly the dashboard
85    /// someone uses to decide a cutover is safe.
86    pub fn percentile(&self, p: f64) -> Option<u64> {
87        percentile_of(&self.ok_micros, p)
88    }
89
90    /// The same percentile over FAILED calls, for reading beside the successes.
91    pub fn error_percentile(&self, p: f64) -> Option<u64> {
92        percentile_of(&self.err_micros, p)
93    }
94}
95
96/// Nearest-rank percentile over a copy of the samples.
97fn 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    // Nearest-rank: ceil(p/100 * n), clamped into the slice.
104    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/// One rendered line: a backend, an operation, and its timings.
109#[derive(Debug, Clone)]
110pub struct Row {
111    pub backend: Backend,
112    pub op: String,
113    pub stats: Stats,
114}
115
116/// Cap on unflushed samples held in memory.
117///
118/// A backstop, not a tuning knob: if the flusher ever stops, this bounds what
119/// the buffer can grow to rather than letting a metrics buffer take the process
120/// down. Oldest are dropped, because the recent ones are the interesting ones.
121const MAX_PENDING: usize = 16_384;
122
123/// The process-wide table. Cheap to clone into handlers via `AppState`.
124#[derive(Debug, Default)]
125pub struct RepoMetrics {
126    stats: Mutex<HashMap<(Backend, &'static str), Stats>>,
127    /// Samples not yet written to SQLite.
128    ///
129    /// Recording buffers rather than writing, because a synchronous insert on
130    /// the request path would put the instrument inside the thing it measures --
131    /// every repo call would carry a database write that the sidecar path never
132    /// had, and the comparison would be of the instrumentation.
133    pending: Mutex<Vec<Sample>>,
134}
135
136/// A buffered sample awaiting its write.
137#[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    /// Record one call: into the live window, and into the write buffer.
151    pub fn record(&self, backend: Backend, op: &'static str, micros: u64, ok: bool) {
152        // Same reasoning as the pinned-client cache: this is on the repo-call
153        // hot path, so panicking on a poisoned lock would convert one panic
154        // anywhere into a permanently broken instance. Metrics are the least
155        // important thing here and must never be the thing that kills it.
156        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    /// This PROCESS's samples, sorted for stable rendering.
176    ///
177    /// Only ever one backend's rows, since a flip is a restart. Use
178    /// [`persisted_rows`] for the cross-flip comparison.
179    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    /// Take everything buffered, leaving the buffer empty.
194    fn drain(&self) -> Vec<Sample> {
195        std::mem::take(&mut *self.pending.lock().unwrap_or_else(|p| p.into_inner()))
196    }
197}
198
199/// Time one repo call and record it.
200///
201/// The single instrumentation point for both backends. Anything measured
202/// elsewhere would be measuring a different boundary, and the comparison would
203/// be between the timers rather than the implementations.
204pub 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
224/// Write the buffered samples and prune each key back to its window.
225///
226/// Called on a timer and before rendering. Errors are returned rather than
227/// logged here so the caller decides: a metrics write failing must never take
228/// down the request that produced the sample.
229pub 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        // All-time counts, kept apart from the window so PRUNING CANNOT LOSE
249        // THEM. Without this a long-running backend would appear to have served
250        // fewer calls than a freshly-flipped one, which is the opposite of the
251        // truth and exactly the kind of number someone would act on.
252        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    // Prune per (backend, op, ok): successes and failures have their OWN
268    // windows, so a burst of failures cannot evict the successes it should be
269    // compared against.
270    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
288/// The (backend, op, ok) keys touched by a batch, deduplicated.
289fn 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
297/// Read the persisted rows for BOTH backends.
298///
299/// This is what makes the comparison possible at all: a flip is a restart, so
300/// the outgoing backend's numbers exist only here.
301pub 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            // A row written by a version that knew a backend this one does not.
312            // Skipped rather than failing the whole table.
313            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
339/// Render a snapshot as a plain-text table.
340///
341/// Milliseconds with one decimal, because the interesting differences here are
342/// tens of milliseconds (a sidecar hop) and microsecond precision would just be
343/// noise from the scheduler.
344///
345/// Both the success and the failure columns are shown. Reading `ok` without
346/// `err` is how a backend that fails half its calls in 2 ms gets mistaken for a
347/// fast one.
348pub 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
371/// `-` rather than `0.0` for an absent measurement: a dash is obviously "no
372/// data", while a zero reads as the fastest row in the table.
373fn 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    /// Nearest-rank percentiles, checked against a distribution whose answers
396    /// can be read off by hand.
397    #[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        // The lowest percentile must still land on a real sample, not index -1.
404        assert_eq!(stats.percentile(0.0), Some(10));
405    }
406
407    /// **A failed call must not enter the success percentiles.**
408    ///
409    /// An erroring backend fails fast — a refused connection returns far quicker
410    /// than a real PDS round trip. Mixing those in makes the broken path look
411    /// like the faster one on the very dashboard used to decide whether the
412    /// cutover is safe.
413    #[test]
414    fn failures_are_counted_but_kept_out_of_the_success_percentiles() {
415        // Ten slow successes, ninety instant failures.
416        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    /// An operation nobody has exercised reports nothing. Zero would render as
430    /// "instant" and read as the best row in the table.
431    #[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    /// An operation that has only ever failed has no success percentile, but its
439    /// failures are still visible — the row must not vanish.
440    #[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    /// The window keeps the MOST RECENT samples. A window that dropped the
449    /// newest instead would freeze the picture at start-up and never show a
450    /// regression.
451    #[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    /// Both backends land in one table under the same operation name, which is
467    /// what makes the two rows comparable.
468    #[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    /// An absent measurement renders as `-`, never `0.0`. A zero in a latency
493    /// column is indistinguishable from "the fastest thing here" at a glance,
494    /// which is the wrong reading of "never succeeded".
495    #[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        // The failure itself is still visible — the row must not be silently dropped.
510        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    /// An empty table says so rather than rendering a bare header that looks
518    /// like a working instrument reporting nothing wrong.
519    #[test]
520    fn an_empty_snapshot_says_so() {
521        assert!(render(&[]).contains("no repo operations recorded"));
522    }
523
524    /// `timed` records the outcome, not just the duration — the wrapper is the
525    /// only instrumentation point, so if it lost the ok/err distinction the
526    /// separation above would never happen in production.
527    #[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    // ── persistence: the reason the comparison is possible at all ────────────
539
540    async fn pool() -> sqlx::SqlitePool {
541        crate::store::init_url("sqlite::memory:").await.unwrap()
542    }
543
544    /// **The whole point.** A backend flip is a RESTART, so the outgoing
545    /// backend's numbers exist only if they were written down. Two processes
546    /// are simulated by two `RepoMetrics` sharing one database, which is exactly
547    /// what a flip produces.
548    ///
549    /// Without persistence the table can only ever show the backend currently
550    /// running, which is not a comparison.
551    #[tokio::test]
552    async fn a_backend_flip_keeps_the_earlier_backends_rows() {
553        let pool = pool().await;
554
555        // Process 1: the sidecar backend serves some traffic.
556        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        // ... the operator flips the flag and restarts. New process, new table.
563        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    /// Successes and failures are timed apart in the DATABASE too, not just in
589    /// memory — otherwise the property the in-memory tests pin would be lost the
590    /// moment it was written down.
591    #[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    /// **All-time counts must survive pruning.** The window is bounded, but the
615    /// totals are not: a long-running backend that had its oldest samples pruned
616    /// would otherwise appear to have served FEWER calls than one freshly
617    /// flipped to — the opposite of the truth, and the kind of number someone
618    /// would act on.
619    #[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        // The window kept the NEWEST samples, so the smallest surviving value is
641        // the 251st recorded, not the 1st.
642        assert_eq!(rows[0].stats.percentile(0.0), Some(251));
643    }
644
645    /// A burst of failures must not evict the successes it is being compared
646    /// against — the two have separate windows.
647    #[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    /// Flushing twice must not double-count: the buffer is drained, not copied.
666    #[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    /// The buffer is bounded. Without the backstop a dead flusher turns a
679    /// metrics buffer into unbounded memory growth.
680    #[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    /// An unknown backend name is SKIPPED, not guessed at. A row written by a
693    /// newer build must not silently be attributed to a backend this one knows,
694    /// which would put another deployment's numbers in your comparison.
695    #[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    /// The p95 column renders p95, not a second p50. Nothing read that column,
705    /// so pointing it at p50 passed every test.
706    #[test]
707    fn the_rendered_table_distinguishes_p50_from_p95() {
708        let metrics = RepoMetrics::new();
709        // 1ms ninety times, 900ms ten times: p50 is 1ms, p95 is 900ms.
710        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}