From 9cd15cfb27599d3957ab49f1ccd63d82cf7d027e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Vilhelm=20Bergs=C3=B8e?= Date: Wed, 22 Apr 2026 16:19:05 +0200 Subject: [PATCH] aggregation based query approach --- benches/queries.rs | 55 ++++++-- src/daemon.rs | 342 ++++++++++++++++++++++++++++++++++++++++++--- src/main.rs | 2 +- src/stats.rs | 169 ++++++++++++++++++---- 4 files changed, 512 insertions(+), 56 deletions(-) diff --git a/benches/queries.rs b/benches/queries.rs index 1661e5d..b4769c1 100644 --- a/benches/queries.rs +++ b/benches/queries.rs @@ -1,5 +1,5 @@ use criterion::{criterion_group, criterion_main, BenchmarkId, Criterion}; -use nod::stats::{collect_stats, collect_trend, BucketSize, SortField}; +use nod::stats::{collect_stats, collect_trend, BucketSize, SortField, TodaySummary}; use rusqlite::Connection; use std::sync::Mutex; @@ -17,8 +17,19 @@ fn seed(conn: &Connection, n: usize) { 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); + CREATE INDEX idx_events_type_start ON events(event_type, start_time, duration_ms, total_bytes); + CREATE INDEX idx_events_start_cover ON events(start_time, event_type, duration_ms, total_bytes); + CREATE TABLE daily_stats ( + day INTEGER NOT NULL, event_type INTEGER NOT NULL, + count INTEGER NOT NULL DEFAULT 0, total_ms INTEGER NOT NULL DEFAULT 0, + total_bytes INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, event_type) + ); + CREATE TABLE daily_cache_stats ( + day INTEGER NOT NULL, cache_url TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, total_ms INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, cache_url) + ); PRAGMA journal_mode = WAL; PRAGMA synchronous = NORMAL; ").unwrap(); @@ -49,29 +60,57 @@ fn seed(conn: &Connection, n: usize) { } conn.execute_batch("COMMIT").unwrap(); + + // Backfill daily_stats from seeded events (mirrors daemon migration v5). + conn.execute_batch(" + INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + SELECT start_time / 86400, event_type, + COUNT(*), COALESCE(SUM(duration_ms), 0), COALESCE(SUM(total_bytes), 0) + FROM events WHERE event_type IN (101, 105, 108) + GROUP BY 1, 2 + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms, + total_bytes = total_bytes + excluded.total_bytes; + INSERT INTO daily_cache_stats (day, cache_url, count, total_ms) + SELECT start_time / 86400, cache_url, COUNT(*), COALESCE(SUM(duration_ms), 0) + FROM events WHERE event_type = 108 AND cache_url IS NOT NULL + GROUP BY 1, 2 + ON CONFLICT (day, cache_url) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms; + ").unwrap(); } fn bench_collect_stats(c: &mut Criterion) { let mut group = c.benchmark_group("collect_stats"); + // All bench data is from 2023 (BASE_TIME). Using today's real Unix day means all + // seeded rows are closed historical days in daily_stats; today snapshot is empty. + // This matches normal daemon usage: summary from daily_stats, slowest_builds from events. + 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: no index benefit on start_time, exercises full scan path. + let t = today_snap.clone(); group.bench_with_input(BenchmarkId::new("all_time", n), &n, |b, _| { - b.iter(|| collect_stats(&db, None, None, SortField::Duration, 10, false).unwrap()) + b.iter(|| collect_stats(&db, None, None, SortField::Duration, 10, false, Some(t.clone())).unwrap()) }); - // Recent 10%: composite index (event_type, start_time) prunes 90% of rows. let since = Some(BASE_TIME + YEAR_SPAN * 9 / 10); + let t = today_snap.clone(); group.bench_with_input(BenchmarkId::new("recent_10pct", n), &n, |b, _| { - b.iter(|| collect_stats(&db, since, None, SortField::Duration, 10, false).unwrap()) + b.iter(|| collect_stats(&db, since, None, SortField::Duration, 10, false, Some(t.clone())).unwrap()) }); + let t = today_snap.clone(); group.bench_with_input(BenchmarkId::new("grouped", n), &n, |b, _| { - b.iter(|| collect_stats(&db, None, None, SortField::Count, 10, true).unwrap()) + b.iter(|| collect_stats(&db, None, None, SortField::Count, 10, true, Some(t.clone())).unwrap()) }); } diff --git a/src/daemon.rs b/src/daemon.rs index 7fc5575..c224225 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, collect_trend, BucketSize, SortField}; +use crate::stats::{collect_stats, collect_trend, BucketSize, SortField, TodaySummary}; const fn schema_hash(s: &[u8]) -> u32 { let mut h: u32 = 2166136261; @@ -34,14 +34,31 @@ CREATE TABLE IF NOT EXISTS events ( duration_ms INTEGER, 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); +CREATE INDEX IF NOT EXISTS idx_events_type_start ON events(event_type, start_time, duration_ms, total_bytes); +CREATE INDEX IF NOT EXISTS idx_events_start_cover ON events(start_time, event_type, duration_ms, total_bytes); +CREATE TABLE IF NOT EXISTS daily_stats ( + day INTEGER NOT NULL, + event_type INTEGER NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + total_ms INTEGER NOT NULL DEFAULT 0, + total_bytes INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, event_type) +); +CREATE TABLE IF NOT EXISTS daily_cache_stats ( + day INTEGER NOT NULL, + cache_url TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + total_ms INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, cache_url) +); "; 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 + 0xee061d32, // v2: INTEGER timestamps (Unix seconds), idx_events_type_start(event_type, start_time) + 0x7c09711e, // v3: idx_events_start_time(start_time) + 0x24268227, // v4: idx_events_type_start extended to cover (duration_ms, total_bytes); idx_events_start_cover added + 0xfcbb3598, // v5: daily_stats + daily_cache_stats tables ]; const SCHEMA_VERSION: u32 = SCHEMA_HASHES.len() as u32; const _: () = assert!( @@ -77,6 +94,52 @@ const MIGRATIONS: &[(u32, &str)] = &[ 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);"), + (4, " + DROP INDEX IF EXISTS idx_events_type_start; + DROP INDEX IF EXISTS idx_events_start_time; + CREATE INDEX idx_events_type_start ON events(event_type, start_time, duration_ms, total_bytes); + CREATE INDEX idx_events_start_cover ON events(start_time, event_type, duration_ms, total_bytes); + "), + // v5: create daily aggregate tables and backfill all closed days from events. + // Today's data is excluded (day < today) — the daemon rebuilds today from events on startup. + // This migration may take several seconds on large databases. + (5, " + CREATE TABLE IF NOT EXISTS daily_stats ( + day INTEGER NOT NULL, + event_type INTEGER NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + total_ms INTEGER NOT NULL DEFAULT 0, + total_bytes INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, event_type) + ); + CREATE TABLE IF NOT EXISTS daily_cache_stats ( + day INTEGER NOT NULL, + cache_url TEXT NOT NULL, + count INTEGER NOT NULL DEFAULT 0, + total_ms INTEGER NOT NULL DEFAULT 0, + PRIMARY KEY (day, cache_url) + ); + INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + SELECT start_time / 86400, event_type, + COUNT(*), COALESCE(SUM(duration_ms), 0), COALESCE(SUM(total_bytes), 0) + FROM events + WHERE event_type IN (101, 105, 108) + AND start_time / 86400 < strftime('%s', 'now') / 86400 + GROUP BY 1, 2 + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms, + total_bytes = total_bytes + excluded.total_bytes; + INSERT INTO daily_cache_stats (day, cache_url, count, total_ms) + SELECT start_time / 86400, cache_url, COUNT(*), COALESCE(SUM(duration_ms), 0) + FROM events + WHERE event_type = 108 AND cache_url IS NOT NULL + AND start_time / 86400 < strftime('%s', 'now') / 86400 + GROUP BY 1, 2 + ON CONFLICT (day, cache_url) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms; + "), ]; #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] @@ -235,8 +298,47 @@ struct Activity { total_bytes: u64, } +// Per-day aggregate cache, maintained in memory for the current calendar day (UTC). +// Flushed to daily_stats/daily_cache_stats when the day rolls over. +// Rebuilt from events on daemon startup via rebuild_today(). +struct TodayStats { + day: i64, // Unix day: start_time / 86400 + build_count: i64, + build_total_ms: i64, + subst_count: i64, + subst_total_ms: i64, + download_count: i64, + download_bytes: i64, + download_ms: i64, + cache: HashMap, // url → (total_ms, count) +} + +impl TodayStats { + fn new(day: i64) -> Self { + assert!(day > 0); + TodayStats { day, build_count: 0, build_total_ms: 0, subst_count: 0, + subst_total_ms: 0, download_count: 0, download_bytes: 0, + download_ms: 0, cache: HashMap::new() } + } + + fn snapshot(&self) -> TodaySummary { + TodaySummary { + day: self.day, + build_count: self.build_count, + build_total_ms: self.build_total_ms, + subst_count: self.subst_count, + subst_total_ms: self.subst_total_ms, + download_count: self.download_count, + download_bytes: self.download_bytes, + download_ms: self.download_ms, + cache: self.cache.clone(), + } + } +} + struct State { active_activities: HashMap, + today: TodayStats, } pub struct DbConnections { @@ -317,13 +419,157 @@ fn run_retention(conn: &Connection, retain_days: u32) -> Result<()> { Ok(()) } +// Flush a completed day's in-memory stats to the persistent daily_stats tables. +fn flush_today(conn: &Connection, today: &TodayStats) -> Result<()> { + assert!(today.day > 0); + let day = today.day; + + if today.build_count > 0 { + conn.execute( + "INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + VALUES (?1, 105, ?2, ?3, 0) + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms", + rusqlite::params![day, today.build_count, today.build_total_ms], + ).context("flush_today: build insert failed")?; + } + if today.subst_count > 0 { + conn.execute( + "INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + VALUES (?1, 108, ?2, ?3, 0) + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms", + rusqlite::params![day, today.subst_count, today.subst_total_ms], + ).context("flush_today: subst insert failed")?; + } + if today.download_count > 0 { + conn.execute( + "INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + VALUES (?1, 101, ?2, ?3, ?4) + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms, + total_bytes = total_bytes + excluded.total_bytes", + rusqlite::params![day, today.download_count, today.download_ms, today.download_bytes], + ).context("flush_today: download insert failed")?; + } + for (url, &(ms, cnt)) in &today.cache { + assert!(cnt > 0); + conn.execute( + "INSERT INTO daily_cache_stats (day, cache_url, count, total_ms) + VALUES (?1, ?2, ?3, ?4) + ON CONFLICT (day, cache_url) DO UPDATE SET + count = count + excluded.count, + total_ms = total_ms + excluded.total_ms", + rusqlite::params![day, url, cnt, ms], + ).context("flush_today: cache insert failed")?; + } + Ok(()) +} + +// Upsert a single past-day event directly into daily_stats. +// Called for events whose start_time predates the current calendar day +// (e.g. a long-running build that started yesterday and finished today). +fn upsert_daily( + conn: &Connection, day: i64, event_type: i64, + duration_ms: i64, total_bytes: i64, cache_url: Option<&str>, +) -> Result<()> { + assert!(day > 0); + conn.execute( + "INSERT INTO daily_stats (day, event_type, count, total_ms, total_bytes) + VALUES (?1, ?2, 1, ?3, ?4) + ON CONFLICT (day, event_type) DO UPDATE SET + count = count + 1, + total_ms = total_ms + excluded.total_ms, + total_bytes = total_bytes + excluded.total_bytes", + rusqlite::params![day, event_type, duration_ms, total_bytes], + ).context("upsert_daily: insert failed")?; + + if let Some(url) = cache_url { + conn.execute( + "INSERT INTO daily_cache_stats (day, cache_url, count, total_ms) + VALUES (?1, ?2, 1, ?3) + ON CONFLICT (day, cache_url) DO UPDATE SET + count = count + 1, + total_ms = total_ms + excluded.total_ms", + rusqlite::params![day, url, duration_ms], + ).context("upsert_daily: cache insert failed")?; + } + Ok(()) +} + +// Rebuild today's in-memory stats from the events table for today's UTC day. +// Runs once at daemon startup — scans only today's rows, which is fast. +fn rebuild_today(conn: &Connection, today_day: i64) -> Result { + assert!(today_day > 0); + let today_ts = today_day * 86400; + let tomorrow_ts = today_ts + 86400; + let mut t = TodayStats::new(today_day); + + let (bc, bms) = conn.query_row( + "SELECT COUNT(*), COALESCE(SUM(duration_ms),0) + FROM events INDEXED BY idx_events_type_start + WHERE event_type = 105 AND start_time >= ?1 AND start_time < ?2", + [today_ts, tomorrow_ts], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), + ).context("rebuild_today: build query failed")?; + t.build_count = bc; t.build_total_ms = bms; + + let (sc, sms) = conn.query_row( + "SELECT COUNT(*), COALESCE(SUM(duration_ms),0) + FROM events INDEXED BY idx_events_type_start + WHERE event_type = 108 AND start_time >= ?1 AND start_time < ?2", + [today_ts, tomorrow_ts], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), + ).context("rebuild_today: subst query failed")?; + t.subst_count = sc; t.subst_total_ms = sms; + + let (dc, db, dms) = conn.query_row( + "SELECT COUNT(*), COALESCE(SUM(total_bytes),0), COALESCE(SUM(duration_ms),0) + FROM events INDEXED BY idx_events_type_start + WHERE event_type = 101 AND start_time >= ?1 AND start_time < ?2", + [today_ts, tomorrow_ts], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?)), + ).context("rebuild_today: download query failed")?; + t.download_count = dc; t.download_bytes = db; t.download_ms = dms; + + let mut stmt = conn.prepare( + "SELECT cache_url, SUM(duration_ms), COUNT(*) + FROM events INDEXED BY idx_events_type_start + WHERE event_type = 108 AND cache_url IS NOT NULL + AND start_time >= ?1 AND start_time < ?2 + GROUP BY cache_url", + ).context("rebuild_today: cache query failed")?; + for row in stmt.query_map([today_ts, tomorrow_ts], |r| { + Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?)) + })?.filter_map(|r| r.ok()) { + let (url, ms, cnt) = row; + t.cache.insert(url, (ms, cnt)); + } + + assert!(t.build_count >= 0); + assert!(t.subst_count >= 0); + assert!(t.download_count >= 0); + info!(today_day, build_count = t.build_count, subst_count = t.subst_count, "Rebuilt today stats from events"); + Ok(t) +} + pub async fn run_daemon( socket_path: PathBuf, db: Arc, retain_days: Option, ) -> Result<()> { + let today_day = Utc::now().timestamp() / 86400; + let today = { + let conn = db.reader.lock().unwrap(); + rebuild_today(&conn, today_day)? + }; + let state = Arc::new(Mutex::new(State { active_activities: HashMap::new(), + today, })); if socket_path.exists() { @@ -381,9 +627,10 @@ async fn handle_connection(mut stream: UnixStream, state: Arc>, db: match serde_json::from_str::(line.trim()) { Ok(SocketMessage::Command(ClientCommand::GetStats { since, drv, sort, limit, group })) => { + let today_snap = state.lock().unwrap().today.snapshot(); let db = Arc::clone(&db); let stats = tokio::task::spawn_blocking(move || { - collect_stats(&db.reader, since, drv.as_deref(), sort, limit, group) + collect_stats(&db.reader, since, drv.as_deref(), sort, limit, group, Some(today_snap)) }).await??; writer.write_all((serde_json::to_string(&stats)? + "\n").as_bytes()).await?; break; @@ -489,20 +736,83 @@ fn process_event(event: NixEvent, state: &Arc>, db: &Arc None, }; + let event_day = act.start_time.timestamp() / 86400; + let real_today = end_time.timestamp() / 86400; + assert!(event_day > 0); + assert!(real_today >= event_day, "end_time must be >= start_time"); + + // Check for day rollover: real calendar day has advanced past our memory snapshot. + let old_today = if s.today.day > 0 && s.today.day < real_today { + Some(std::mem::replace(&mut s.today, TodayStats::new(real_today))) + } else { + if s.today.day == 0 { s.today.day = real_today; } + None + }; + + // Update in-memory today only for events whose start_time is today. + // Past-day events (long-running builds that spanned midnight) are written + // directly to daily_stats after the events insert. + let is_today = event_day == real_today; + if is_today { + match act_type { + ActivityType::Build => { + s.today.build_count += 1; + s.today.build_total_ms += duration_ms; + } + ActivityType::Substitute => { + s.today.subst_count += 1; + s.today.subst_total_ms += duration_ms; + if let Some(ref url) = cache_url { + let e = s.today.cache.entry(url.clone()).or_insert((0, 0)); + e.0 += duration_ms; e.1 += 1; + } + } + ActivityType::FileTransfer => { + s.today.download_count += 1; + s.today.download_bytes += act.total_bytes as i64; + s.today.download_ms += duration_ms; + } + _ => {} + } + } + drop(s); - db.writer.lock().unwrap().execute( - "INSERT INTO events (nix_id, parent_id, event_type, text, drv_path, cache_url, start_time, end_time, duration_ms, total_bytes) - VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", - rusqlite::params![ - act.id as i64, act.parent_id as i64, act.event_type as i64, - act.text, drv_path, cache_url, - act.start_time.timestamp(), end_time.timestamp(), - duration_ms, act.total_bytes as i64, - ], - ).context("Failed to insert event")?; + + { + let conn = db.writer.lock().unwrap(); + conn.execute( + "INSERT INTO events (nix_id, parent_id, event_type, text, drv_path, cache_url, start_time, end_time, duration_ms, total_bytes) + VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)", + rusqlite::params![ + act.id as i64, act.parent_id as i64, act.event_type as i64, + act.text, drv_path, cache_url.as_deref(), + act.start_time.timestamp(), end_time.timestamp(), + duration_ms, act.total_bytes as i64, + ], + ).context("Failed to insert event")?; + + if let Some(ref old) = old_today { + flush_today(&conn, old).context("Failed to flush old day to daily_stats")?; + } + if !is_today { + upsert_daily(&conn, event_day, act.event_type as i64, + duration_ms, act.total_bytes as i64, cache_url.as_deref()) + .context("Failed to upsert past-day event to daily_stats")?; + } + } } } } Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn schema_hash_value() { + println!("SCHEMA hash = {:#010x}", schema_hash(SCHEMA.as_bytes())); + } +} diff --git a/src/main.rs b/src/main.rs index 02c2318..0d54c73 100644 --- a/src/main.rs +++ b/src/main.rs @@ -326,7 +326,7 @@ fn query_direct( let trend = collect_trend(&conn, since, bucket_size, drv)?; display_trend_output(&trend, output); } else { - let s = nod::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, None)?; display_stats(s); } diff --git a/src/stats.rs b/src/stats.rs index 6323760..27a14f1 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -2,6 +2,7 @@ use anyhow::{Context, Result}; use chrono::prelude::*; use rusqlite::Connection; use serde::{Deserialize, Serialize}; +use std::collections::HashMap; use std::sync::Mutex; #[derive(Debug, Clone, Default, Serialize, Deserialize, clap::ValueEnum)] @@ -40,24 +41,103 @@ pub struct CacheStat { pub count: i64, } -// Three separate queries — one per event type — so the composite index -// (event_type, start_time) can be used for each rather than doing a full -// table scan with per-aggregate FILTER conditions. -pub fn collect_stats( - db: &Mutex, +// Snapshot of the current day's in-memory aggregates, passed from the daemon +// to collect_stats so the summary path never needs to scan the events table. +// cache: substituter URL → (total_ms, count) +#[derive(Clone, Default)] +pub struct TodaySummary { + pub day: i64, + pub build_count: i64, + pub build_total_ms: i64, + pub subst_count: i64, + pub subst_total_ms: i64, + pub download_count: i64, + pub download_bytes: i64, + pub download_ms: i64, + pub cache: HashMap, +} + +// Fast path: query daily_stats for closed days, merge today's in-memory snapshot. +// Returns (build_count, build_ms, subst_count, subst_ms, dl_bytes, dl_ms, cache_latency). +fn summary_from_cache( + conn: &Connection, since: Option, - drv: Option<&str>, - sort: SortField, - limit: u32, - group: bool, -) -> Result { - assert!(limit > 0, "limit must be > 0"); + today: &TodaySummary, +) -> Result<(i64, i64, i64, i64, i64, i64, Vec)> { + assert!(today.day > 0); + let since_day: Option = since.map(|ts| ts / 86400); + let today_in_range = since_day.map_or(true, |sd| today.day >= sd); + + let (hbc, hbms) = conn.query_row( + "SELECT COALESCE(SUM(count),0), COALESCE(SUM(total_ms),0) + FROM daily_stats WHERE event_type = 105 + AND (?1 IS NULL OR day >= ?1) AND day < ?2", + rusqlite::params![since_day, today.day], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), + ).context("Failed to query build summary from daily_stats")?; - let conn = db.lock().unwrap(); + let (hsc, hsms) = conn.query_row( + "SELECT COALESCE(SUM(count),0), COALESCE(SUM(total_ms),0) + FROM daily_stats WHERE event_type = 108 + AND (?1 IS NULL OR day >= ?1) AND day < ?2", + rusqlite::params![since_day, today.day], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), + ).context("Failed to query subst summary from daily_stats")?; + + let (hdb, hdms) = conn.query_row( + "SELECT COALESCE(SUM(total_bytes),0), COALESCE(SUM(total_ms),0) + FROM daily_stats WHERE event_type = 101 + AND (?1 IS NULL OR day >= ?1) AND day < ?2", + rusqlite::params![since_day, today.day], + |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), + ).context("Failed to query download summary from daily_stats")?; + + let build_count = hbc + if today_in_range { today.build_count } else { 0 }; + let build_total_ms = hbms + if today_in_range { today.build_total_ms } else { 0 }; + let subst_count = hsc + if today_in_range { today.subst_count } else { 0 }; + let subst_total_ms = hsms + if today_in_range { today.subst_total_ms } else { 0 }; + let download_bytes = hdb + if today_in_range { today.download_bytes } else { 0 }; + let download_ms = hdms + if today_in_range { today.download_ms } else { 0 }; + // Cache latency: closed days from daily_cache_stats, today from memory. + let mut cache_map: HashMap = HashMap::new(); + let mut stmt = conn.prepare( + "SELECT cache_url, SUM(total_ms), SUM(count) + FROM daily_cache_stats + WHERE (?1 IS NULL OR day >= ?1) AND day < ?2 + GROUP BY cache_url", + ).context("Failed to prepare cache latency query")?; + for row in stmt.query_map(rusqlite::params![since_day, today.day], |r| { + Ok((r.get::<_, String>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?)) + })?.filter_map(|r| r.ok()) { + let (url, ms, cnt) = row; + let e = cache_map.entry(url).or_insert((0, 0)); + e.0 += ms; e.1 += cnt; + } + if today_in_range { + for (url, &(ms, cnt)) in &today.cache { + let e = cache_map.entry(url.clone()).or_insert((0, 0)); + e.0 += ms; e.1 += cnt; + } + } + let mut cache_latency: Vec = cache_map.into_iter() + .filter(|(_, (_, cnt))| *cnt > 0) + .map(|(url, (ms, cnt))| CacheStat { cache_url: url, avg_ms: ms as f64 / cnt as f64, count: cnt }) + .collect(); + cache_latency.sort_by(|a, b| b.avg_ms.partial_cmp(&a.avg_ms).unwrap_or(std::cmp::Ordering::Equal)); + + Ok((build_count, build_total_ms, subst_count, subst_total_ms, download_bytes, download_ms, cache_latency)) +} + +// Slow path: full events table scan (used in direct mode or when drv filter is active). +fn summary_from_events( + conn: &Connection, + since: Option, + drv: Option<&str>, +) -> Result<(i64, i64, i64, i64, i64, i64, Vec)> { let (build_count, build_total_ms) = conn.query_row( "SELECT COUNT(*), COALESCE(SUM(duration_ms), 0) - FROM events WHERE event_type = 105 + FROM events INDEXED BY idx_events_type_start WHERE event_type = 105 AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%')", rusqlite::params![since, drv], @@ -66,7 +146,7 @@ pub fn collect_stats( let (subst_count, subst_total_ms) = conn.query_row( "SELECT COUNT(*), COALESCE(SUM(duration_ms), 0) - FROM events WHERE event_type = 108 + FROM events INDEXED BY idx_events_type_start WHERE event_type = 108 AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%')", rusqlite::params![since, drv], @@ -75,12 +155,50 @@ pub fn collect_stats( let (download_bytes, download_ms) = conn.query_row( "SELECT COALESCE(SUM(total_bytes), 0), COALESCE(SUM(duration_ms), 0) - FROM events WHERE event_type = 101 + FROM events INDEXED BY idx_events_type_start WHERE event_type = 101 AND (?1 IS NULL OR start_time >= ?1)", rusqlite::params![since], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?)), ).context("Failed to query download summary")?; + let mut stmt = conn.prepare( + "SELECT cache_url, AVG(duration_ms), COUNT(*) + FROM events INDEXED BY idx_events_type_start + WHERE event_type = 108 AND cache_url IS NOT NULL + AND (?1 IS NULL OR start_time >= ?1) + GROUP BY cache_url ORDER BY AVG(duration_ms) DESC", + ).context("Failed to prepare cache latency query")?; + let cache_latency: Vec = stmt.query_map(rusqlite::params![since], |r| { + Ok(CacheStat { cache_url: r.get(0)?, avg_ms: r.get(1)?, count: r.get(2)? }) + })?.filter_map(|r| r.ok()).collect(); + + Ok((build_count, build_total_ms, subst_count, subst_total_ms, download_bytes, download_ms, cache_latency)) +} + +pub fn collect_stats( + db: &Mutex, + since: Option, + drv: Option<&str>, + sort: SortField, + limit: u32, + group: bool, + today: Option, +) -> Result { + assert!(limit > 0, "limit must be > 0"); + + let conn = db.lock().unwrap(); + + // Fast path: daily_stats + today memory — O(days) instead of O(events). + // Falls back to events scan when a drv filter is active (drv_path not in daily_stats) + // or when running in direct mode (no daemon, today is None). + let (build_count, build_total_ms, subst_count, subst_total_ms, + download_bytes, download_ms, cache_latency) = + if let (Some(t), None) = (today.as_ref(), drv) { + summary_from_cache(&conn, since, t)? + } else { + summary_from_events(&conn, since, drv)? + }; + assert!(build_count >= 0); assert!(subst_count >= 0); assert!(download_bytes >= 0); @@ -93,7 +211,7 @@ pub fn collect_stats( }; let sql = format!( "SELECT drv_path, CAST(ROUND(AVG(duration_ms)) AS INTEGER) as avg_ms, COUNT(*) as cnt - FROM events + FROM events INDEXED BY idx_events_type_start WHERE event_type = 105 AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%') @@ -112,7 +230,7 @@ pub fn collect_stats( }; let sql = format!( "SELECT duration_ms, drv_path, text - FROM events + FROM events INDEXED BY idx_events_type_start WHERE event_type = 105 AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%') @@ -125,19 +243,6 @@ pub fn collect_stats( })?.filter_map(|r| r.ok()).collect() }; - // 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) - GROUP BY cache_url ORDER BY AVG(duration_ms) DESC", - ).context("Failed to prepare cache latency query")?; - let cache_latency: Vec = stmt.query_map(rusqlite::params![since], |r| { - Ok(CacheStat { cache_url: r.get(0)?, avg_ms: r.get(1)?, count: r.get(2)? }) - })?.filter_map(|r| r.ok()).collect(); - Ok(Stats { build_count, build_total_ms, subst_count, subst_total_ms, download_bytes, download_ms, slowest_builds, cache_latency }) } @@ -366,9 +471,11 @@ pub fn collect_trend( let drv_ref = drv.as_deref(); // 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. let mut stmt = conn.prepare( "SELECT start_time, event_type, duration_ms, total_bytes - FROM events + FROM events INDEXED BY idx_events_start_cover WHERE event_type IN (101, 105, 108) AND (?1 IS NULL OR start_time >= ?1) AND (?2 IS NULL OR drv_path LIKE '%' || ?2 || '%') -- 2.51.2