diff --git a/src/daemon.rs b/src/daemon.rs index dfb6dc9..d45c247 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -10,7 +10,7 @@ use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader}; use tokio::net::{UnixListener, UnixStream}; use tracing::{error, info}; -use crate::stats::collect_stats; +use crate::stats::{collect_stats, collect_trend, BucketSize}; const fn schema_hash(s: &[u8]) -> u32 { let mut h: u32 = 2166136261; @@ -158,6 +158,7 @@ enum NixEvent { #[serde(tag = "action", rename_all = "snake_case")] enum ClientCommand { GetStats { since: Option }, + GetTrend { since: Option, bucket: BucketSize, drv: Option }, Clean, } @@ -260,6 +261,13 @@ async fn handle_connection(mut stream: UnixStream, state: Arc>) -> writer.write_all((serde_json::to_string(&stats)? + "\n").as_bytes()).await?; break; } + Ok(SocketMessage::Command(ClientCommand::GetTrend { since, bucket, drv })) => { + let db = state.lock().unwrap().db.clone(); + let trend = tokio::task::spawn_blocking(move || collect_trend(&db, since, bucket, drv)) + .await??; + writer.write_all((serde_json::to_string(&trend)? + "\n").as_bytes()).await?; + break; + } Ok(SocketMessage::Command(ClientCommand::Clean)) => { let db = state.lock().unwrap().db.clone(); tokio::task::spawn_blocking(move || -> Result<()> { diff --git a/src/main.rs b/src/main.rs index 44a02f4..8b501fb 100644 --- a/src/main.rs +++ b/src/main.rs @@ -12,7 +12,15 @@ use tokio::net::UnixStream; use tracing::{error, info}; use daemon::{open_db, run_daemon}; -use stats::{display_stats, Stats}; +use stats::{display_stats, display_trend, display_trend_test, output_csv_trend, BucketSize, Stats, Trend}; + +#[derive(Clone, clap::ValueEnum)] +enum BucketArg { + Hour, + Day, + Week, + Month, +} #[derive(Parser)] struct Cli { @@ -46,6 +54,33 @@ enum Commands { #[arg(short = 'y', long)] years: Option, }, + /// Show how builds, substitutions, and downloads trend over time + Trend { + /// Path to the Unix socket + #[arg(short, long, env = "NOD_SOCKET")] + socket: Option, + /// Limit to last N days + #[arg(short = 'd', long)] + days: Option, + /// Limit to last N months + #[arg(short = 'm', long)] + months: Option, + /// Limit to last N years + #[arg(short = 'y', long)] + years: Option, + /// Time bucket granularity (auto-detected from period if omitted) + #[arg(short = 'b', long, value_enum)] + bucket: Option, + /// Filter builds and substitutions to derivations matching this substring + #[arg(long)] + drv: Option, + /// Run Mann-Whitney U test between adjacent periods instead of showing the plain table + #[arg(long)] + test: bool, + /// Output as CSV + #[arg(long)] + csv: bool, + }, /// Clear all data from the database Clean { /// Path to the Unix socket @@ -61,7 +96,7 @@ async fn main() -> Result<()> { let project_dirs = ProjectDirs::from("org", "nixos", "nod"); - let command = cli.command.unwrap_or(Commands::Daemon { db: None, socket: None }); + let command = cli.command.unwrap_or(Commands::Stats { socket: None, days: None, months: None, years: None }); match command { Commands::Daemon { db, socket } => { @@ -92,6 +127,12 @@ async fn main() -> Result<()> { } } + // Check if a daemon is already listening before we remove and rebind the socket. + // A successful connect means another process owns it — refuse to start. + if UnixStream::connect(&socket_path).await.is_ok() { + anyhow::bail!("Daemon already running at {}", socket_path.display()); + } + let conn = open_db(&db_path)?; run_daemon(socket_path, Arc::new(Mutex::new(conn))).await.context("Daemon failed")? } @@ -126,6 +167,55 @@ async fn main() -> Result<()> { let stats: Stats = serde_json::from_str(&line).context("Invalid response from daemon")?; display_stats(stats); } + Commands::Trend { socket, days, months, years, bucket, drv, test, csv } => { + let socket_path = socket.unwrap_or_else(|| { + project_dirs.as_ref() + .and_then(|d| d.runtime_dir()) + .map(|d| d.join("nod.sock")) + .unwrap_or_else(|| PathBuf::from("/tmp/nod.sock")) + }); + + let since: Option = if days.is_some() || months.is_some() || years.is_some() { + let mut t = Utc::now(); + if let Some(y) = years { t = t - chrono::Months::new(y * 12); } + if let Some(m) = months { t = t - chrono::Months::new(m); } + if let Some(d) = days { t = t - chrono::Duration::days(d as i64); } + Some(t.timestamp()) + } else { + None + }; + + // Auto-detect bucket granularity from the period length when not specified. + let bucket_size = match bucket { + Some(BucketArg::Hour) => BucketSize::Hour, + Some(BucketArg::Day) => BucketSize::Day, + Some(BucketArg::Week) => BucketSize::Week, + Some(BucketArg::Month) => BucketSize::Month, + None => match since { + None => BucketSize::Month, + Some(ts) => { + let days_span = (Utc::now().timestamp() - ts) / 86400; + if days_span <= 2 { BucketSize::Hour } + else if days_span <= 60 { BucketSize::Day } + else if days_span <= 365 { BucketSize::Week } + else { BucketSize::Month } + } + }, + }; + + let mut stream = UnixStream::connect(&socket_path) + .await + .with_context(|| format!("Failed to connect to daemon at {}", socket_path.display()))?; + + let cmd = serde_json::json!({"action": "get_trend", "since": since, "bucket": bucket_size, "drv": drv}); + stream.write_all((cmd.to_string() + "\n").as_bytes()).await?; + + let mut reader = BufReader::new(stream); + let mut line = String::new(); + reader.read_line(&mut line).await.context("Daemon closed connection without response")?; + let trend: Trend = serde_json::from_str(&line).context("Invalid response from daemon")?; + if csv { output_csv_trend(&trend); } else if test { display_trend_test(&trend); } else { display_trend(&trend); } + } Commands::Clean { socket } => { let socket_path = socket.unwrap_or_else(|| { project_dirs.as_ref() diff --git a/src/stats.rs b/src/stats.rs index 3fe058c..9cb5b04 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -3,7 +3,6 @@ use rusqlite::Connection; use serde::{Deserialize, Serialize}; use std::sync::{Arc, Mutex}; -// Sent over the socket as JSON - all computed values, no raw types #[derive(Debug, Serialize, Deserialize)] pub struct Stats { pub build_count: i64, @@ -33,8 +32,7 @@ pub struct CacheStat { pub fn collect_stats(db: &Arc>, since: Option) -> Result { let conn = db.lock().unwrap(); - // Convert unix timestamp to RFC3339 for comparison with stored start_time strings. - // None means no filter - SQL NULL makes the condition vacuously true. + // SQL NULL makes the WHERE condition vacuously true, giving us "no filter". let since_str: Option = since .and_then(|ts| chrono::DateTime::from_timestamp(ts, 0)) .map(|dt| dt.to_rfc3339()); @@ -68,8 +66,8 @@ pub fn collect_stats(db: &Arc>, since: Option) -> Result< }) })?.filter_map(|r| r.ok()).collect(); - // Use Substitute (108) events for cache latency — these measure full download time - // per substituter, not just metadata query time (QueryPathInfo). + // Substitute (108) measures full substitution time per cache; QueryPathInfo (109) + // only measures metadata lookup and would undercount latency. let mut stmt = conn.prepare( "SELECT cache_url, AVG(duration_ms), COUNT(*) FROM events WHERE event_type = 108 AND cache_url IS NOT NULL AND (?1 IS NULL OR start_time >= ?1) @@ -86,6 +84,104 @@ pub fn collect_stats(db: &Arc>, since: Option) -> Result< Ok(Stats { build_count, build_total_ms, subst_count, subst_total_ms, download_bytes, download_ms, slowest_builds, cache_latency }) } +// Mann-Whitney U is non-parametric and makes no distributional assumptions, +// which is appropriate for build times that are right-skewed. +pub struct MannWhitneyResult { + pub p_value: f64, +} + +pub fn mann_whitney_u(a: &[i64], b: &[i64]) -> Option { + if a.is_empty() || b.is_empty() { return None; } + + let n1 = a.len(); + let n2 = b.len(); + let n1f = n1 as f64; + let n2f = n2 as f64; + + let mut combined: Vec<(i64, usize)> = a.iter().map(|&v| (v, 0)) + .chain(b.iter().map(|&v| (v, 1))) + .collect(); + combined.sort_unstable_by_key(|&(v, _)| v); + + let n = combined.len(); + + // Average ranks within tie groups and accumulate the tie-correction term Σ(t³ - t). + let mut rank_sum_a = 0.0f64; + let mut tie_correction = 0.0f64; + let mut i = 0; + while i < n { + let mut j = i; + while j < n && combined[j].0 == combined[i].0 { j += 1; } + let avg_rank = (i as f64 + 1.0 + j as f64) / 2.0; + for k in i..j { + if combined[k].1 == 0 { rank_sum_a += avg_rank; } + } + let t = (j - i) as f64; + if t > 1.0 { tie_correction += t * t * t - t; } + i = j; + } + + assert!(rank_sum_a >= 0.0); + + let u_a = rank_sum_a - n1f * (n1f + 1.0) / 2.0; + let u_b = n1f * n2f - u_a; + + assert!(u_a >= 0.0); + assert!(u_b >= 0.0); + assert!((u_a + u_b - n1f * n2f).abs() < 1e-6, "U_A + U_B must equal n1*n2"); + + let cliffs_delta = (u_a - u_b) / (n1f * n2f); + assert!(cliffs_delta >= -1.0 - 1e-9); + assert!(cliffs_delta <= 1.0 + 1e-9); + let _ = cliffs_delta; + + // Var[U] = n1*n2/12 * [(n+1) - Σ(t³-t)/(n*(n-1))] + let nf = n as f64; + let variance = (n1f * n2f / 12.0) * ((nf + 1.0) - tie_correction / (nf * (nf - 1.0))); + + let p_value = if variance <= 0.0 { + 1.0 + } else { + let u_min = u_a.min(u_b); + let mean_u = n1f * n2f / 2.0; + // Continuity correction: +0.5 shifts U_min toward the mean, making z conservative. + let z = (u_min - mean_u + 0.5) / variance.sqrt(); + assert!(z <= 0.0 + 1e-9, "z must be non-positive for U_min ≤ mean_U"); + (2.0 * normal_cdf(z)).min(1.0) + }; + + assert!(p_value >= 0.0); + assert!(p_value <= 1.0); + + Some(MannWhitneyResult { p_value }) +} + +// Abramowitz & Stegun 7.1.26, max error ≈ 1.5×10⁻⁷. +fn normal_cdf(z: f64) -> f64 { + 0.5 * (1.0 + erf_approx(z / std::f64::consts::SQRT_2)) +} + +fn erf_approx(x: f64) -> f64 { + let t = 1.0 / (1.0 + 0.3275911 * x.abs()); + let poly = t * (0.254829592 + + t * (-0.284496736 + + t * (1.421413741 + + t * (-1.453152027 + + t * 1.061405429)))); + let result = 1.0 - poly * (-x * x).exp(); + if x >= 0.0 { result } else { -result } +} + +fn median_sorted(sorted: &[i64]) -> f64 { + assert!(!sorted.is_empty()); + let n = sorted.len(); + if n % 2 == 0 { + (sorted[n / 2 - 1] + sorted[n / 2]) as f64 / 2.0 + } else { + sorted[n / 2] as f64 + } +} + fn fmt_ms(ms: i64) -> String { if ms < 1000 { format!("{}ms", ms) @@ -96,11 +192,212 @@ fn fmt_ms(ms: i64) -> String { } } -// Strip /nix/store/ prefix, keep the hash and package name. fn drv_name(path: &str) -> &str { path.strip_prefix("/nix/store/").unwrap_or(path) } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum BucketSize { + Hour, + Day, + Week, + Month, +} + +impl BucketSize { + fn strftime_fmt(&self) -> &'static str { + match self { + BucketSize::Hour => "%Y-%m-%dT%H", + BucketSize::Day => "%Y-%m-%d", + BucketSize::Week => "%Y-W%W", + BucketSize::Month => "%Y-%m", + } + } + + fn col_width(&self) -> usize { + match self { + BucketSize::Hour => 13, + BucketSize::Day => 10, + BucketSize::Week => 8, + BucketSize::Month => 7, + } + } +} + +// 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, + pub build_durations: Vec, + pub subst_durations: Vec, + pub download_bytes: i64, +} + +#[derive(Debug, Serialize, Deserialize)] +pub struct Trend { + pub buckets: Vec, + pub bucket_size: BucketSize, + pub drv_filter: Option, +} + +pub fn collect_trend( + db: &Arc>, + since: Option, + bucket: BucketSize, + drv: Option, +) -> Result { + if let Some(ref d) = drv { + assert!(!d.is_empty(), "drv filter must not be empty"); + } + + let conn = db.lock().unwrap(); + + let since_str: Option = since + .and_then(|ts| chrono::DateTime::from_timestamp(ts, 0)) + .map(|dt| dt.to_rfc3339()); + let p = since_str.as_deref(); + let drv_ref = drv.as_deref(); + let fmt = bucket.strftime_fmt(); + + // Rows come out ordered by bucket then time, so grouping by sequential scan is valid. + // FileTransfer (101) has NULL drv_path and is intentionally excluded by the drv filter. + let mut stmt = conn.prepare( + "SELECT strftime(?3, start_time), event_type, duration_ms, total_bytes + FROM events + 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) ASC, start_time ASC", + ).context("Failed to prepare trend query")?; + + let mut buckets: Vec = vec![]; + for row in stmt.query_map(rusqlite::params![p, drv_ref, 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 }); + } + let last = buckets.last_mut().unwrap(); + match etype { + 105 => { assert!(dur >= 0); last.build_durations.push(dur); } + 108 => { assert!(dur >= 0); last.subst_durations.push(dur); } + 101 => { assert!(bytes >= 0); last.download_bytes += bytes; } + _ => {} + } + } + + 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 }) +} + +pub fn display_trend(trend: &Trend) { + let has_downloads = trend.drv_filter.is_none() + && trend.buckets.iter().any(|b| b.download_bytes > 0); + let bw = trend.bucket_size.col_width(); + + if let Some(ref drv) = trend.drv_filter { + println!("filter: {}", drv); + } + + print!("{:6} {:>10} {:>6} {:>10}", "period", "builds", "build med", "subst", "subst med"); + if has_downloads { print!(" {:>8}", "dl (MB)"); } + println!(); + + if trend.buckets.is_empty() { + println!("(no data)"); + return; + } + + 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)); + if has_downloads { print!(" {:>8.1}", b.download_bytes as f64 / 1_048_576.0); } + println!(); + } +} + +fn print_test_section(label: &str, buckets: &[TrendBucket], get_durs: F, bw: usize) +where + F: Fn(&TrendBucket) -> &[i64], +{ + let any_data = buckets.iter().any(|b| !get_durs(b).is_empty()); + if !any_data { return; } + + println!("{}", label); + println!("{:5} {:>10} {:>7} {:>8}", "period", "n", "median", "Δ", "p-value"); + + for (i, b) in buckets.iter().enumerate() { + let durs = get_durs(b); + let mut sorted = durs.to_vec(); sorted.sort_unstable(); + let med = if sorted.is_empty() { 0.0 } else { median_sorted(&sorted) }; + + let (delta_str, p_str) = if i == 0 || get_durs(&buckets[i - 1]).is_empty() { + (String::new(), String::new()) + } else { + let prev = get_durs(&buckets[i - 1]); + let mut prev_sorted = prev.to_vec(); prev_sorted.sort_unstable(); + let prev_med = median_sorted(&prev_sorted); + + let delta = if prev_med > 0.0 { + let pct = (med - prev_med) / prev_med * 100.0; + let sign = if pct >= 0.0 { "+" } else { "" }; + format!("{}{:.0}%", sign, pct) + } else { + String::new() + }; + + let p = match mann_whitney_u(prev, durs) { + Some(r) => format!("{:.3}", r.p_value), + None => String::new(), + }; + + (delta, p) + }; + + println!("{:5} {:>10} {:>7} {:>8}", + b.bucket, durs.len(), fmt_ms(med as i64), delta_str, p_str); + } + println!(); +} + +pub fn display_trend_test(trend: &Trend) { + let bw = trend.bucket_size.col_width(); + + if let Some(ref drv) = trend.drv_filter { + println!("filter: {}", drv); + } + + print_test_section("builds", &trend.buckets, |b| &b.build_durations, bw); + print_test_section("substitutions", &trend.buckets, |b| &b.subst_durations, bw); + + println!("Mann-Whitney U (two-tailed). H0: adjacent periods have the same duration distribution."); +} + +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); + } +} + pub fn display_stats(stats: Stats) { let build_avg = if stats.build_count > 0 { stats.build_total_ms / stats.build_count } else { 0 }; let subst_avg = if stats.subst_count > 0 { stats.subst_total_ms / stats.subst_count } else { 0 };