diff --git a/src/daemon.rs b/src/daemon.rs index 1635f00..008b319 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -12,8 +12,12 @@ use tracing::{error, info}; use crate::stats::collect_stats; -// Schema version - increment when the schema changes incompatibly. -const SCHEMA_VERSION: u32 = 1; +const fn schema_hash(s: &[u8]) -> u32 { + let mut h: u32 = 2166136261; + let mut i = 0; + while i < s.len() { h ^= s[i] as u32; h = h.wrapping_mul(16777619); i += 1; } + h +} const SCHEMA: &str = " CREATE TABLE IF NOT EXISTS events ( @@ -31,6 +35,16 @@ CREATE TABLE IF NOT EXISTS events ( ); "; +// Append a new hash entry when the schema changes - SCHEMA_VERSION auto-increments. +const SCHEMA_HASHES: &[u32] = &[ + 0x9bc94a70, // v1 +]; +const SCHEMA_VERSION: u32 = SCHEMA_HASHES.len() as u32; +const _: () = assert!( + schema_hash(SCHEMA.as_bytes()) == SCHEMA_HASHES[SCHEMA_VERSION as usize - 1], + "schema changed - append new hash to SCHEMA_HASHES" +); + #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[repr(u64)] pub enum ActivityType { @@ -114,19 +128,46 @@ impl fmt::Display for ResultType { } } -#[derive(Debug, Deserialize, Serialize)] -struct NixEvent { - action: String, - id: u64, - // Raw number: ActivityType for "start", ResultType for "result", absent on "stop" - #[serde(default, rename = "type")] - event_type: u64, - #[serde(default)] - text: String, - #[serde(default)] - fields: Vec, - #[serde(default)] - parent: u64, +#[derive(Debug, Deserialize)] +#[serde(tag = "action", rename_all = "snake_case")] +enum NixEvent { + Start { + id: u64, + #[serde(rename = "type", default)] + event_type: u64, + #[serde(default)] + text: String, + #[serde(default)] + fields: Vec, + #[serde(default)] + parent: u64, + }, + Stop { + id: u64, + }, + Result { + id: u64, + #[serde(rename = "type", default)] + event_type: u64, + #[serde(default)] + fields: Vec, + }, +} + +#[derive(Debug, Deserialize)] +#[serde(tag = "action", rename_all = "snake_case")] +enum ClientCommand { + GetStats { since: Option }, + Clean, +} + +// Untagged outer enum: serde tries ClientCommand first, then NixEvent. +// The action values never overlap so routing is unambiguous. +#[derive(Debug, Deserialize)] +#[serde(untagged)] +enum SocketMessage { + Command(ClientCommand), + Event(NixEvent), } struct Activity { @@ -211,32 +252,32 @@ async fn handle_connection(mut stream: UnixStream, state: Arc>) -> loop { line.clear(); if reader.read_line(&mut line).await? == 0 { break; } - let cmd = line.trim(); - - if cmd.starts_with("get_stats") { - let since = cmd.split_whitespace().nth(1) - .unwrap_or("1970-01-01T00:00:00+00:00") - .to_string(); - let db = state.lock().unwrap().db.clone(); - let stats = tokio::task::spawn_blocking(move || collect_stats(&db, &since)) - .await??; - writer.write_all((serde_json::to_string(&stats)? + "\n").as_bytes()).await?; - break; - } - if cmd == "clean" { - let db = state.lock().unwrap().db.clone(); - tokio::task::spawn_blocking(move || -> Result<()> { - let conn = db.lock().unwrap(); - conn.execute_batch("DELETE FROM events; VACUUM; PRAGMA wal_checkpoint(TRUNCATE);")?; - Ok(()) - }).await??; - writer.write_all(b"ok\n").await?; - info!("Database cleared via socket command"); - break; - } - if let Ok(event) = serde_json::from_str::(cmd) { - process_event(event, &state)?; + match serde_json::from_str::(line.trim()) { + Ok(SocketMessage::Command(ClientCommand::GetStats { since })) => { + let db = state.lock().unwrap().db.clone(); + let stats = tokio::task::spawn_blocking(move || collect_stats(&db, since)) + .await??; + writer.write_all((serde_json::to_string(&stats)? + "\n").as_bytes()).await?; + break; + } + Ok(SocketMessage::Command(ClientCommand::Clean)) => { + let db = state.lock().unwrap().db.clone(); + tokio::task::spawn_blocking(move || -> Result<()> { + let conn = db.lock().unwrap(); + conn.execute_batch("DELETE FROM events; VACUUM; PRAGMA wal_checkpoint(TRUNCATE);")?; + Ok(()) + }).await??; + writer.write_all(b"ok\n").await?; + info!("Database cleared via socket command"); + break; + } + Ok(SocketMessage::Event(event)) => { + if let Err(e) = process_event(event, &state) { + error!("Failed to process event: {}", e); + } + } + Err(e) => error!("Invalid message: {}", e), } } Ok(()) @@ -245,56 +286,56 @@ async fn handle_connection(mut stream: UnixStream, state: Arc>) -> fn process_event(event: NixEvent, state: &Arc>) -> Result<()> { let mut s = state.lock().unwrap(); - match event.action.as_str() { - "start" => { - let act_type = ActivityType::from(event.event_type); - let text = if !event.text.is_empty() { - event.text.clone() + match event { + NixEvent::Start { id, event_type, text, fields, parent } => { + let act_type = ActivityType::from(event_type); + let text = if !text.is_empty() { + text } else { - event.fields.get(0).and_then(|v| v.as_str()).unwrap_or("").to_string() + fields.get(0).and_then(|v| v.as_str()).unwrap_or("").to_string() }; info!( - id = event.id, - parent = event.parent, + id, + parent, act_type = %act_type, text = %text, - fields = ?event.fields, + fields = ?fields, "start" ); - s.active_activities.insert(event.id, Activity { - id: event.id, - parent_id: event.parent, - event_type: event.event_type, + s.active_activities.insert(id, Activity { + id, + parent_id: parent, + event_type, text, start_time: Utc::now(), - fields: event.fields, + fields, total_bytes: 0, }); } - "result" => { - let res_type = ResultType::from(event.event_type); + NixEvent::Result { id, event_type, fields } => { + let res_type = ResultType::from(event_type); if res_type != ResultType::BuildLogLine && res_type != ResultType::Progress && res_type != ResultType::SetExpected { info!( - id = event.id, + id, res_type = %res_type, - fields = ?event.fields, + fields = ?fields, "result" ); } - if let Some(act) = s.active_activities.get_mut(&event.id) { + if let Some(act) = s.active_activities.get_mut(&id) { if res_type == ResultType::Progress { - if let Some(total) = event.fields.get(1).and_then(|v| v.as_u64()) { + if let Some(total) = fields.get(1).and_then(|v| v.as_u64()) { if total > 0 { act.total_bytes = total; } } } } } - "stop" => { - if let Some(act) = s.active_activities.remove(&event.id) { + NixEvent::Stop { id } => { + if let Some(act) = s.active_activities.remove(&id) { let act_type = ActivityType::from(act.event_type); let end_time = Utc::now(); let duration_ms = end_time.signed_duration_since(act.start_time).num_milliseconds(); @@ -338,9 +379,6 @@ fn process_event(event: NixEvent, state: &Arc>) -> Result<()> { ).context("Failed to insert event")?; } } - _ => { - info!(action = %event.action, id = event.id, "unknown action"); - } } Ok(()) } diff --git a/src/main.rs b/src/main.rs index 8cd2fbf..44a02f4 100644 --- a/src/main.rs +++ b/src/main.rs @@ -103,20 +103,22 @@ async fn main() -> Result<()> { .unwrap_or_else(|| PathBuf::from("/tmp/nod.sock")) }); - // Compute the since timestamp. Flags are additive (e.g. -y 1 -d 3 = 1 year + 3 days ago). - let mut since = Utc::now(); - if let Some(y) = years { since = since - chrono::Months::new(y * 12); } - if let Some(m) = months { since = since - chrono::Months::new(m); } - if let Some(d) = days { since = since - chrono::Duration::days(d as i64); } - let has_filter = days.is_some() || months.is_some() || years.is_some(); - // Epoch as the "no filter" sentinel — a WHERE start_time >= epoch matches everything. - let since_str = if has_filter { since.to_rfc3339() } else { "1970-01-01T00:00:00+00:00".to_string() }; + 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 + }; let mut stream = UnixStream::connect(&socket_path) .await .with_context(|| format!("Failed to connect to daemon at {}", socket_path.display()))?; - stream.write_all(format!("get_stats {}\n", since_str).as_bytes()).await?; + let cmd = serde_json::json!({"action": "get_stats", "since": since}); + stream.write_all((cmd.to_string() + "\n").as_bytes()).await?; let mut reader = BufReader::new(stream); let mut line = String::new(); @@ -136,7 +138,7 @@ async fn main() -> Result<()> { .await .with_context(|| format!("Failed to connect to daemon at {}", socket_path.display()))?; - stream.write_all(b"clean\n").await?; + stream.write_all(b"{\"action\":\"clean\"}\n").await?; let mut reader = BufReader::new(stream); let mut line = String::new(); diff --git a/src/stats.rs b/src/stats.rs index 9c417cf..3fe058c 100644 --- a/src/stats.rs +++ b/src/stats.rs @@ -30,9 +30,16 @@ pub struct CacheStat { pub count: i64, } -pub fn collect_stats(db: &Arc>, since: &str) -> Result { +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. + 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 (build_count, build_total_ms, subst_count, subst_total_ms, download_bytes, download_ms) = conn.query_row( "SELECT @@ -42,18 +49,18 @@ pub fn collect_stats(db: &Arc>, since: &str) -> Result COALESCE(SUM(duration_ms) FILTER (WHERE event_type = 108), 0), COALESCE(SUM(total_bytes) FILTER (WHERE event_type = 101), 0), COALESCE(SUM(duration_ms) FILTER (WHERE event_type = 101), 0) - FROM events WHERE start_time >= ?1", - [since], + FROM events WHERE (?1 IS NULL OR start_time >= ?1)", + rusqlite::params![p], |r| Ok((r.get::<_, i64>(0)?, r.get::<_, i64>(1)?, r.get::<_, i64>(2)?, r.get::<_, i64>(3)?, r.get::<_, i64>(4)?, r.get::<_, i64>(5)?)), ).context("Failed to query summary stats")?; let mut stmt = conn.prepare( "SELECT duration_ms, drv_path, text - FROM events WHERE event_type = 105 AND start_time >= ?1 + FROM events WHERE event_type = 105 AND (?1 IS NULL OR start_time >= ?1) ORDER BY duration_ms DESC LIMIT 10", ).context("Failed to prepare slowest builds query")?; - let slowest_builds: Vec = stmt.query_map([since], |r| { + let slowest_builds: Vec = stmt.query_map(rusqlite::params![p], |r| { Ok(SlowBuild { duration_ms: r.get(0)?, drv_path: r.get(1)?, @@ -65,10 +72,10 @@ pub fn collect_stats(db: &Arc>, since: &str) -> Result // per substituter, not just metadata query time (QueryPathInfo). 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 start_time >= ?1 + 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([since], |r| { + let cache_latency: Vec = stmt.query_map(rusqlite::params![p], |r| { Ok(CacheStat { cache_url: r.get(0)?, avg_ms: r.get(1)?,