diff --git a/Cargo.lock b/Cargo.lock index 9c690dd..7fa4c99 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -486,9 +486,9 @@ checksum = "b5b646652bf6661599e1da8901b3b9522896f01e736bad5f723fe7a3a27f899d" [[package]] name = "libredox" -version = "0.1.14" +version = "0.1.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1744e39d1d6a9948f4f388969627434e31128196de472883b39f148769bfe30a" +checksum = "e02f3bb43d335493c96bf3fd3a321600bf6bd07ed34bc64118e9293bdffea46c" dependencies = [ "libc", ] diff --git a/Cargo.toml b/Cargo.toml index e82bad1..5c2510d 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -23,9 +23,9 @@ rusqlite = { version = "0.31", features = ["bundled"] } chrono = "0.4" clap = { version = "4", features = ["derive", "env"] } anyhow = "1.0" -directories = "5.0" serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" +directories = "5.0" [dev-dependencies] criterion = { version = "0.5", features = ["html_reports"] } diff --git a/README.md b/README.md index ef34ad4..22ef811 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@ ## nod -A simple daemon that collects Nix build and substitution statistics using structured JSON logs. +A daemon that collects Nix build and substitution statistics using structured JSON logs. ## requirements @@ -23,30 +23,51 @@ nod daemon Point Nix at the socket in `nix.conf`: ``` -json-log-path = /run/user/1000/nod/nod.sock +json-log-path = /tmp/nod.sock ``` -Then use Nix normally. View accumulated stats with: +Then use Nix normally. Query accumulated stats: ```bash -nod stats # all time -nod stats -d 7 # last 7 days -nod stats -m 3 # last 3 months -nod stats -y 1 # last year -nod clean # reset the database +nod # aggregate stats, all time +nod -d 30 # last 30 days +nod -m 3 # last 3 months +nod --drv firefox # filter to derivations matching "firefox" +nod --sort count --group # group by derivation, sort by frequency +nod --bucket day # time series by day +nod --bucket month # time series by month +nod --bucket day --output test # Mann-Whitney significance test vs prior period +nod --bucket day --output csv # CSV output +nod clean # delete all events +nod clean --before-days 90 # delete events older than 90 days ``` +## NixOS + +```nix +{ + services.nod = { + enable = true; + retainDays = 180; # default + }; +} +``` + +The module sets `nix.settings.json-log-path` automatically and exports `NOD_SOCKET` into `/etc/environment` so every session (login, SSH, scripts) finds the daemon without `--socket`. + ## configuration | flag | env | default | |------|-----|---------| -| `--socket` | `NOD_SOCKET` | `$XDG_RUNTIME_DIR/nod/nod.sock` | -| `--db` | `NOD_DB` | `$XDG_DATA_HOME/nod/nod.db` | +| `--socket` | `NOD_SOCKET` | `/tmp/nod.sock` | +| `--db` | `NOD_DB` | `nod.db` (current directory) | -Both directories are created automatically. The socket path in `nix.conf` must match `--socket`. +When using the NixOS module both are pinned to fixed system paths and exported automatically. ## how? -Nix 2.30 added `json-log-path`, which writes a stream of structured activity events (start/result/stop) to a file or Unix socket while a build runs. nod listens on that socket, tracks in-flight activities by ID, and on each stop event inserts a completed row into a local SQLite database. `nod stats` queries that database through the daemon. +Nix 2.30 added `json-log-path`, which writes a stream of structured activity events (start/result/stop) to a file or Unix socket while a build runs. nod listens on that socket, tracks in-flight activities by ID, and on each stop event inserts a completed row into a local SQLite database. `nod` queries that database through the daemon socket. + +Stats queries are served from a pre-aggregated `daily_stats` table (one row per day per event type) maintained in lockstep with inserts, so summary queries are O(days) regardless of total event count. The current day's aggregates are held in memory and flushed on day rollover. Relevant Nix source: [logging.hh](https://github.com/NixOS/nix/blob/b4de973847370204cf28fe2092abdd21f25ee0e8/src/libutil/include/nix/util/logging.hh) diff --git a/benches/queries.rs b/benches/queries.rs index b4769c1..56d8a9a 100644 --- a/benches/queries.rs +++ b/benches/queries.rs @@ -120,25 +120,40 @@ fn bench_collect_stats(c: &mut Criterion) { fn bench_collect_trend(c: &mut Criterion) { let mut group = c.benchmark_group("collect_trend"); + let today_day = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH).unwrap().as_secs() as i64 / 86400; + let today_snap = TodaySummary { day: today_day, ..Default::default() }; + 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); - // All-time monthly: exercises full scan + Rust-side bucketing. - group.bench_with_input(BenchmarkId::new("all_time/month", n), &n, |b, _| { - b.iter(|| collect_trend(&db, None, BucketSize::Month, None).unwrap()) + // full=false: daily_stats path (table/csv output). + let t = today_snap.clone(); + group.bench_with_input(BenchmarkId::new("agg/all_time/month", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Month, None, Some(t.clone()), false).unwrap()) + }); + + let t = today_snap.clone(); + group.bench_with_input(BenchmarkId::new("agg/all_time/day", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Day, None, Some(t.clone()), false).unwrap()) + }); + + let since = Some(BASE_TIME + YEAR_SPAN * 9 / 10); + let t = today_snap.clone(); + group.bench_with_input(BenchmarkId::new("agg/recent_10pct/day", n), &n, |b, _| { + b.iter(|| collect_trend(&db, since, BucketSize::Day, None, Some(t.clone()), false).unwrap()) }); - // All-time daily: more bucket transitions, otherwise identical scan. - group.bench_with_input(BenchmarkId::new("all_time/day", n), &n, |b, _| { - b.iter(|| collect_trend(&db, None, BucketSize::Day, None).unwrap()) + // full=true: events scan (--output test). + group.bench_with_input(BenchmarkId::new("full/all_time/month", n), &n, |b, _| { + b.iter(|| collect_trend(&db, None, BucketSize::Month, None, None, true).unwrap()) }); - // Recent 10%: start_time index reduces rows scanned by 90%. let since = Some(BASE_TIME + YEAR_SPAN * 9 / 10); - group.bench_with_input(BenchmarkId::new("recent_10pct/day", n), &n, |b, _| { - b.iter(|| collect_trend(&db, since, BucketSize::Day, None).unwrap()) + group.bench_with_input(BenchmarkId::new("full/recent_10pct/day", n), &n, |b, _| { + b.iter(|| collect_trend(&db, since, BucketSize::Day, None, None, true).unwrap()) }); } diff --git a/flake.nix b/flake.nix index 69f8f2a..d5d3f5e 100644 --- a/flake.nix +++ b/flake.nix @@ -84,12 +84,17 @@ socketPath = lib.mkOption { type = lib.types.path; default = "/run/nod/nod.sock"; - description = "Path to the Unix socket. Exposed via NOD_SOCKET in the session environment."; + description = "Path to the Unix socket. Propagated to all sessions via /etc/environment so nod always finds the daemon without --socket."; }; databasePath = lib.mkOption { type = lib.types.path; default = "/var/lib/nod/nod.db"; - description = "Path to the SQLite database"; + description = "Path to the SQLite database."; + }; + retainDays = lib.mkOption { + type = lib.types.nullOr lib.types.ints.positive; + default = null; + description = "Override the retention period in days. When null the daemon default of 180 days is used."; }; }; @@ -101,22 +106,32 @@ }; users.groups.${cfg.group} = {}; - # Tell nix to forward its internal JSON log to the socket + # Forward Nix's internal JSON activity log to the daemon socket. + # The nix-daemon runs as root so the socket directory must be world-searchable + # and the socket itself must be group-writable (handled by RuntimeDirectoryMode + # and UMask below). Users that only need to query nod require no group membership. nix.settings.json-log-path = cfg.socketPath; - # Make the socket path available to interactive shells so users - # can run `nod stats` without passing --socket explicitly. - environment.sessionVariables.NOD_SOCKET = cfg.socketPath; + # Expose the socket path to every session (login, SSH, scripts) via /etc/environment + # so that `nod` always resolves the socket without needing --socket or NOD_SOCKET set + # manually. sessionVariables only reaches interactive login shells and would cause + # "cannot connect to socket" errors in non-login SSH sessions and cron jobs. + environment.variables.NOD_SOCKET = cfg.socketPath; + environment.variables.NOD_DB = cfg.databasePath; + + # Make `nod` available to all users without manual systemPackages entries. + environment.systemPackages = [ cfg.package ]; systemd.services.nod = { description = "Nix Observability Daemon"; wantedBy = ["multi-user.target"]; - after = ["network.target"]; + after = ["local-fs.target"]; serviceConfig = { User = cfg.user; Group = cfg.group; - ExecStart = "${cfg.package}/bin/nod daemon --db ${cfg.databasePath} --socket ${cfg.socketPath}"; + ExecStart = "${cfg.package}/bin/nod daemon --db ${cfg.databasePath} --socket ${cfg.socketPath}" + + lib.optionalString (cfg.retainDays != null) " --retain-days ${toString cfg.retainDays}"; Restart = "always"; StateDirectory = "nod"; StateDirectoryMode = "0750"; diff --git a/src/daemon.rs b/src/daemon.rs index c224225..fc3052a 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -269,6 +269,10 @@ enum ClientCommand { since: Option, bucket: BucketSize, drv: Option, + // true = return individual duration samples (needed for --output test / Mann-Whitney). + // false = return aggregates only via the daily_stats fast path. + #[serde(default)] + full: bool, }, Clean { // Unix timestamp; None means delete everything. @@ -404,9 +408,14 @@ pub fn open_db(path: &PathBuf) -> Result { fn run_retention(conn: &Connection, retain_days: u32) -> Result<()> { assert!(retain_days > 0, "retain_days must be > 0"); - let cutoff = Utc::now().timestamp() - retain_days as i64 * 86400; + let cutoff = Utc::now().timestamp() - retain_days as i64 * 86400; + let cutoff_day = cutoff / 86400; let deleted = conn.execute("DELETE FROM events WHERE start_time < ?1", [cutoff]) .context("Retention DELETE failed")?; + conn.execute("DELETE FROM daily_stats WHERE day < ?1", [cutoff_day]) + .context("Retention DELETE daily_stats failed")?; + conn.execute("DELETE FROM daily_cache_stats WHERE day < ?1", [cutoff_day]) + .context("Retention DELETE daily_cache_stats 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() { @@ -559,8 +568,9 @@ fn rebuild_today(conn: &Connection, today_day: i64) -> Result { pub async fn run_daemon( socket_path: PathBuf, db: Arc, - retain_days: Option, + retain_days: u32, ) -> Result<()> { + assert!(retain_days > 0, "retain_days must be > 0"); let today_day = Utc::now().timestamp() / 86400; let today = { let conn = db.reader.lock().unwrap(); @@ -578,11 +588,10 @@ pub async fn run_daemon( } // Run retention on startup, then every 6 hours. - if let Some(days) = retain_days { - assert!(days > 0, "retain_days must be > 0"); + { let db_startup = Arc::clone(&db); tokio::task::spawn_blocking(move || { - run_retention(&db_startup.writer.lock().unwrap(), days) + run_retention(&db_startup.writer.lock().unwrap(), retain_days) }).await??; let db_timer = Arc::clone(&db); @@ -593,7 +602,7 @@ pub async fn run_daemon( interval.tick().await; let db2 = Arc::clone(&db_timer); if let Err(e) = tokio::task::spawn_blocking(move || { - run_retention(&db2.writer.lock().unwrap(), days) + run_retention(&db2.writer.lock().unwrap(), retain_days) }).await { error!("Retention timer task failed: {}", e); } @@ -635,10 +644,16 @@ 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, full })) => { + let today = if !full && drv.is_none() && !matches!(bucket, BucketSize::Hour) { + Some(state.lock().unwrap().today.snapshot()) + } else { + None + }; let db = Arc::clone(&db); - let trend = tokio::task::spawn_blocking(move || collect_trend(&db.reader, since, bucket, drv)) - .await??; + let trend = tokio::task::spawn_blocking(move || { + collect_trend(&db.reader, since, bucket, drv, today, full) + }).await??; writer.write_all((serde_json::to_string(&trend)? + "\n").as_bytes()).await?; break; } diff --git a/src/main.rs b/src/main.rs index 0d54c73..69006c4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -79,9 +79,9 @@ struct Cli { enum Commands { /// Run the observability daemon Daemon { - /// Delete events older than N days; cleanup runs on startup and every 6 hours - #[arg(long, env = "NOD_RETAIN_DAYS")] - retain_days: Option, + /// Delete events older than N days; cleanup runs on startup and every 6 hours (default: 180) + #[arg(long, env = "NOD_RETAIN_DAYS", default_value = "180")] + retain_days: u32, }, /// Clear data from the database Clean { @@ -91,21 +91,16 @@ enum Commands { }, } -fn resolve_db_path(db: Option, project_dirs: &Option) -> PathBuf { +fn resolve_db_path(db: Option) -> PathBuf { db.unwrap_or_else(|| { - project_dirs.as_ref() - .map(|d| d.data_dir().join("nod.db")) + ProjectDirs::from("", "", "nod") + .map(|dirs| dirs.data_local_dir().join("nod.db")) .unwrap_or_else(|| PathBuf::from("nod.db")) }) } -fn resolve_socket_path(socket: Option, project_dirs: &Option) -> PathBuf { - 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")) - }) +fn resolve_socket_path(socket: Option) -> PathBuf { + socket.unwrap_or_else(|| PathBuf::from("/tmp/nod.sock")) } fn is_connection_error(e: &std::io::Error) -> bool { @@ -153,16 +148,12 @@ async fn main() -> Result<()> { tracing_subscriber::fmt::init(); let cli = Cli::parse(); - let project_dirs = ProjectDirs::from("org", "nixos", "nod"); - match cli.command { Some(Commands::Daemon { retain_days }) => { - if let Some(days) = retain_days { - assert!(days > 0, "--retain-days must be > 0"); - } + assert!(retain_days > 0, "--retain-days must be > 0"); - let db_path = resolve_db_path(cli.db, &project_dirs); - let socket_path = resolve_socket_path(cli.socket, &project_dirs); + let db_path = resolve_db_path(cli.db); + let socket_path = resolve_socket_path(cli.socket); if let Some(parent) = db_path.parent() { if !parent.as_os_str().is_empty() { @@ -193,7 +184,7 @@ async fn main() -> Result<()> { } Some(Commands::Clean { before_days }) => { - let socket_path = resolve_socket_path(cli.socket, &project_dirs); + let socket_path = resolve_socket_path(cli.socket); let before: Option = before_days.map(|d| { assert!(d > 0, "--before-days must be > 0"); (Utc::now() - chrono::Duration::days(d as i64)).timestamp() @@ -216,7 +207,7 @@ async fn main() -> Result<()> { } Err(e) if is_connection_error(&e) => { // Daemon not running — operate directly on the database. - let db_path = resolve_db_path(cli.db, &project_dirs); + let db_path = resolve_db_path(cli.db); let conn = open_db(&db_path)?; if let Some(ts) = before { let deleted = conn.execute( @@ -242,14 +233,14 @@ async fn main() -> Result<()> { } let since = compute_since(cli.days, cli.months, cli.years); - let socket_path = resolve_socket_path(cli.socket, &project_dirs); + let socket_path = resolve_socket_path(cli.socket); match UnixStream::connect(&socket_path).await { Ok(stream) => { query_via_socket(stream, since, cli.drv, cli.bucket, cli.sort, cli.limit, cli.group, cli.output).await?; } Err(e) if is_connection_error(&e) => { - let db_path = resolve_db_path(cli.db, &project_dirs); + let db_path = resolve_db_path(cli.db); query_direct(&db_path, since, cli.drv, cli.bucket, cli.sort, cli.limit, cli.group, &cli.output)?; } Err(e) => return Err(e).context("Failed to connect to daemon"), @@ -278,6 +269,7 @@ async fn query_via_socket( "since": since, "bucket": bucket_size, "drv": drv, + "full": matches!(output, OutputFormat::Test), }); stream.write_all((cmd.to_string() + "\n").as_bytes()).await?; @@ -323,7 +315,8 @@ fn query_direct( let conn = Mutex::new(conn); if let Some(bucket_size) = bucket { - let trend = collect_trend(&conn, since, bucket_size, drv)?; + let full = matches!(output, OutputFormat::Test); + let trend = collect_trend(&conn, since, bucket_size, drv, None, full)?; display_trend_output(&trend, output); } else { let s = nod::stats::collect_stats(&conn, since, drv.as_deref(), sort, limit, group, None)?; diff --git a/src/stats.rs b/src/stats.rs index 27a14f1..110c9fb 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -441,9 +441,16 @@ fn bucket_end(ts: i64, size: &BucketSize) -> i64 { #[derive(Debug, Serialize, Deserialize)] pub struct TrendBucket { pub bucket: String, + pub build_count: i64, + pub build_total_ms: i64, + pub subst_count: i64, + pub subst_total_ms: i64, + pub download_bytes: i64, + // Only populated when full duration data is requested (--output test). + #[serde(default)] pub build_durations: Vec, + #[serde(default)] pub subst_durations: Vec, - pub download_bytes: i64, } #[derive(Debug, Serialize, Deserialize)] @@ -453,23 +460,19 @@ pub struct Trend { pub drv_filter: Option, } -// Query returns raw start_time integers — strftime is computed in Rust via -// bucket_label()/bucket_end() to avoid N SQLite string allocations. Bucket -// boundaries are detected with a single integer comparison per row (ts >= next_bucket) -// so string allocs happen only once per bucket, not once per row. -pub fn collect_trend( - db: &Mutex, +// Events-table scan. Returns buckets with individual duration samples populated. +// Required when individual samples are needed (--output test / Mann-Whitney) or +// when a drv filter or hour granularity rules out the daily_stats path. +// +// Raw start_time integers from SQLite — bucket boundaries detected with a single +// integer comparison per row (ts >= next_bucket) so string allocs happen only +// once per bucket, not once per row. +fn trend_from_events( + conn: &Connection, 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 drv_ref = drv.as_deref(); - + bucket: &BucketSize, + drv: Option<&str>, +) -> Result> { // FileTransfer (101) has NULL drv_path and is intentionally excluded by the drv filter. // INDEXED BY forces the covering start_time-first index so all four projected columns // (start_time, event_type, duration_ms, total_bytes) are served without table lookups. @@ -485,7 +488,7 @@ pub fn collect_trend( let mut buckets: Vec = vec![]; let mut next_bucket: i64 = 0; - for row in stmt.query_map(rusqlite::params![since, drv_ref], |r| { + for row in stmt.query_map(rusqlite::params![since, drv], |r| { Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?)) })?.filter_map(|r| r.ok()) { let (ts, etype, dur, bytes) = row; @@ -493,23 +496,134 @@ pub fn collect_trend( if ts >= next_bucket { buckets.push(TrendBucket { - bucket: bucket_label(ts, &bucket), + bucket: bucket_label(ts, bucket), + build_count: 0, build_total_ms: 0, + subst_count: 0, subst_total_ms: 0, + download_bytes: 0, build_durations: vec![], subst_durations: vec![], - download_bytes: 0, }); - next_bucket = bucket_end(ts, &bucket); + next_bucket = bucket_end(ts, bucket); } 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); } + 105 => { assert!(dur >= 0); last.build_count += 1; last.build_total_ms += dur; last.build_durations.push(dur); } + 108 => { assert!(dur >= 0); last.subst_count += 1; last.subst_total_ms += dur; last.subst_durations.push(dur); } 101 => { assert!(bytes >= 0); last.download_bytes += bytes; } _ => {} } } + Ok(buckets) +} + +// daily_stats query. O(days) — does not populate build_durations/subst_durations. +// Closed days come from the table; today's partial data is merged from the in-memory +// snapshot so the current day is always included when running via the daemon. +fn trend_from_daily_stats( + conn: &Connection, + since: Option, + bucket: &BucketSize, + today: Option, +) -> Result> { + let since_day: Option = since.map(|ts| ts / 86400); + let today_day = today.as_ref().map(|t| t.day) + .unwrap_or_else(|| Utc::now().timestamp() / 86400); + + let mut stmt = conn.prepare( + "SELECT day * 86400, event_type, count, total_ms, total_bytes + FROM daily_stats + WHERE event_type IN (101, 105, 108) + AND (?1 IS NULL OR day >= ?1) + AND day < ?2 + ORDER BY day ASC", + ).context("Failed to prepare trend query")?; + + let mut buckets: Vec = vec![]; + let mut next_bucket: i64 = 0; + + for row in stmt.query_map(rusqlite::params![since_day, today_day], |r| { + Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?, r.get::<_, i64>(4)?)) + })?.filter_map(|r| r.ok()) { + let (ts, etype, count, total_ms, total_bytes) = row; + assert!(ts >= 0); + assert!(count >= 0); + + if ts >= next_bucket { + buckets.push(TrendBucket { + bucket: bucket_label(ts, bucket), + build_count: 0, build_total_ms: 0, + subst_count: 0, subst_total_ms: 0, + download_bytes: 0, + build_durations: vec![], + subst_durations: vec![], + }); + next_bucket = bucket_end(ts, bucket); + } + + let last = buckets.last_mut().unwrap(); + match etype { + 105 => { last.build_count += count; last.build_total_ms += total_ms; } + 108 => { last.subst_count += count; last.subst_total_ms += total_ms; } + 101 => { last.download_bytes += total_bytes; } + _ => {} + } + } + + if let Some(t) = today { + if t.build_count > 0 || t.subst_count > 0 || t.download_bytes > 0 { + let today_ts = t.day * 86400; + if since.map_or(true, |s| today_ts >= s) { + if today_ts >= next_bucket { + buckets.push(TrendBucket { + bucket: bucket_label(today_ts, bucket), + build_count: t.build_count, + build_total_ms: t.build_total_ms, + subst_count: t.subst_count, + subst_total_ms: t.subst_total_ms, + download_bytes: t.download_bytes, + build_durations: vec![], + subst_durations: vec![], + }); + } else if let Some(last) = buckets.last_mut() { + // Today falls in the same week/month bucket as the last historical day. + last.build_count += t.build_count; + last.build_total_ms += t.build_total_ms; + last.subst_count += t.subst_count; + last.subst_total_ms += t.subst_total_ms; + last.download_bytes += t.download_bytes; + } + } + } + } + + Ok(buckets) +} + +pub fn collect_trend( + db: &Mutex, + since: Option, + bucket: BucketSize, + drv: Option, + today: Option, + full: bool, +) -> Result { + if let Some(ref d) = drv { + assert!(!d.is_empty(), "drv filter must not be empty"); + } + + let conn = db.lock().unwrap(); + + // daily_stats path when individual samples are not needed, no drv filter is active + // (daily_stats has no per-drv breakdown), and bucket granularity is at least a day + // (daily_stats has no intra-day resolution). + let buckets = if !full && drv.is_none() && !matches!(bucket, BucketSize::Hour) { + trend_from_daily_stats(&conn, since, &bucket, today)? + } else { + trend_from_events(&conn, since, &bucket, drv.as_deref())? + }; + for i in 1..buckets.len() { assert!(buckets[i].bucket > buckets[i - 1].bucket, "buckets must be strictly ascending"); } @@ -526,7 +640,7 @@ pub fn display_trend(trend: &Trend) { println!("filter: {}", drv); } - print!("{:6} {:>10} {:>6} {:>10}", "period", "builds", "build med", "subst", "subst med"); + print!("{:6} {:>10} {:>6} {:>10}", "period", "builds", "build avg", "subst", "subst avg"); if has_downloads { print!(" {:>8}", "dl (MB)"); } println!(); @@ -536,14 +650,12 @@ 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 }; + let build_avg = if b.build_count > 0 { b.build_total_ms / b.build_count } else { 0 }; + let subst_avg = if b.subst_count > 0 { b.subst_total_ms / b.subst_count } else { 0 }; 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(build_avg), + b.subst_count, fmt_ms(subst_avg)); if has_downloads { print!(" {:>8.1}", b.download_bytes as f64 / 1_048_576.0); } println!(); } @@ -606,15 +718,13 @@ 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"); + println!("period,build_count,build_avg_ms,subst_count,subst_avg_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); + let build_avg = if b.build_count > 0 { b.build_total_ms / b.build_count } else { 0 }; + let subst_avg = if b.subst_count > 0 { b.subst_total_ms / b.subst_count } else { 0 }; + assert!(build_avg >= 0); + assert!(subst_avg >= 0); + println!("{},{},{},{},{},{}", b.bucket, b.build_count, build_avg, b.subst_count, subst_avg, b.download_bytes); } }