use std::collections::HashMap;
use std::sync::Mutex;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum Backend {
Sidecar,
Rust,
}
impl Backend {
pub fn parse(raw: &str) -> Option<Self> {
match raw {
"sidecar" => Some(Backend::Sidecar),
"rust" => Some(Backend::Rust),
_ => None,
}
}
pub fn as_str(self) -> &'static str {
match self {
Backend::Sidecar => "sidecar",
Backend::Rust => "rust",
}
}
}
const WINDOW: usize = 1024;
#[derive(Debug, Default, Clone)]
pub struct Stats {
ok_micros: Vec<u64>,
err_micros: Vec<u64>,
pub ok_count: u64,
pub err_count: u64,
}
impl Stats {
fn record(&mut self, micros: u64, ok: bool) {
let (samples, count) = if ok {
(&mut self.ok_micros, &mut self.ok_count)
} else {
(&mut self.err_micros, &mut self.err_count)
};
*count += 1;
if samples.len() == WINDOW {
samples.remove(0);
}
samples.push(micros);
}
pub fn percentile(&self, p: f64) -> Option<u64> {
percentile_of(&self.ok_micros, p)
}
pub fn error_percentile(&self, p: f64) -> Option<u64> {
percentile_of(&self.err_micros, p)
}
}
fn percentile_of(samples: &[u64], p: f64) -> Option<u64> {
if samples.is_empty() {
return None;
}
let mut sorted = samples.to_vec();
sorted.sort_unstable();
let rank = ((p / 100.0) * sorted.len() as f64).ceil().max(1.0) as usize;
Some(sorted[rank.min(sorted.len()) - 1])
}
#[derive(Debug, Clone)]
pub struct Row {
pub backend: Backend,
pub op: String,
pub stats: Stats,
}
const MAX_PENDING: usize = 16_384;
#[derive(Debug, Default)]
pub struct RepoMetrics {
stats: Mutex<HashMap<(Backend, &'static str), Stats>>,
pending: Mutex<Vec<Sample>>,
}
#[derive(Debug, Clone, Copy)]
struct Sample {
backend: Backend,
op: &'static str,
micros: u64,
ok: bool,
}
impl RepoMetrics {
pub fn new() -> Self {
Self::default()
}
pub fn record(&self, backend: Backend, op: &'static str, micros: u64, ok: bool) {
self.stats
.lock()
.unwrap_or_else(|p| p.into_inner())
.entry((backend, op))
.or_default()
.record(micros, ok);
let mut pending = self.pending.lock().unwrap_or_else(|p| p.into_inner());
if pending.len() >= MAX_PENDING {
pending.remove(0);
}
pending.push(Sample {
backend,
op,
micros,
ok,
});
}
pub fn snapshot(&self) -> Vec<Row> {
let stats = self.stats.lock().unwrap_or_else(|p| p.into_inner());
let mut rows: Vec<Row> = stats
.iter()
.map(|((backend, op), stats)| Row {
backend: *backend,
op: (*op).to_string(),
stats: stats.clone(),
})
.collect();
rows.sort_by(|a, b| (a.backend, &a.op).cmp(&(b.backend, &b.op)));
rows
}
fn drain(&self) -> Vec<Sample> {
std::mem::take(&mut *self.pending.lock().unwrap_or_else(|p| p.into_inner()))
}
}
pub async fn timed<T, E, F>(
metrics: &RepoMetrics,
backend: Backend,
op: &'static str,
call: F,
) -> Result<T, E>
where
F: std::future::Future<Output = Result<T, E>>,
{
let started = std::time::Instant::now();
let result = call.await;
metrics.record(
backend,
op,
started.elapsed().as_micros() as u64,
result.is_ok(),
);
result
}
pub async fn flush(metrics: &RepoMetrics, pool: &sqlx::SqlitePool, now: i64) -> anyhow::Result<()> {
let samples = metrics.drain();
if samples.is_empty() {
return Ok(());
}
let mut tx = pool.begin().await?;
for sample in &samples {
sqlx::query(
"INSERT INTO repo_timing (backend, op, micros, ok, at) VALUES (?1, ?2, ?3, ?4, ?5)",
)
.bind(sample.backend.as_str())
.bind(sample.op)
.bind(sample.micros as i64)
.bind(i64::from(sample.ok))
.bind(now)
.execute(&mut *tx)
.await?;
sqlx::query(
"INSERT INTO repo_timing_total (backend, op, ok_count, err_count) \
VALUES (?1, ?2, ?3, ?4) \
ON CONFLICT(backend, op) DO UPDATE SET \
ok_count = ok_count + excluded.ok_count, \
err_count = err_count + excluded.err_count",
)
.bind(sample.backend.as_str())
.bind(sample.op)
.bind(i64::from(sample.ok))
.bind(i64::from(!sample.ok))
.execute(&mut *tx)
.await?;
}
for (backend, op, ok) in distinct_keys(&samples) {
sqlx::query(
"DELETE FROM repo_timing WHERE backend = ?1 AND op = ?2 AND ok = ?3 AND id NOT IN \
(SELECT id FROM repo_timing WHERE backend = ?1 AND op = ?2 AND ok = ?3 \
ORDER BY id DESC LIMIT ?4)",
)
.bind(backend.as_str())
.bind(op)
.bind(i64::from(ok))
.bind(WINDOW as i64)
.execute(&mut *tx)
.await?;
}
tx.commit().await?;
Ok(())
}
fn distinct_keys(samples: &[Sample]) -> Vec<(Backend, &'static str, bool)> {
let mut keys: Vec<(Backend, &'static str, bool)> =
samples.iter().map(|s| (s.backend, s.op, s.ok)).collect();
keys.sort();
keys.dedup();
keys
}
pub async fn persisted_rows(pool: &sqlx::SqlitePool) -> anyhow::Result<Vec<Row>> {
let totals: Vec<(String, String, i64, i64)> = sqlx::query_as(
"SELECT backend, op, ok_count, err_count FROM repo_timing_total ORDER BY backend, op",
)
.fetch_all(pool)
.await?;
let mut rows = Vec::with_capacity(totals.len());
for (backend, op, ok_count, err_count) in totals {
let Some(backend) = Backend::parse(&backend) else {
continue;
};
let samples: Vec<(i64, i64)> =
sqlx::query_as("SELECT micros, ok FROM repo_timing WHERE backend = ?1 AND op = ?2")
.bind(backend.as_str())
.bind(&op)
.fetch_all(pool)
.await?;
let mut stats = Stats {
ok_count: ok_count as u64,
err_count: err_count as u64,
..Stats::default()
};
for (micros, ok) in samples {
if ok == 1 {
stats.ok_micros.push(micros as u64);
} else {
stats.err_micros.push(micros as u64);
}
}
rows.push(Row { backend, op, stats });
}
Ok(rows)
}
pub fn render(rows: &[Row]) -> String {
let mut out = String::from(
"backend operation ok p50ms p95ms err errp50ms\n",
);
if rows.is_empty() {
out.push_str("(no repo operations recorded yet)\n");
return out;
}
for row in rows {
out.push_str(&format!(
"{:<8} {:<28} {:>4} {:>7} {:>7} {:>6} {:>9}\n",
row.backend.as_str(),
row.op,
row.stats.ok_count,
render_micros(row.stats.percentile(50.0)),
render_micros(row.stats.percentile(95.0)),
row.stats.err_count,
render_micros(row.stats.error_percentile(50.0)),
));
}
out
}
fn render_micros(micros: Option<u64>) -> String {
match micros {
Some(micros) => format!("{:.1}", micros as f64 / 1000.0),
None => "-".to_string(),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn stats_of(ok: &[u64], err: &[u64]) -> Stats {
let mut stats = Stats::default();
for micros in ok {
stats.record(*micros, true);
}
for micros in err {
stats.record(*micros, false);
}
stats
}
#[test]
fn percentiles_come_from_the_recorded_samples() {
let stats = stats_of(&[10, 20, 30, 40, 50, 60, 70, 80, 90, 100], &[]);
assert_eq!(stats.percentile(50.0), Some(50));
assert_eq!(stats.percentile(95.0), Some(100));
assert_eq!(stats.percentile(100.0), Some(100));
assert_eq!(stats.percentile(0.0), Some(10));
}
#[test]
fn failures_are_counted_but_kept_out_of_the_success_percentiles() {
let stats = stats_of(&[1000; 10], &[1; 90]);
assert_eq!(
stats.percentile(50.0),
Some(1000),
"fast failures dragged the success percentile down, which is how a \
broken backend passes for a fast one"
);
assert_eq!(stats.ok_count, 10);
assert_eq!(stats.err_count, 90);
assert_eq!(stats.error_percentile(50.0), Some(1));
}
#[test]
fn an_unexercised_operation_reports_nothing_rather_than_zero() {
let stats = Stats::default();
assert_eq!(stats.percentile(50.0), None);
assert_eq!(stats.error_percentile(50.0), None);
}
#[test]
fn an_operation_that_only_ever_fails_still_reports_its_failures() {
let stats = stats_of(&[], &[5, 7, 9]);
assert_eq!(stats.percentile(50.0), None);
assert_eq!(stats.err_count, 3);
assert_eq!(stats.error_percentile(50.0), Some(7));
}
#[test]
fn the_window_keeps_the_most_recent_samples() {
let mut stats = Stats::default();
for i in 0..(WINDOW as u64 + 10) {
stats.record(i, true);
}
assert_eq!(stats.ok_count, WINDOW as u64 + 10, "the total counts all");
assert_eq!(
stats.percentile(0.0),
Some(10),
"the oldest samples should have aged out of the window"
);
assert_eq!(stats.percentile(100.0), Some(WINDOW as u64 + 9));
}
#[test]
fn the_two_backends_are_recorded_side_by_side() {
let metrics = RepoMetrics::new();
metrics.record(Backend::Sidecar, "list_subscriptions", 5_000, true);
metrics.record(Backend::Rust, "list_subscriptions", 2_000, true);
let snapshot = metrics.snapshot();
assert_eq!(snapshot.len(), 2);
let sidecar = snapshot
.iter()
.find(|r| r.backend == Backend::Sidecar)
.expect("the sidecar row");
let rust = snapshot
.iter()
.find(|r| r.backend == Backend::Rust)
.expect("the rust row");
assert_eq!(
sidecar.op, rust.op,
"the same operation name, or the rows cannot be compared"
);
assert_eq!(sidecar.stats.percentile(50.0), Some(5_000));
assert_eq!(rust.stats.percentile(50.0), Some(2_000));
}
#[test]
fn an_absent_measurement_renders_as_a_dash_rather_than_zero() {
let metrics = RepoMetrics::new();
metrics.record(Backend::Rust, "list_folders", 3_000, false);
let table = render(&metrics.snapshot());
assert!(
!table.contains("0.0"),
"a never-succeeded operation rendered as 0.0ms, which reads as instant:\n{table}"
);
assert!(
table.contains('-'),
"expected a dash for the absent p50:\n{table}"
);
assert!(table.contains("list_folders"), "the row vanished:\n{table}");
assert!(
table.contains("3.0"),
"the failure latency is missing:\n{table}"
);
}
#[test]
fn an_empty_snapshot_says_so() {
assert!(render(&[]).contains("no repo operations recorded"));
}
#[tokio::test]
async fn timed_records_success_and_failure_distinctly() {
let metrics = RepoMetrics::new();
let _: Result<(), ()> = timed(&metrics, Backend::Rust, "op", async { Ok(()) }).await;
let _: Result<(), ()> = timed(&metrics, Backend::Rust, "op", async { Err(()) }).await;
let snapshot = metrics.snapshot();
assert_eq!(snapshot[0].stats.ok_count, 1);
assert_eq!(snapshot[0].stats.err_count, 1);
}
async fn pool() -> sqlx::SqlitePool {
crate::store::init_url("sqlite::memory:").await.unwrap()
}
#[tokio::test]
async fn a_backend_flip_keeps_the_earlier_backends_rows() {
let pool = pool().await;
let before = RepoMetrics::new();
for _ in 0..5 {
before.record(Backend::Sidecar, "list_subscriptions_sorted", 9_000, true);
}
flush(&before, &pool, 1_700_000_000).await.unwrap();
let after = RepoMetrics::new();
for _ in 0..5 {
after.record(Backend::Rust, "list_subscriptions_sorted", 3_000, true);
}
flush(&after, &pool, 1_700_000_100).await.unwrap();
assert_eq!(
after.snapshot().len(),
1,
"in-process memory only ever holds the running backend"
);
let rows = persisted_rows(&pool).await.unwrap();
assert_eq!(
rows.len(),
2,
"both backends must survive the flip: {rows:?}"
);
let sidecar = rows.iter().find(|r| r.backend == Backend::Sidecar).unwrap();
let rust = rows.iter().find(|r| r.backend == Backend::Rust).unwrap();
assert_eq!(sidecar.stats.percentile(50.0), Some(9_000));
assert_eq!(rust.stats.percentile(50.0), Some(3_000));
assert_eq!(sidecar.op, rust.op, "rows must be comparable by operation");
}
#[tokio::test]
async fn persisted_failures_stay_out_of_the_success_percentiles() {
let pool = pool().await;
let metrics = RepoMetrics::new();
for _ in 0..10 {
metrics.record(Backend::Rust, "add_subscription", 1_000, true);
}
for _ in 0..90 {
metrics.record(Backend::Rust, "add_subscription", 1, false);
}
flush(&metrics, &pool, 1_700_000_000).await.unwrap();
let rows = persisted_rows(&pool).await.unwrap();
let stats = &rows[0].stats;
assert_eq!(
stats.percentile(50.0),
Some(1_000),
"fast failures dragged the persisted success percentile down"
);
assert_eq!(stats.ok_count, 10);
assert_eq!(stats.err_count, 90);
}
#[tokio::test]
async fn pruning_bounds_the_window_without_losing_the_totals() {
let pool = pool().await;
let metrics = RepoMetrics::new();
let total = WINDOW + 250;
for i in 0..total {
metrics.record(Backend::Rust, "list_folders_sorted", i as u64 + 1, true);
}
flush(&metrics, &pool, 1_700_000_000).await.unwrap();
let kept: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM repo_timing")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(kept, WINDOW as i64, "the window is not bounded");
let rows = persisted_rows(&pool).await.unwrap();
assert_eq!(
rows[0].stats.ok_count, total as u64,
"pruning ate the all-time count"
);
assert_eq!(rows[0].stats.percentile(0.0), Some(251));
}
#[tokio::test]
async fn a_burst_of_failures_does_not_evict_the_successes() {
let pool = pool().await;
let metrics = RepoMetrics::new();
metrics.record(Backend::Rust, "put_read_state", 5_000, true);
for _ in 0..(WINDOW + 100) {
metrics.record(Backend::Rust, "put_read_state", 2, false);
}
flush(&metrics, &pool, 1_700_000_000).await.unwrap();
let rows = persisted_rows(&pool).await.unwrap();
assert_eq!(
rows[0].stats.percentile(50.0),
Some(5_000),
"the only success was evicted by a flood of failures"
);
}
#[tokio::test]
async fn flushing_twice_does_not_double_count() {
let pool = pool().await;
let metrics = RepoMetrics::new();
metrics.record(Backend::Rust, "remove_saved", 1_000, true);
flush(&metrics, &pool, 1_700_000_000).await.unwrap();
flush(&metrics, &pool, 1_700_000_001).await.unwrap();
let rows = persisted_rows(&pool).await.unwrap();
assert_eq!(rows[0].stats.ok_count, 1, "the sample was counted twice");
}
#[test]
fn the_pending_buffer_is_bounded() {
let metrics = RepoMetrics::new();
for _ in 0..(MAX_PENDING + 100) {
metrics.record(Backend::Rust, "op", 1, true);
}
assert!(
metrics.pending.lock().unwrap().len() <= MAX_PENDING,
"the write buffer grew past its bound"
);
}
#[test]
fn an_unknown_backend_name_is_not_guessed() {
assert_eq!(Backend::parse("sidecar"), Some(Backend::Sidecar));
assert_eq!(Backend::parse("rust"), Some(Backend::Rust));
for unknown in ["", "RUST", "postgres", "rust "] {
assert_eq!(Backend::parse(unknown), None, "guessed at {unknown:?}");
}
}
#[test]
fn the_rendered_table_distinguishes_p50_from_p95() {
let metrics = RepoMetrics::new();
for _ in 0..90 {
metrics.record(Backend::Rust, "op", 1_000, true);
}
for _ in 0..10 {
metrics.record(Backend::Rust, "op", 900_000, true);
}
let table = render(&metrics.snapshot());
assert!(table.contains("1.0"), "p50 missing from:\n{table}");
assert!(
table.contains("900.0"),
"the p95 column is not showing p95:\n{table}"
);
}
}