diff --git a/.tangled/workflows/main.yml b/.tangled/workflows/main.yml new file mode 100644 index 0000000..b09fd3f --- /dev/null +++ b/.tangled/workflows/main.yml @@ -0,0 +1,21 @@ +when: + - event: ["pull_request", "push", "manual"] + branch: ["master"] + +engine: microvm +image: nixos + +dependencies: + - gcc + - cargo + - rustc + - clippy + - rustfmt + +steps: + - name: "Check formatting" + command: cargo fmt --check + - name: "Clippy" + command: cargo clippy + - name: "Test" + command: cargo test --all diff --git a/src/daemon.rs b/src/daemon.rs index 603e658..766801f 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -1,6 +1,6 @@ use crate::db; -use crate::model::{self, gc_record, Activity, ActivityId, Record}; -use crate::writer::{run_writer, Config}; +use crate::model::{self, Activity, ActivityId, Record, gc_record}; +use crate::writer::{Config, run_writer}; use anyhow::{Context, Result}; use chrono::{DateTime, Utc}; use std::collections::HashMap; @@ -8,7 +8,7 @@ use std::io::{BufRead, BufReader}; use std::os::unix::io::AsRawFd; use std::os::unix::net::{UnixListener, UnixStream}; use std::path::{Path, PathBuf}; -use std::sync::mpsc::{channel, Sender}; +use std::sync::mpsc::{Sender, channel}; use std::thread; use std::time::Duration; @@ -52,7 +52,10 @@ pub fn run_daemon( .context("failed to spawn gc thread")?; info!("GC tracking enabled ({})", nix_state_dir.display()); } else { - warn!("gc.lock not found at {} — GC tracking disabled", gc_lock.display()); + warn!( + "gc.lock not found at {} — GC tracking disabled", + gc_lock.display() + ); } for stream in listener.incoming() { @@ -80,7 +83,9 @@ fn handle_conn(stream: UnixStream, tx: Sender) -> Result<()> { for line in reader.lines() { let line = line.context("read line")?; - let Some(ev) = model::parse(&line) else { continue }; + let Some(ev) = model::parse(&line) else { + continue; + }; if let Some(rec) = model::apply(&mut active, ev, Utc::now()) && tx.send(rec).is_err() { diff --git a/src/db.rs b/src/db.rs index 6b0bb17..f5ce2ac 100644 --- a/src/db.rs +++ b/src/db.rs @@ -43,7 +43,10 @@ pub fn open(path: &Path) -> Result { let version: u32 = conn .query_row("PRAGMA user_version", [], |r| r.get(0)) .unwrap_or(0); - assert!(version <= SCHEMA_VERSION + 1_000_000, "implausible user_version"); + assert!( + version <= SCHEMA_VERSION + 1_000_000, + "implausible user_version" + ); if version > SCHEMA_VERSION { anyhow::bail!( @@ -55,7 +58,8 @@ pub fn open(path: &Path) -> Result { ); } - conn.execute_batch(SCHEMA).context("Failed to create schema")?; + conn.execute_batch(SCHEMA) + .context("Failed to create schema")?; conn.execute_batch(&format!("PRAGMA user_version = {SCHEMA_VERSION}")) .context("Failed to set schema version")?; @@ -108,7 +112,9 @@ mod tests { assert_eq!(mode.to_lowercase(), "wal"); // auto_vacuum INCREMENTAL == 2. - let av: i64 = conn.query_row("PRAGMA auto_vacuum", [], |r| r.get(0)).unwrap(); + let av: i64 = conn + .query_row("PRAGMA auto_vacuum", [], |r| r.get(0)) + .unwrap(); assert_eq!(av, 2, "auto_vacuum must be INCREMENTAL"); } diff --git a/src/lib.rs b/src/lib.rs index 5ba8c85..907c612 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -3,8 +3,8 @@ #[macro_use] pub mod log; -pub mod db; pub mod daemon; +pub mod db; pub mod model; pub mod query; pub mod writer; diff --git a/src/main.rs b/src/main.rs index 8e3d811..8306d2a 100644 --- a/src/main.rs +++ b/src/main.rs @@ -4,11 +4,11 @@ use clap::{Parser, Subcommand}; use std::path::PathBuf; use std::time::Duration; -use nod::db; use nod::daemon::{default_db_path, default_socket_path, run_daemon}; +use nod::db; use nod::query::{ - collect_stats, collect_trend, display_stats, display_trend, - output_csv_trend, BucketSize, SortField, + BucketSize, SortField, collect_stats, collect_trend, display_stats, display_trend, + output_csv_trend, }; use nod::writer::Config; @@ -85,7 +85,9 @@ fn parse_since(spec: &str) -> Result> { .find(|c: char| !c.is_ascii_digit()) .with_context(|| format!("invalid --since '{spec}': expected , e.g. 30d"))?; let (num, unit) = spec.split_at(split); - let n: i64 = num.parse().with_context(|| format!("invalid --since '{spec}'"))?; + let n: i64 = num + .parse() + .with_context(|| format!("invalid --since '{spec}'"))?; let now = Utc::now(); let t = match unit { "h" => now - chrono::Duration::hours(n), @@ -103,20 +105,31 @@ fn main() -> Result<()> { let db_path = cli.db.clone().unwrap_or_else(default_db_path); match cli.command { - Some(Commands::Daemon { retain_days, nix_state_dir, maint_secs, debug }) => { + Some(Commands::Daemon { + retain_days, + nix_state_dir, + maint_secs, + debug, + }) => { nod::log::set_debug(debug); assert!(retain_days > 0, "--retain-days must be > 0"); assert!(maint_secs > 0); let socket_path = cli.socket.clone().unwrap_or_else(default_socket_path); - for p in [db_path.parent(), socket_path.parent()].into_iter().flatten() { + for p in [db_path.parent(), socket_path.parent()] + .into_iter() + .flatten() + { if !p.as_os_str().is_empty() { std::fs::create_dir_all(p) .with_context(|| format!("failed to create {}", p.display()))?; } } - let cfg = Config { retain_days, interval: Duration::from_secs(maint_secs) }; + let cfg = Config { + retain_days, + interval: Duration::from_secs(maint_secs), + }; run_daemon(&socket_path, &db_path, cfg, &nix_state_dir).context("daemon failed")?; } @@ -137,8 +150,12 @@ fn main() -> Result<()> { None => { assert!(cli.limit > 0, "--limit must be > 0"); let since = cli.since.as_deref().map(parse_since).transpose()?; - let conn = db::open_readonly(&db_path) - .with_context(|| format!("no database at {} (is the daemon running?)", db_path.display()))?; + let conn = db::open_readonly(&db_path).with_context(|| { + format!( + "no database at {} (is the daemon running?)", + db_path.display() + ) + })?; if let Some(bucket) = cli.bucket { let trend = collect_trend(&conn, since, bucket, cli.drv, cli.exclude)?; @@ -147,7 +164,14 @@ fn main() -> Result<()> { OutputFormat::Csv => output_csv_trend(&trend), } } else { - let stats = collect_stats(&conn, since, cli.drv.as_deref(), cli.exclude.as_deref(), cli.sort, cli.limit)?; + let stats = collect_stats( + &conn, + since, + cli.drv.as_deref(), + cli.exclude.as_deref(), + cli.sort, + cli.limit, + )?; display_stats(&stats); } } diff --git a/src/model.rs b/src/model.rs index 6798d29..7140307 100644 --- a/src/model.rs +++ b/src/model.rs @@ -147,7 +147,10 @@ pub struct Activity { impl Activity { fn field(&self, i: usize) -> Option { - self.fields.get(i).and_then(|v| v.as_str()).map(String::from) + self.fields + .get(i) + .and_then(|v| v.as_str()) + .map(String::from) } } @@ -159,16 +162,35 @@ pub fn parse(line: &str) -> Option { serde_json::from_str::(line).ok() } -pub fn apply(active: &mut HashMap, ev: NixEvent, now: DateTime) -> Option { +pub fn apply( + active: &mut HashMap, + ev: NixEvent, + now: DateTime, +) -> Option { match ev { - NixEvent::Start { id, event_type, fields, .. } => { + NixEvent::Start { + id, + event_type, + fields, + .. + } => { active.insert( ActivityId(id), - Activity { event_type: ActivityType::from_code(event_type), start: now, fields, total_bytes: 0, phases: Vec::new() }, + Activity { + event_type: ActivityType::from_code(event_type), + start: now, + fields, + total_bytes: 0, + phases: Vec::new(), + }, ); None } - NixEvent::Result { id, event_type, fields } => { + NixEvent::Result { + id, + event_type, + fields, + } => { if let Some(act) = active.get_mut(&ActivityId(id)) { match ResultType::from_code(event_type) { ResultType::Progress => { @@ -180,7 +202,8 @@ pub fn apply(active: &mut HashMap, ev: NixEvent, now: Date } ResultType::SetPhase => { if let Some(name) = fields.first().and_then(|v| v.as_str()) { - act.phases.push((name.to_string(), (now - act.start).num_milliseconds())); + act.phases + .push((name.to_string(), (now - act.start).num_milliseconds())); } } _ => {} @@ -203,7 +226,13 @@ pub fn apply(active: &mut HashMap, ev: NixEvent, now: Date let phases = if act.phases.is_empty() { None } else { - Some(act.phases.iter().map(|(n, s)| format!("{n}:{s}")).collect::>().join(",")) + Some( + act.phases + .iter() + .map(|(n, s)| format!("{n}:{s}")) + .collect::>() + .join(","), + ) }; Some(Record { @@ -266,9 +295,19 @@ mod tests { #[test] fn substitution_captures_cache_url() { let mut active = HashMap::new(); - let s = parse(&start(7, ActivityType::Substitute, &["/nix/store/x-bar", "https://cache.nixos.org"])).unwrap(); + let s = parse(&start( + 7, + ActivityType::Substitute, + &["/nix/store/x-bar", "https://cache.nixos.org"], + )) + .unwrap(); apply(&mut active, s, at(0)); - let rec = apply(&mut active, parse(r#"{"action":"stop","id":7}"#).unwrap(), at(1000)).unwrap(); + let rec = apply( + &mut active, + parse(r#"{"action":"stop","id":7}"#).unwrap(), + at(1000), + ) + .unwrap(); assert_eq!(rec.cache_url.as_deref(), Some("https://cache.nixos.org")); } @@ -276,23 +315,50 @@ mod tests { fn unpersisted_type_yields_no_record() { let mut active = HashMap::new(); // Builds (104) is a grouping activity, not persisted. - apply(&mut active, parse(&start(3, ActivityType::Builds, &[])).unwrap(), at(0)); - assert!(apply(&mut active, parse(r#"{"action":"stop","id":3}"#).unwrap(), at(5)).is_none()); + apply( + &mut active, + parse(&start(3, ActivityType::Builds, &[])).unwrap(), + at(0), + ); + assert!( + apply( + &mut active, + parse(r#"{"action":"stop","id":3}"#).unwrap(), + at(5) + ) + .is_none() + ); } #[test] fn stop_without_start_is_ignored() { let mut active = HashMap::new(); - assert!(apply(&mut active, parse(r#"{"action":"stop","id":99}"#).unwrap(), at(5)).is_none()); + assert!( + apply( + &mut active, + parse(r#"{"action":"stop","id":99}"#).unwrap(), + at(5) + ) + .is_none() + ); } #[test] fn progress_sets_bytes_for_download() { let mut active = HashMap::new(); - apply(&mut active, parse(&start(2, ActivityType::FileTransfer, &[])).unwrap(), at(0)); + apply( + &mut active, + parse(&start(2, ActivityType::FileTransfer, &[])).unwrap(), + at(0), + ); let prog = parse(r#"{"action":"result","id":2,"type":105,"fields":[100,4096]}"#).unwrap(); apply(&mut active, prog, at(100)); - let rec = apply(&mut active, parse(r#"{"action":"stop","id":2}"#).unwrap(), at(500)).unwrap(); + let rec = apply( + &mut active, + parse(r#"{"action":"stop","id":2}"#).unwrap(), + at(500), + ) + .unwrap(); assert_eq!(rec.bytes, 4096); assert_eq!(rec.event_type, ActivityType::FileTransfer); } @@ -313,8 +379,22 @@ mod tests { #[test] fn file_transfer_does_not_store_url_as_drv() { let mut active = HashMap::new(); - apply(&mut active, parse(&start(4, ActivityType::FileTransfer, &["https://example/foo.nar"])).unwrap(), at(0)); - let rec = apply(&mut active, parse(r#"{"action":"stop","id":4}"#).unwrap(), at(100)).unwrap(); + apply( + &mut active, + parse(&start( + 4, + ActivityType::FileTransfer, + &["https://example/foo.nar"], + )) + .unwrap(), + at(0), + ); + let rec = apply( + &mut active, + parse(r#"{"action":"stop","id":4}"#).unwrap(), + at(100), + ) + .unwrap(); assert_eq!(rec.drv, None); assert_eq!(rec.cache_url, None); } @@ -322,18 +402,30 @@ mod tests { #[test] fn build_captures_phases() { let mut active = HashMap::new(); - apply(&mut active, parse(&start(1, ActivityType::Build, &["/nix/store/x.drv"])).unwrap(), at(0)); + apply( + &mut active, + parse(&start(1, ActivityType::Build, &["/nix/store/x.drv"])).unwrap(), + at(0), + ); // result type 104 = SetPhase, fields[0] = phase name. let phase = r#"{"action":"result","id":1,"type":104,"fields":["buildPhase"]}"#; apply(&mut active, parse(phase).unwrap(), at(2000)); - let rec = apply(&mut active, parse(r#"{"action":"stop","id":1}"#).unwrap(), at(5000)).unwrap(); + let rec = apply( + &mut active, + parse(r#"{"action":"stop","id":1}"#).unwrap(), + at(5000), + ) + .unwrap(); assert_eq!(rec.phases.as_deref(), Some("buildPhase:2000")); } #[test] fn activity_type_roundtrips_known_codes() { for code in [0u64, 100, 101, 105, 108, 109, 112, 200] { - assert_eq!(ActivityType::from_code(code) as i64 as u64, if code == 0 { 0 } else { code }); + assert_eq!( + ActivityType::from_code(code) as i64 as u64, + if code == 0 { 0 } else { code } + ); } assert_eq!(ActivityType::from_code(9999), ActivityType::Unknown); } @@ -348,7 +440,8 @@ mod tests { leaf.prop_recursive(3, 16, 5, |inner| { prop_oneof![ prop::collection::vec(inner.clone(), 0..4).prop_map(Value::from), - prop::collection::vec((".*", inner), 0..4).prop_map(|kvs| Value::Object(kvs.into_iter().collect())), + prop::collection::vec((".*", inner), 0..4) + .prop_map(|kvs| Value::Object(kvs.into_iter().collect())), ] }) } @@ -362,17 +455,37 @@ mod tests { // wrong-type near-misses, and plain non-events. fn event_object() -> impl Strategy { ( - opt(prop_oneof![Just(Value::from("start")), Just(Value::from("stop")), Just(Value::from("result")), ".*".prop_map(Value::from)]), - opt(prop_oneof![any::().prop_map(Value::from), ".*".prop_map(Value::from), Just(Value::Null)]), - opt(prop_oneof![(0u64..250).prop_map(Value::from), any::().prop_map(Value::from)]), + opt(prop_oneof![ + Just(Value::from("start")), + Just(Value::from("stop")), + Just(Value::from("result")), + ".*".prop_map(Value::from) + ]), + opt(prop_oneof![ + any::().prop_map(Value::from), + ".*".prop_map(Value::from), + Just(Value::Null) + ]), + opt(prop_oneof![ + (0u64..250).prop_map(Value::from), + any::().prop_map(Value::from) + ]), opt(prop::collection::vec(arb_json(), 0..6).prop_map(Value::from)), ) .prop_map(|(action, id, ty, fields)| { let mut m = serde_json::Map::new(); - if let Some(v) = action { m.insert("action".into(), v); } - if let Some(v) = id { m.insert("id".into(), v); } - if let Some(v) = ty { m.insert("type".into(), v); } - if let Some(v) = fields { m.insert("fields".into(), v); } + if let Some(v) = action { + m.insert("action".into(), v); + } + if let Some(v) = id { + m.insert("id".into(), v); + } + if let Some(v) = ty { + m.insert("type".into(), v); + } + if let Some(v) = fields { + m.insert("fields".into(), v); + } Value::Object(m).to_string() }) } diff --git a/src/query.rs b/src/query.rs index 2b60a25..2a256ca 100644 --- a/src/query.rs +++ b/src/query.rs @@ -41,7 +41,13 @@ fn since_secs(since: Option>) -> Option { since.map(|d| d.timestamp()) } -fn type_totals(conn: &Connection, ty: ActivityType, since: Option>, drv: Option<&str>, exclude: Option<&str>) -> Result<(i64, Duration)> { +fn type_totals( + conn: &Connection, + ty: ActivityType, + since: Option>, + drv: Option<&str>, + exclude: Option<&str>, +) -> Result<(i64, Duration)> { conn.query_row( "SELECT COUNT(*), COALESCE(SUM(duration_ms),0) FROM events INDEXED BY idx_events_type_ts @@ -50,7 +56,12 @@ fn type_totals(conn: &Connection, ty: ActivityType, since: Option> AND (?3 IS NULL OR drv LIKE '%' || ?3 || '%') AND (?4 IS NULL OR drv IS NULL OR drv NOT LIKE '%' || ?4 || '%')", rusqlite::params![ty, since_secs(since), drv, exclude], - |r| Ok((r.get::<_, i64>(0)?, Duration::milliseconds(r.get::<_, i64>(1)?))), + |r| { + Ok(( + r.get::<_, i64>(0)?, + Duration::milliseconds(r.get::<_, i64>(1)?), + )) + }, ) .context("type_totals") } @@ -61,7 +72,12 @@ fn download_totals(conn: &Connection, since: Option>) -> Result<(i FROM events INDEXED BY idx_events_type_ts WHERE event_type = ?1 AND (?2 IS NULL OR ts >= ?2)", rusqlite::params![ActivityType::FileTransfer, since_secs(since)], - |r| Ok((r.get::<_, i64>(0)?, Duration::milliseconds(r.get::<_, i64>(1)?))), + |r| { + Ok(( + r.get::<_, i64>(0)?, + Duration::milliseconds(r.get::<_, i64>(1)?), + )) + }, ) .context("download_totals") } @@ -74,9 +90,16 @@ fn cache_latency(conn: &Connection, since: Option>) -> Result 0, "limit must be > 0"); let (build_count, build_total) = type_totals(conn, ActivityType::Build, since, drv, exclude)?; - let (subst_count, subst_total) = type_totals(conn, ActivityType::Substitute, since, drv, exclude)?; + let (subst_count, subst_total) = + type_totals(conn, ActivityType::Substitute, since, drv, exclude)?; // Downloads and GC have no derivation, so the name filters don't apply to them // (they're shown in full unless a positive drv filter narrows to one package). - let (download_bytes, download_time) = - if drv.is_none() { download_totals(conn, since)? } else { (0, Duration::zero()) }; - let (gc_count, gc_total) = - if drv.is_none() { type_totals(conn, ActivityType::GarbageCollect, since, None, None)? } else { (0, Duration::zero()) }; - let cache_latency = if drv.is_none() { cache_latency(conn, since)? } else { Vec::new() }; + let (download_bytes, download_time) = if drv.is_none() { + download_totals(conn, since)? + } else { + (0, Duration::zero()) + }; + let (gc_count, gc_total) = if drv.is_none() { + type_totals(conn, ActivityType::GarbageCollect, since, None, None)? + } else { + (0, Duration::zero()) + }; + let cache_latency = if drv.is_none() { + cache_latency(conn, since)? + } else { + Vec::new() + }; let slowest_builds = slowest_builds(conn, since, drv, exclude, &sort, limit)?; @@ -142,9 +176,15 @@ fn slowest_builds( ); let mut stmt = conn.prepare(&sql)?; Ok(stmt - .query_map(rusqlite::params![since_secs(since), drv, limit, exclude], |r| { - Ok(SlowBuild { duration: Duration::milliseconds(r.get::<_, i64>(0)?), drv: r.get(1)? }) - })? + .query_map( + rusqlite::params![since_secs(since), drv, limit, exclude], + |r| { + Ok(SlowBuild { + duration: Duration::milliseconds(r.get::<_, i64>(0)?), + drv: r.get(1)?, + }) + }, + )? .filter_map(|r| r.ok()) .collect()) } @@ -215,21 +255,27 @@ fn bucket_label(dt: DateTime, size: &BucketSize) -> String { // Start of the bucket after the one containing dt. fn bucket_end(dt: DateTime, size: &BucketSize) -> DateTime { match size { - BucketSize::Hour => Utc - .with_ymd_and_hms(dt.year(), dt.month(), dt.day(), dt.hour(), 0, 0) - .unwrap() - + Duration::hours(1), - BucketSize::Day => Utc - .with_ymd_and_hms(dt.year(), dt.month(), dt.day(), 0, 0, 0) - .unwrap() - + Duration::days(1), + BucketSize::Hour => { + Utc.with_ymd_and_hms(dt.year(), dt.month(), dt.day(), dt.hour(), 0, 0) + .unwrap() + + Duration::hours(1) + } + BucketSize::Day => { + Utc.with_ymd_and_hms(dt.year(), dt.month(), dt.day(), 0, 0, 0) + .unwrap() + + Duration::days(1) + } BucketSize::Week => { let days_since_monday = dt.weekday().num_days_from_monday() as i64; let monday = dt.date_naive() - Duration::days(days_since_monday); Utc.from_utc_datetime(&monday.and_hms_opt(0, 0, 0).unwrap()) + Duration::weeks(1) } BucketSize::Month => { - let (y, m) = if dt.month() == 12 { (dt.year() + 1, 1) } else { (dt.year(), dt.month() + 1) }; + let (y, m) = if dt.month() == 12 { + (dt.year() + 1, 1) + } else { + (dt.year(), dt.month() + 1) + }; Utc.with_ymd_and_hms(y, m, 1, 0, 0, 0).unwrap() } } @@ -249,20 +295,38 @@ pub fn collect_trend( let buckets = build_trend(conn, since, &bucket, drv.as_deref(), exclude.as_deref())?; for i in 1..buckets.len() { - assert!(buckets[i].bucket > buckets[i - 1].bucket, "buckets must be strictly ascending"); + assert!( + buckets[i].bucket > buckets[i - 1].bucket, + "buckets must be strictly ascending" + ); } - Ok(Trend { buckets, bucket_size: bucket, drv_filter: drv }) + Ok(Trend { + buckets, + bucket_size: bucket, + drv_filter: drv, + }) } -fn push_or_merge(buckets: &mut Vec, next: &mut Option>, ts: DateTime, size: &BucketSize) { +fn push_or_merge( + buckets: &mut Vec, + next: &mut Option>, + ts: DateTime, + size: &BucketSize, +) { if next.is_none_or(|nb| ts >= nb) { buckets.push(TrendBucket::empty(bucket_label(ts, size))); *next = Some(bucket_end(ts, size)); } } -fn build_trend(conn: &Connection, since: Option>, bucket: &BucketSize, drv: Option<&str>, exclude: Option<&str>) -> Result> { +fn build_trend( + conn: &Connection, + since: Option>, + bucket: &BucketSize, + drv: Option<&str>, + exclude: Option<&str>, +) -> Result> { let mut stmt = conn.prepare( "SELECT ts, event_type, duration_ms, bytes FROM events INDEXED BY idx_events_ts WHERE event_type IN (?1, ?2, ?3) @@ -276,7 +340,14 @@ fn build_trend(conn: &Connection, since: Option>, bucket: &BucketS let mut next: Option> = None; for row in stmt .query_map( - rusqlite::params![ActivityType::FileTransfer, ActivityType::Build, ActivityType::Substitute, since_secs(since), drv, exclude], + rusqlite::params![ + ActivityType::FileTransfer, + ActivityType::Build, + ActivityType::Substitute, + since_secs(since), + drv, + exclude + ], |r| { Ok(( DateTime::from_timestamp(r.get::<_, i64>(0)?, 0).unwrap_or_default(), @@ -308,7 +379,11 @@ fn build_trend(conn: &Connection, since: Option>, bucket: &BucketS } fn avg(total: Duration, n: i64) -> Duration { - if n > 0 { Duration::milliseconds(total.num_milliseconds() / n) } else { Duration::zero() } + if n > 0 { + Duration::milliseconds(total.num_milliseconds() / n) + } else { + Duration::zero() + } } pub fn fmt_dur(d: Duration) -> String { @@ -320,9 +395,21 @@ pub fn fmt_dur(d: Duration) -> String { } else if ms < 3_600_000 { format!("{}m{:.1}s", ms / 60_000, (ms % 60_000) as f64 / 1000.0) } else { - let (h, m, s) = (ms / 3_600_000, (ms % 3_600_000) / 60_000, (ms % 60_000) / 1000); - let m = if m > 0 { format!("{m}m") } else { String::new() }; - let s = if s > 0 { format!("{s}s") } else { String::new() }; + let (h, m, s) = ( + ms / 3_600_000, + (ms % 3_600_000) / 60_000, + (ms % 60_000) / 1000, + ); + let m = if m > 0 { + format!("{m}m") + } else { + String::new() + }; + let s = if s > 0 { + format!("{s}s") + } else { + String::new() + }; format!("{h}h{m}{s}") } } @@ -334,13 +421,41 @@ fn drv_name(path: &str) -> &str { pub fn display_stats(stats: &Stats) { let mb = stats.download_bytes as f64 / 1_048_576.0; let dl_ms = stats.download_time.num_milliseconds(); - let dl_speed = if dl_ms > 0 { mb / (dl_ms as f64 / 1000.0) } else { 0.0 }; + let dl_speed = if dl_ms > 0 { + mb / (dl_ms as f64 / 1000.0) + } else { + 0.0 + }; - println!("{:<14} {:>6} total {:>9} avg {:>9}", "built", stats.build_count, fmt_dur(stats.build_total), fmt_dur(avg(stats.build_total, stats.build_count))); - println!("{:<14} {:>6} total {:>9} avg {:>9}", "substituted", stats.subst_count, fmt_dur(stats.subst_total), fmt_dur(avg(stats.subst_total, stats.subst_count))); - println!("{:<14} {:>6} total {:>9} avg {:>9}", "downloaded", "", format!("{mb:.1}MB"), format!("{dl_speed:.1}MB/s")); + println!( + "{:<14} {:>6} total {:>9} avg {:>9}", + "built", + stats.build_count, + fmt_dur(stats.build_total), + fmt_dur(avg(stats.build_total, stats.build_count)) + ); + println!( + "{:<14} {:>6} total {:>9} avg {:>9}", + "substituted", + stats.subst_count, + fmt_dur(stats.subst_total), + fmt_dur(avg(stats.subst_total, stats.subst_count)) + ); + println!( + "{:<14} {:>6} total {:>9} avg {:>9}", + "downloaded", + "", + format!("{mb:.1}MB"), + format!("{dl_speed:.1}MB/s") + ); if stats.gc_count > 0 { - println!("{:<14} {:>6} total {:>9} avg {:>9}", "gc", stats.gc_count, fmt_dur(stats.gc_total), fmt_dur(avg(stats.gc_total, stats.gc_count))); + println!( + "{:<14} {:>6} total {:>9} avg {:>9}", + "gc", + stats.gc_count, + fmt_dur(stats.gc_total), + fmt_dur(avg(stats.gc_total, stats.gc_count)) + ); } if !stats.slowest_builds.is_empty() { @@ -352,22 +467,37 @@ pub fn display_stats(stats: &Stats) { } if !stats.cache_latency.is_empty() { - let url_w = stats.cache_latency.iter().map(|r| r.cache_url.len()).max().unwrap_or(0); + let url_w = stats + .cache_latency + .iter() + .map(|r| r.cache_url.len()) + .max() + .unwrap_or(0); println!(); for row in &stats.cache_latency { - println!("{:7} {:>6} queries", row.cache_url, fmt_dur(Duration::milliseconds(row.avg_ms as i64)), row.count, width = url_w); + println!( + "{:7} {:>6} queries", + row.cache_url, + fmt_dur(Duration::milliseconds(row.avg_ms as i64)), + row.count, + width = url_w + ); } } } pub fn display_trend(trend: &Trend) { - let has_downloads = trend.drv_filter.is_none() && trend.buckets.iter().any(|b| b.download_bytes > 0); + 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 avg", "subst", "subst avg"); + print!( + "{:6} {:>10} {:>6} {:>10}", + "period", "builds", "build avg", "subst", "subst avg" + ); if has_downloads { print!(" {:>8}", "dl (MB)"); } @@ -383,12 +513,24 @@ pub fn display_trend(trend: &Trend) { "{:6} {:>10} {:>6} {:>10}", b.bucket, or_dash(b.build_count > 0, b.build_count.to_string()), - or_dash(b.build_count > 0, fmt_dur(avg(b.build_total, b.build_count))), + or_dash( + b.build_count > 0, + fmt_dur(avg(b.build_total, b.build_count)) + ), or_dash(b.subst_count > 0, b.subst_count.to_string()), - or_dash(b.subst_count > 0, fmt_dur(avg(b.subst_total, b.subst_count))), + or_dash( + b.subst_count > 0, + fmt_dur(avg(b.subst_total, b.subst_count)) + ), ); if has_downloads { - print!(" {:>8}", or_dash(b.download_bytes > 0, format!("{:.1}", b.download_bytes as f64 / 1_048_576.0))); + print!( + " {:>8}", + or_dash( + b.download_bytes > 0, + format!("{:.1}", b.download_bytes as f64 / 1_048_576.0) + ) + ); } println!(); } @@ -399,7 +541,10 @@ pub fn output_csv_trend(trend: &Trend) { for b in &trend.buckets { let build_avg = avg(b.build_total, b.build_count).num_milliseconds(); let subst_avg = avg(b.subst_total, b.subst_count).num_milliseconds(); - println!("{},{},{},{},{},{}", b.bucket, b.build_count, build_avg, b.subst_count, subst_avg, b.download_bytes); + println!( + "{},{},{},{},{},{}", + b.bucket, b.build_count, build_avg, b.subst_count, subst_avg, b.download_bytes + ); } } @@ -418,7 +563,15 @@ mod tests { conn } - fn ins(conn: &Connection, ts: i64, ty: ActivityType, dur_ms: i64, drv: Option<&str>, cache: Option<&str>, bytes: i64) { + fn ins( + conn: &Connection, + ts: i64, + ty: ActivityType, + dur_ms: i64, + drv: Option<&str>, + cache: Option<&str>, + bytes: i64, + ) { conn.execute( "INSERT INTO events (ts, event_type, duration_ms, drv, cache_url, bytes) VALUES (?1,?2,?3,?4,?5,?6)", rusqlite::params![ts, ty, dur_ms, drv, cache, bytes], @@ -454,7 +607,11 @@ mod tests { (TODAY_DAY - WINDOW)..=TODAY_DAY, 0i64..86_400, 0i64..600_000, - prop::sample::select(vec!["/nix/store/a-foo", "/nix/store/b-bar", "/nix/store/c-baz"]), + prop::sample::select(vec![ + "/nix/store/a-foo", + "/nix/store/b-bar", + "/nix/store/c-baz", + ]), prop::sample::select(vec!["https://cache.nixos.org", "https://my.cache"]), 0i64..500_000_000, ) @@ -464,9 +621,21 @@ mod tests { ts, ty, dur, - drv: if matches!(ty, ActivityType::Build | ActivityType::Substitute) { Some(drv.into()) } else { None }, - cache: if ty == ActivityType::Substitute { Some(cache.into()) } else { None }, - bytes: if ty == ActivityType::FileTransfer { bytes } else { 0 }, + drv: if matches!(ty, ActivityType::Build | ActivityType::Substitute) { + Some(drv.into()) + } else { + None + }, + cache: if ty == ActivityType::Substitute { + Some(cache.into()) + } else { + None + }, + bytes: if ty == ActivityType::FileTransfer { + bytes + } else { + 0 + }, } }) } @@ -478,7 +647,10 @@ mod tests { .prepare("INSERT INTO events (ts, event_type, duration_ms, drv, cache_url, bytes) VALUES (?1,?2,?3,?4,?5,?6)") .unwrap(); for e in events { - stmt.execute(rusqlite::params![e.ts, e.ty, e.dur, e.drv, e.cache, e.bytes]).unwrap(); + stmt.execute(rusqlite::params![ + e.ts, e.ty, e.dur, e.drv, e.cache, e.bytes + ]) + .unwrap(); } drop(stmt); conn @@ -495,7 +667,10 @@ mod tests { } fn day_label(ts: i64) -> String { - DateTime::::from_timestamp(ts, 0).unwrap().format("%Y-%m-%d").to_string() + DateTime::::from_timestamp(ts, 0) + .unwrap() + .format("%Y-%m-%d") + .to_string() } // Substring filters against the three fixed derivations (a-foo / b-bar / c-baz): @@ -521,8 +696,24 @@ mod tests { #[test] fn cache_latency_averages_per_url() { let conn = db(); - ins(&conn, DAY, ActivityType::Substitute, 100, Some("/nix/store/a"), Some("https://c"), 0); - ins(&conn, DAY + 1, ActivityType::Substitute, 300, Some("/nix/store/b"), Some("https://c"), 0); + ins( + &conn, + DAY, + ActivityType::Substitute, + 100, + Some("/nix/store/a"), + Some("https://c"), + 0, + ); + ins( + &conn, + DAY + 1, + ActivityType::Substitute, + 300, + Some("/nix/store/b"), + Some("https://c"), + 0, + ); let s = collect_stats(&conn, None, None, None, SortField::Duration, 10).unwrap(); assert_eq!(s.cache_latency.len(), 1); assert_eq!(s.cache_latency[0].cache_url, "https://c"); diff --git a/src/writer.rs b/src/writer.rs index 514304a..5e63a3c 100644 --- a/src/writer.rs +++ b/src/writer.rs @@ -15,7 +15,12 @@ pub struct Config { pub interval: Duration, } -pub fn run_writer(conn: Connection, rx: Receiver, cfg: Config, now: impl Fn() -> DateTime) -> Result<()> { +pub fn run_writer( + conn: Connection, + rx: Receiver, + cfg: Config, + now: impl Fn() -> DateTime, +) -> Result<()> { assert!(cfg.retain_days > 0, "retain_days must be > 0"); let mut batch: Vec = Vec::with_capacity(MAX_BATCH); @@ -143,11 +148,16 @@ mod tests { #[test] fn flush_inserts_all() { let conn = fresh(); - let mut batch = vec![rec(1, ActivityType::Build, 10), rec(2, ActivityType::Build, 20)]; + let mut batch = vec![ + rec(1, ActivityType::Build, 10), + rec(2, ActivityType::Build, 20), + ]; flush(&conn, &mut batch).unwrap(); assert!(batch.is_empty()); let (n, sum): (i64, i64) = conn - .query_row("SELECT COUNT(*), SUM(duration_ms) FROM events", [], |r| Ok((r.get(0)?, r.get(1)?))) + .query_row("SELECT COUNT(*), SUM(duration_ms) FROM events", [], |r| { + Ok((r.get(0)?, r.get(1)?)) + }) .unwrap(); assert_eq!((n, sum), (2, 30)); } @@ -155,23 +165,41 @@ mod tests { #[test] fn prune_keeps_window_drops_older() { let conn = fresh(); - let mut batch = vec![rec(DAY, ActivityType::Build, 10), rec(9 * DAY, ActivityType::Build, 10)]; + let mut batch = vec![ + rec(DAY, ActivityType::Build, 10), + rec(9 * DAY, ActivityType::Build, 10), + ]; flush(&conn, &mut batch).unwrap(); prune(&conn, 3, at(10 * DAY)).unwrap(); - let remaining: i64 = conn.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0)).unwrap(); + let remaining: i64 = conn + .query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0)) + .unwrap(); assert_eq!(remaining, 1); - let min_ts: i64 = conn.query_row("SELECT MIN(ts) FROM events", [], |r| r.get(0)).unwrap(); + let min_ts: i64 = conn + .query_row("SELECT MIN(ts) FROM events", [], |r| r.get(0)) + .unwrap(); assert_eq!(min_ts, 9 * DAY); } #[test] fn maintenance_prunes_old_rows() { let conn = fresh(); - let mut batch = vec![rec(2 * DAY, ActivityType::Build, 100), rec(10 * DAY, ActivityType::Build, 200)]; + let mut batch = vec![ + rec(2 * DAY, ActivityType::Build, 100), + rec(10 * DAY, ActivityType::Build, 200), + ]; flush(&conn, &mut batch).unwrap(); - let cfg = Config { retain_days: 3, interval: Duration::from_secs(1) }; + let cfg = Config { + retain_days: 3, + interval: Duration::from_secs(1), + }; run_maintenance(&conn, &cfg, at(10 * DAY + 100)).unwrap(); - let remaining: i64 = conn.query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0)).unwrap(); - assert_eq!(remaining, 1, "only rows within the retention window survive"); + let remaining: i64 = conn + .query_row("SELECT COUNT(*) FROM events", [], |r| r.get(0)) + .unwrap(); + assert_eq!( + remaining, 1, + "only rows within the retention window survive" + ); } } diff --git a/tests/integration.rs b/tests/integration.rs index 5b27c67..f233e47 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -1,8 +1,8 @@ use chrono::Utc; use nod::db::{open, open_readonly}; use nod::model::{ActivityType, Record}; -use nod::query::{collect_stats, SortField}; -use nod::writer::{run_writer, Config}; +use nod::query::{SortField, collect_stats}; +use nod::writer::{Config, run_writer}; use std::sync::mpsc::channel; use std::thread; use std::time::{Duration, Instant}; @@ -51,13 +51,19 @@ fn reads_stay_consistent_under_constant_writes() { max_latency = max_latency.max(t0.elapsed()); reads += 1; - assert!(s.build_count >= last, "build_count must be monotonic (snapshot isolation)"); + assert!( + s.build_count >= last, + "build_count must be monotonic (snapshot isolation)" + ); assert!(s.build_count <= N, "must never exceed what was sent"); last = s.build_count; if last >= N { break; } - assert!(Instant::now() < deadline, "did not converge; stuck at {last}"); + assert!( + Instant::now() < deadline, + "did not converge; stuck at {last}" + ); thread::sleep(Duration::from_millis(1)); } @@ -67,6 +73,11 @@ fn reads_stay_consistent_under_constant_writes() { let s = collect_stats(&rconn, None, None, None, SortField::Duration, 10).unwrap(); assert_eq!(s.build_count, N, "all rows must be durable after drain"); - eprintln!("concurrency: {N} rows, {reads} reads during ingest, max read latency {max_latency:?}"); - assert!(max_latency < Duration::from_secs(2), "reads must stay responsive under write load"); + eprintln!( + "concurrency: {N} rows, {reads} reads during ingest, max read latency {max_latency:?}" + ); + assert!( + max_latency < Duration::from_secs(2), + "reads must stay responsive under write load" + ); }