diff --git a/Cargo.lock b/Cargo.lock index 03c5ee6..9c690dd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -14,6 +14,15 @@ dependencies = [ "zerocopy", ] +[[package]] +name = "aho-corasick" +version = "1.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ddd31a130427c27518df266943a5308ed92d4b226cc639f5a8f1002816174301" +dependencies = [ + "memchr", +] + [[package]] name = "android_system_properties" version = "0.1.5" @@ -23,6 +32,12 @@ dependencies = [ "libc", ] +[[package]] +name = "anes" +version = "0.1.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "4b46cbb362ab8752921c97e041f5e366ee6297bd428a31275b9fcf1e380f7299" + [[package]] name = "anstream" version = "0.6.21" @@ -103,6 +118,12 @@ version = "1.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1e748733b7cbc798e1434b6ac524f0c1ff2ab456fe201501e6497c8417a4fc33" +[[package]] +name = "cast" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "37b2a672a2cb129a2e41c10b1224bb368f9f37a2b16b612598138befd7b37eb5" + [[package]] name = "cc" version = "1.2.56" @@ -132,6 +153,33 @@ dependencies = [ "windows-link", ] +[[package]] +name = "ciborium" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "42e69ffd6f0917f5c029256a24d0161db17cea3997d185db0d35926308770f0e" +dependencies = [ + "ciborium-io", + "ciborium-ll", + "serde", +] + +[[package]] +name = "ciborium-io" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "05afea1e0a06c9be33d539b876f1ce3692f4afea2cb41f740e7743225ed1c757" + +[[package]] +name = "ciborium-ll" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "57663b653d948a338bfb3eeba9bb2fd5fcfaecb9e199e87e1eda4d9e8b240fd9" +dependencies = [ + "ciborium-io", + "half", +] + [[package]] name = "clap" version = "4.5.60" @@ -184,6 +232,73 @@ version = "0.8.7" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "773648b94d0e5d620f64f280777445740e61fe701025087ec8b57f45c791888b" +[[package]] +name = "criterion" +version = "0.5.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f2b12d017a929603d80db1831cd3a24082f8137ce19c69e6447f54f5fc8d692f" +dependencies = [ + "anes", + "cast", + "ciborium", + "clap", + "criterion-plot", + "is-terminal", + "itertools", + "num-traits", + "once_cell", + "oorandom", + "plotters", + "rayon", + "regex", + "serde", + "serde_derive", + "serde_json", + "tinytemplate", + "walkdir", +] + +[[package]] +name = "criterion-plot" +version = "0.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6b50826342786a51a89e2da3a28f1c32b06e387201bc2d19791f622c673706b1" +dependencies = [ + "cast", + "itertools", +] + +[[package]] +name = "crossbeam-deque" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9dd111b7b7f7d55b72c0a6ae361660ee5853c9af73f70c3c2ef6858b950e2e51" +dependencies = [ + "crossbeam-epoch", + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-epoch" +version = "0.9.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5b82ac4a3c2ca9c3460964f020e1402edd5753411d7737aa39c3714ad1b5420e" +dependencies = [ + "crossbeam-utils", +] + +[[package]] +name = "crossbeam-utils" +version = "0.8.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d0a5c400df2834b80a4c3327b3aad3a4c4cd4de0629063962b03235697506a28" + +[[package]] +name = "crunchy" +version = "0.2.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" + [[package]] name = "directories" version = "5.0.1" @@ -205,6 +320,12 @@ dependencies = [ "windows-sys 0.48.0", ] +[[package]] +name = "either" +version = "1.15.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "48c757948c5ede0e46177b7add2e67155f70e33c07fea8284df6576da70b3719" + [[package]] name = "errno" version = "0.3.14" @@ -244,6 +365,17 @@ dependencies = [ "wasi", ] +[[package]] +name = "half" +version = "2.7.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6ea2d84b969582b4b1864a92dc5d27cd2b77b622a8d79306834f1be5ba20d84b" +dependencies = [ + "cfg-if", + "crunchy", + "zerocopy", +] + [[package]] name = "hashbrown" version = "0.14.5" @@ -268,6 +400,12 @@ version = "0.5.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2304e00983f87ffb38b55b444b5e3b60a884b5d30c0fca7d82fe33449bbe55ea" +[[package]] +name = "hermit-abi" +version = "0.5.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c" + [[package]] name = "iana-time-zone" version = "0.1.65" @@ -292,12 +430,32 @@ dependencies = [ "cc", ] +[[package]] +name = "is-terminal" +version = "0.4.17" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3640c1c38b8e4e43584d8df18be5fc6b0aa314ce6ebf51b53313d4306cca8e46" +dependencies = [ + "hermit-abi", + "libc", + "windows-sys 0.61.2", +] + [[package]] name = "is_terminal_polyfill" version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "a6cb138bb79a146c1bd460005623e142ef0181e3d0219cb493e02f7d08a35695" +[[package]] +name = "itertools" +version = "0.10.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0fd2260e829bddf4cb6ea802289de2f86d6a7a690192fbe91b3f46e0f2c8473" +dependencies = [ + "either", +] + [[package]] name = "itoa" version = "1.0.17" @@ -385,6 +543,7 @@ dependencies = [ "anyhow", "chrono", "clap", + "criterion", "directories", "rusqlite", "serde", @@ -424,6 +583,12 @@ version = "1.70.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "384b8ab6d37215f3c5301a95a4accb5d64aa607f1fcb26a11b5303878451b4fe" +[[package]] +name = "oorandom" +version = "11.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e" + [[package]] name = "option-ext" version = "0.2.0" @@ -465,6 +630,34 @@ version = "0.3.32" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7edddbd0b52d732b21ad9a5fab5c704c14cd949e5e9a1ec5929a24fded1b904c" +[[package]] +name = "plotters" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5aeb6f403d7a4911efb1e33402027fc44f29b5bf6def3effcc22d7bb75f2b747" +dependencies = [ + "num-traits", + "plotters-backend", + "plotters-svg", + "wasm-bindgen", + "web-sys", +] + +[[package]] +name = "plotters-backend" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df42e13c12958a16b3f7f4386b9ab1f3e7933914ecea48da7139435263a4172a" + +[[package]] +name = "plotters-svg" +version = "0.3.7" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "51bae2ac328883f7acdfea3d66a7c35751187f870bc81f94563733a154d7a670" +dependencies = [ + "plotters-backend", +] + [[package]] name = "proc-macro2" version = "1.0.106" @@ -483,6 +676,26 @@ dependencies = [ "proc-macro2", ] +[[package]] +name = "rayon" +version = "1.11.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "368f01d005bf8fd9b1206fb6fa653e6c4a81ceb1466406b81792d87c5677a58f" +dependencies = [ + "either", + "rayon-core", +] + +[[package]] +name = "rayon-core" +version = "1.13.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22e18b0f0062d30d4230b2e85ff77fdfe4326feb054b9783a3460d8435c8ab91" +dependencies = [ + "crossbeam-deque", + "crossbeam-utils", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -503,6 +716,35 @@ dependencies = [ "thiserror", ] +[[package]] +name = "regex" +version = "1.12.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e10754a14b9137dd7b1e3e5b0493cc9171fdd105e0ab477f51b72e7f3ac0e276" +dependencies = [ + "aho-corasick", + "memchr", + "regex-automata", + "regex-syntax", +] + +[[package]] +name = "regex-automata" +version = "0.4.14" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6e1dd4122fc1595e8162618945476892eefca7b88c52820e74af6262213cae8f" +dependencies = [ + "aho-corasick", + "memchr", + "regex-syntax", +] + +[[package]] +name = "regex-syntax" +version = "0.8.10" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc897dd8d9e8bd1ed8cdad82b5966c3e0ecae09fb1907d58efaa013543185d0a" + [[package]] name = "rusqlite" version = "0.31.0" @@ -523,6 +765,15 @@ version = "1.0.22" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d" +[[package]] +name = "same-file" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" +dependencies = [ + "winapi-util", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -659,6 +910,16 @@ dependencies = [ "cfg-if", ] +[[package]] +name = "tinytemplate" +version = "1.2.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "be4d6b5f19ff7664e8c98d03e2139cb510db9b0a60b55f8e8709b689d939b6bc" +dependencies = [ + "serde", + "serde_json", +] + [[package]] name = "tokio" version = "1.50.0" @@ -774,6 +1035,16 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "walkdir" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29790946404f91d9c5d06f9874efddea1dc06c5efe94541a7d6863108e3a5e4b" +dependencies = [ + "same-file", + "winapi-util", +] + [[package]] name = "wasi" version = "0.11.1+wasi-snapshot-preview1" @@ -825,6 +1096,25 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "web-sys" +version = "0.3.91" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "854ba17bb104abfb26ba36da9729addc7ce7f06f5c0f90f3c391f8461cca21f9" +dependencies = [ + "js-sys", + "wasm-bindgen", +] + +[[package]] +name = "winapi-util" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "windows-core" version = "0.62.2" diff --git a/Cargo.toml b/Cargo.toml index 7f08f6f..e82bad1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -3,6 +3,18 @@ name = "nod" version = "0.1.0" edition = "2024" +[lib] +name = "nod" +path = "src/lib.rs" + +[[bin]] +name = "nod" +path = "src/main.rs" + +[[bench]] +name = "queries" +harness = false + [dependencies] tokio = { version = "1", features = ["full"] } tracing = "0.1" @@ -14,3 +26,6 @@ anyhow = "1.0" directories = "5.0" serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" + +[dev-dependencies] +criterion = { version = "0.5", features = ["html_reports"] } diff --git a/benches/queries.rs b/benches/queries.rs new file mode 100644 index 0000000..1d70be5 --- /dev/null +++ b/benches/queries.rs @@ -0,0 +1,109 @@ +use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; +use nod::stats::{collect_stats, collect_trend, BucketSize, SortField}; +use rusqlite::Connection; +use std::sync::Mutex; + +// Seed an in-memory database with n rows spread evenly across one year. +// Row mix: 60% builds (105), 30% substitutions (108), 10% downloads (101). +// Uses a single transaction and a prepared statement for speed. +fn seed(conn: &Connection, n: usize) { + conn.execute_batch(" + CREATE TABLE events ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + nix_id INTEGER, + parent_id INTEGER, + event_type INTEGER, + text TEXT, + drv_path TEXT, + cache_url TEXT, + start_time INTEGER, + end_time INTEGER, + duration_ms INTEGER, + total_bytes INTEGER + ); + CREATE INDEX idx_events_type_start ON events(event_type, start_time); + CREATE INDEX idx_events_start_time ON events(start_time); + PRAGMA journal_mode = WAL; + PRAGMA synchronous = NORMAL; + ").unwrap(); + + let base: i64 = 1_700_000_000; // 2023-11-14 + let span: i64 = 365 * 86400; + + conn.execute_batch("BEGIN").unwrap(); + let mut stmt = conn.prepare( + "INSERT INTO events (event_type, drv_path, cache_url, start_time, end_time, duration_ms, total_bytes) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7)", + ).unwrap(); + + for i in 0..n { + let event_type: i64 = if i % 10 < 6 { 105 } else if i % 10 < 9 { 108 } else { 101 }; + let start = base + (i as i64 * span / n as i64); + // Varied durations 1ms–10min. wrapping_mul avoids overflow; abs() ensures positive. + let duration_ms: i64 = 1 + (i as i64).wrapping_mul(6364136223846793005).abs() % 600_000; + let total_bytes: i64 = if event_type == 101 { (i as i64).wrapping_mul(104729).abs() % 500_000_000 } else { 0 }; + let drv_path: Option = if event_type != 101 { + Some(format!("/nix/store/{:032x}-pkg-{}.drv", i as u128, i % 200)) + } else { + None + }; + let cache_url: Option<&str> = if event_type == 108 { Some("https://cache.nixos.org") } else { None }; + + stmt.execute(rusqlite::params![ + event_type, drv_path, cache_url, + start, start + duration_ms / 1000, + duration_ms, total_bytes, + ]).unwrap(); + } + + conn.execute_batch("COMMIT").unwrap(); +} + +fn bench_collect_stats(c: &mut Criterion) { + let mut group = c.benchmark_group("collect_stats"); + + for n in [10_000usize, 100_000, 1_000_000] { + let conn = Connection::open_in_memory().unwrap(); + seed(&conn, n); + let db = Mutex::new(conn); + + group.bench_with_input(BenchmarkId::new("no_filter", n), &n, |b, _| { + b.iter(|| collect_stats(&db, None, None, SortField::Duration, 10, false).unwrap()) + }); + + group.bench_with_input(BenchmarkId::new("grouped", n), &n, |b, _| { + b.iter(|| collect_stats(&db, None, None, SortField::Count, 10, true).unwrap()) + }); + } + + group.finish(); +} + +fn bench_collect_trend(c: &mut Criterion) { + let mut group = c.benchmark_group("collect_trend"); + + for n in [10_000usize, 100_000, 1_000_000] { + let conn = Connection::open_in_memory().unwrap(); + seed(&conn, n); + let db = Mutex::new(conn); + + // aggregate (raw=false): SQL window function median — the default code path. + group.bench_with_input(BenchmarkId::new("aggregate/month", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Month, None, false).unwrap()) + }); + + group.bench_with_input(BenchmarkId::new("aggregate/day", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Day, None, false).unwrap()) + }); + + // raw (raw=true): sends all durations for Mann-Whitney — the memory-heavy path. + group.bench_with_input(BenchmarkId::new("raw/month", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Month, None, true).unwrap()) + }); + } + + group.finish(); +} + +criterion_group!(benches, bench_collect_stats, bench_collect_trend); +criterion_main!(benches); diff --git a/src/daemon.rs b/src/daemon.rs index 6394438..a50545e 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -35,11 +35,13 @@ CREATE TABLE IF NOT EXISTS events ( total_bytes INTEGER ); CREATE INDEX IF NOT EXISTS idx_events_type_start ON events(event_type, start_time); +CREATE INDEX IF NOT EXISTS idx_events_start_time ON events(start_time); "; const SCHEMA_HASHES: &[u32] = &[ 0x9bc94a70, // v1: TEXT timestamps, no indexes 0xee061d32, // v2: INTEGER timestamps (Unix seconds), idx_events_type_start + 0x7c09711e, // v3: idx_events_start_time ]; const SCHEMA_VERSION: u32 = SCHEMA_HASHES.len() as u32; const _: () = assert!( @@ -47,7 +49,8 @@ const _: () = assert!( "schema changed - append new hash to SCHEMA_HASHES and add a migration in MIGRATIONS" ); -// (target_version, sql). Table rebuild because SQLite does not support ALTER COLUMN. +// (target_version, sql). Must be ordered by target version ascending. +// Table rebuild used for v2 because SQLite does not support ALTER COLUMN. const MIGRATIONS: &[(u32, &str)] = &[ (2, " CREATE TABLE events_new ( @@ -73,6 +76,7 @@ const MIGRATIONS: &[(u32, &str)] = &[ ALTER TABLE events_new RENAME TO events; CREATE INDEX IF NOT EXISTS idx_events_type_start ON events(event_type, start_time); "), + (3, "CREATE INDEX IF NOT EXISTS idx_events_start_time ON events(start_time);"), ]; #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] @@ -202,6 +206,8 @@ enum ClientCommand { since: Option, bucket: BucketSize, drv: Option, + #[serde(default)] + raw: bool, }, Clean { // Unix timestamp; None means delete everything. @@ -285,12 +291,12 @@ pub fn open_db(path: &PathBuf) -> Result { } conn.execute_batch(" - PRAGMA journal_mode = WAL; - PRAGMA synchronous = NORMAL; - PRAGMA temp_store = MEMORY; - PRAGMA mmap_size = 134217728; - PRAGMA cache_size = -8000; - PRAGMA wal_autocheckpoint = 0; + PRAGMA journal_mode = WAL; + PRAGMA synchronous = NORMAL; + PRAGMA temp_store = MEMORY; + PRAGMA mmap_size = 134217728; + PRAGMA cache_size = -8000; + PRAGMA wal_autocheckpoint = 1000; ").context("Failed to configure database")?; Ok(conn) @@ -301,8 +307,12 @@ fn run_retention(conn: &Connection, retain_days: u32) -> Result<()> { let cutoff = Utc::now().timestamp() - retain_days as i64 * 86400; let deleted = conn.execute("DELETE FROM events WHERE start_time < ?1", [cutoff]) .context("Retention DELETE failed")?; - conn.execute_batch("PRAGMA wal_checkpoint(PASSIVE)") - .context("Retention checkpoint failed")?; + // TRUNCATE resets and shrinks the WAL file; use RESTART as fallback if readers + // are active (TRUNCATE fails when a reader holds the WAL open). + if conn.execute_batch("PRAGMA wal_checkpoint(TRUNCATE)").is_err() { + conn.execute_batch("PRAGMA wal_checkpoint(RESTART)") + .context("Retention checkpoint failed")?; + } if deleted > 0 { info!(deleted, retain_days, "Retention cleanup removed rows"); } @@ -380,9 +390,9 @@ async fn handle_connection(mut stream: UnixStream, state: Arc>, db: writer.write_all((serde_json::to_string(&stats)? + "\n").as_bytes()).await?; break; } - Ok(SocketMessage::Command(ClientCommand::GetTrend { since, bucket, drv })) => { + Ok(SocketMessage::Command(ClientCommand::GetTrend { since, bucket, drv, raw })) => { let db = Arc::clone(&db); - let trend = tokio::task::spawn_blocking(move || collect_trend(&db.reader, since, bucket, drv)) + let trend = tokio::task::spawn_blocking(move || collect_trend(&db.reader, since, bucket, drv, raw)) .await??; writer.write_all((serde_json::to_string(&trend)? + "\n").as_bytes()).await?; break; diff --git a/src/lib.rs b/src/lib.rs new file mode 100644 index 0000000..c04b866 --- /dev/null +++ b/src/lib.rs @@ -0,0 +1,2 @@ +pub mod daemon; +pub mod stats; diff --git a/src/main.rs b/src/main.rs index dba4b6b..569bdf4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,6 +1,3 @@ -mod daemon; -mod stats; - use anyhow::{Context, Result}; use chrono::Utc; use clap::{Parser, Subcommand}; @@ -14,8 +11,8 @@ use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::UnixStream; use tracing::{error, info}; -use daemon::{open_db, run_daemon, DbConnections}; -use stats::{ +use nod::daemon::{open_db, run_daemon, DbConnections}; +use nod::stats::{ collect_trend, display_stats, display_trend, display_trend_test, output_csv_trend, BucketSize, SortField, Stats, Trend, }; @@ -276,11 +273,14 @@ async fn query_via_socket( assert!(limit > 0); if let Some(bucket_size) = bucket { + // raw=true only when Mann-Whitney test output is requested — sends all durations. + let raw = matches!(output, OutputFormat::Test); let cmd = serde_json::json!({ "action": "get_trend", "since": since, "bucket": bucket_size, "drv": drv, + "raw": raw, }); stream.write_all((cmd.to_string() + "\n").as_bytes()).await?; @@ -326,10 +326,11 @@ fn query_direct( let conn = Mutex::new(conn); if let Some(bucket_size) = bucket { - let trend = collect_trend(&conn, since, bucket_size, drv)?; + let raw = matches!(output, OutputFormat::Test); + let trend = collect_trend(&conn, since, bucket_size, drv, raw)?; display_trend_output(&trend, output); } else { - let s = stats::collect_stats(&conn, since, drv.as_deref(), sort, limit, group)?; + let s = nod::stats::collect_stats(&conn, since, drv.as_deref(), sort, limit, group)?; display_stats(s); } diff --git a/src/stats.rs b/src/stats.rs index 04e12df..a2b39fe 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -227,7 +227,7 @@ fn erf_approx(x: f64) -> f64 { if x >= 0.0 { result } else { -result } } -fn median_sorted(sorted: &[i64]) -> f64 { +pub fn median_sorted(sorted: &[i64]) -> f64 { assert!(!sorted.is_empty()); let n = sorted.len(); if n % 2 == 0 { @@ -237,7 +237,7 @@ fn median_sorted(sorted: &[i64]) -> f64 { } } -fn fmt_ms(ms: i64) -> String { +pub fn fmt_ms(ms: i64) -> String { if ms < 1000 { format!("{}ms", ms) } else if ms < 60_000 { @@ -261,7 +261,7 @@ pub enum BucketSize { } impl BucketSize { - fn strftime_fmt(&self) -> &'static str { + pub fn strftime_fmt(&self) -> &'static str { match self { BucketSize::Hour => "%Y-%m-%dT%H", BucketSize::Day => "%Y-%m-%d", @@ -280,14 +280,25 @@ impl BucketSize { } } -// Raw durations are carried per-bucket so adjacent buckets can be compared -// without a second round-trip to the daemon. #[derive(Debug, Serialize, Deserialize)] pub struct TrendBucket { pub bucket: String, + // Precomputed aggregates — always populated; computed client-side in raw mode, + // computed SQL-side in aggregate mode. + #[serde(default)] + pub build_count: i64, + #[serde(default)] + pub build_median_ms: i64, + #[serde(default)] + pub subst_count: i64, + #[serde(default)] + pub subst_median_ms: i64, + pub download_bytes: i64, + // Raw durations — only populated when raw=true (needed for Mann-Whitney). + #[serde(default)] pub build_durations: Vec, + #[serde(default)] pub subst_durations: Vec, - pub download_bytes: i64, } #[derive(Debug, Serialize, Deserialize)] @@ -297,21 +308,47 @@ pub struct Trend { pub drv_filter: Option, } +// raw=false: compute per-bucket median in SQL via window functions — no raw durations +// sent over the socket. Use this for all display modes except Mann-Whitney. +// raw=true: return all durations for Mann-Whitney (--output test). Memory-intensive +// for large datasets; only use when genuinely needed. pub fn collect_trend( db: &Mutex, since: Option, bucket: BucketSize, drv: Option, + raw: bool, ) -> Result { if let Some(ref d) = drv { assert!(!d.is_empty(), "drv filter must not be empty"); } let conn = db.lock().unwrap(); - let drv_ref = drv.as_deref(); let fmt = bucket.strftime_fmt(); + let buckets = if raw { + collect_trend_raw(&conn, since, drv_ref, fmt)? + } else { + collect_trend_aggregate(&conn, since, drv_ref, fmt)? + }; + + for i in 1..buckets.len() { + assert!(buckets[i].bucket > buckets[i - 1].bucket, "buckets must be strictly ascending"); + } + + Ok(Trend { buckets, bucket_size: bucket, drv_filter: drv }) +} + +// Fetches all raw durations. ORDER BY start_time (not strftime) so the index +// (event_type, start_time) can be used for ordering — strftime is monotone in +// start_time so bucket grouping is preserved. +fn collect_trend_raw( + conn: &Connection, + since: Option, + drv: Option<&str>, + fmt: &str, +) -> Result> { // FileTransfer (101) has NULL drv_path and is intentionally excluded by the drv filter. let mut stmt = conn.prepare( "SELECT strftime(?3, start_time, 'unixepoch'), event_type, duration_ms, total_bytes @@ -319,16 +356,20 @@ pub fn collect_trend( WHERE event_type IN (101, 105, 108) AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%') - ORDER BY strftime(?3, start_time, 'unixepoch') ASC, start_time ASC", - ).context("Failed to prepare trend query")?; + ORDER BY start_time ASC", + ).context("Failed to prepare raw trend query")?; let mut buckets: Vec = vec![]; - for row in stmt.query_map(rusqlite::params![since, drv_ref, fmt], |r| { + for row in stmt.query_map(rusqlite::params![since, drv, fmt], |r| { Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?)) })?.filter_map(|r| r.ok()) { let (b, etype, dur, bytes) = row; if buckets.last().map(|x: &TrendBucket| x.bucket.as_str()) != Some(&b) { - buckets.push(TrendBucket { bucket: b, build_durations: vec![], subst_durations: vec![], download_bytes: 0 }); + buckets.push(TrendBucket { + bucket: b, build_count: 0, build_median_ms: 0, + subst_count: 0, subst_median_ms: 0, download_bytes: 0, + build_durations: vec![], subst_durations: vec![], + }); } let last = buckets.last_mut().unwrap(); match etype { @@ -339,11 +380,110 @@ pub fn collect_trend( } } - for i in 1..buckets.len() { - assert!(buckets[i].bucket > buckets[i - 1].bucket, "buckets must be strictly ascending"); + // Compute aggregates from raw durations so display functions can use them uniformly. + for b in &mut buckets { + if !b.build_durations.is_empty() { + let mut s = b.build_durations.clone(); s.sort_unstable(); + b.build_count = s.len() as i64; + b.build_median_ms = median_sorted(&s) as i64; + } + if !b.subst_durations.is_empty() { + let mut s = b.subst_durations.clone(); s.sort_unstable(); + b.subst_count = s.len() as i64; + b.subst_median_ms = median_sorted(&s) as i64; + } } - Ok(Trend { buckets, bucket_size: bucket, drv_filter: drv }) + Ok(buckets) +} + +// Computes per-bucket median entirely in SQL using window functions. Returns no +// raw durations — only counts and medians. Memory usage is proportional to the +// number of distinct buckets, not the number of rows. +fn collect_trend_aggregate( + conn: &Connection, + since: Option, + drv: Option<&str>, + fmt: &str, +) -> Result> { + // Window function median: ROW_NUMBER orders rows by duration within each + // (bucket, event_type) partition; we take the one or two middle rows and AVG them. + // Integer division: odd n → single middle row; even n → average of two middle rows. + let mut stmt = conn.prepare( + "SELECT bucket, event_type, + CAST(AVG(duration_ms) AS INTEGER) AS median_ms, + MAX(n) AS cnt + FROM ( + SELECT strftime(?3, start_time, 'unixepoch') AS bucket, + event_type, duration_ms, + ROW_NUMBER() OVER ( + PARTITION BY strftime(?3, start_time, 'unixepoch'), event_type + ORDER BY duration_ms + ) AS rn, + COUNT(*) OVER ( + PARTITION BY strftime(?3, start_time, 'unixepoch'), event_type + ) AS n + FROM events + WHERE event_type IN (105, 108) + AND (?1 IS NULL OR start_time >= ?1) + AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%') + ) + WHERE rn IN ((n + 1) / 2, (n + 2) / 2) + GROUP BY bucket, event_type + ORDER BY bucket ASC", + ).context("Failed to prepare aggregate trend query")?; + + let mut buckets: Vec = vec![]; + for row in stmt.query_map(rusqlite::params![since, drv, fmt], |r| { + Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?)) + })?.filter_map(|r| r.ok()) { + let (b, etype, median_ms, cnt) = row; + assert!(median_ms >= 0); + assert!(cnt > 0); + if buckets.last().map(|x: &TrendBucket| x.bucket.as_str()) != Some(&b) { + buckets.push(TrendBucket { + bucket: b, build_count: 0, build_median_ms: 0, + subst_count: 0, subst_median_ms: 0, download_bytes: 0, + build_durations: vec![], subst_durations: vec![], + }); + } + let last = buckets.last_mut().unwrap(); + match etype { + 105 => { last.build_median_ms = median_ms; last.build_count = cnt; } + 108 => { last.subst_median_ms = median_ms; last.subst_count = cnt; } + _ => {} + } + } + + // Separate query for download bytes — FileTransfer has no meaningful duration median. + let mut dl_stmt = conn.prepare( + "SELECT strftime(?2, start_time, 'unixepoch') AS bucket, SUM(total_bytes) + FROM events + WHERE event_type = 101 + AND (?1 IS NULL OR start_time >= ?1) + GROUP BY bucket + ORDER BY bucket ASC", + ).context("Failed to prepare download bytes query")?; + + for row in dl_stmt.query_map(rusqlite::params![since, fmt], |r| { + Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?)) + })?.filter_map(|r| r.ok()) { + let (b, bytes) = row; + assert!(bytes >= 0); + // Merge into existing bucket or create a download-only bucket. + if let Some(existing) = buckets.iter_mut().find(|x| x.bucket == b) { + existing.download_bytes = bytes; + } else { + buckets.push(TrendBucket { + bucket: b, build_count: 0, build_median_ms: 0, + subst_count: 0, subst_median_ms: 0, download_bytes: bytes, + build_durations: vec![], subst_durations: vec![], + }); + } + } + + buckets.sort_unstable_by(|a, b| a.bucket.cmp(&b.bucket)); + Ok(buckets) } pub fn display_trend(trend: &Trend) { @@ -365,14 +505,9 @@ pub fn display_trend(trend: &Trend) { } for b in &trend.buckets { - let mut bs = b.build_durations.clone(); bs.sort_unstable(); - let mut ss = b.subst_durations.clone(); ss.sort_unstable(); - let build_med = if bs.is_empty() { 0 } else { median_sorted(&bs) as i64 }; - let subst_med = if ss.is_empty() { 0 } else { median_sorted(&ss) as i64 }; - print!("{:6} {:>10} {:>6} {:>10}", - b.bucket, b.build_durations.len(), fmt_ms(build_med), - b.subst_durations.len(), fmt_ms(subst_med)); + b.bucket, b.build_count, fmt_ms(b.build_median_ms), + b.subst_count, fmt_ms(b.subst_median_ms)); if has_downloads { print!(" {:>8.1}", b.download_bytes as f64 / 1_048_576.0); } println!(); } @@ -438,13 +573,9 @@ pub fn display_trend_test(trend: &Trend) { pub fn output_csv_trend(trend: &Trend) { println!("period,build_count,build_median_ms,subst_count,subst_median_ms,download_bytes"); for b in &trend.buckets { - let mut bs = b.build_durations.clone(); bs.sort_unstable(); - let mut ss = b.subst_durations.clone(); ss.sort_unstable(); - let build_med = if bs.is_empty() { 0 } else { median_sorted(&bs) as i64 }; - let subst_med = if ss.is_empty() { 0 } else { median_sorted(&ss) as i64 }; - assert!(build_med >= 0); - assert!(subst_med >= 0); - println!("{},{},{},{},{},{}", b.bucket, b.build_durations.len(), build_med, b.subst_durations.len(), subst_med, b.download_bytes); + assert!(b.build_median_ms >= 0); + assert!(b.subst_median_ms >= 0); + println!("{},{},{},{},{},{}", b.bucket, b.build_count, b.build_median_ms, b.subst_count, b.subst_median_ms, b.download_bytes); } }