diff --git a/src/events.rs b/src/events.rs new file mode 100644 index 0000000..51fb2ba --- /dev/null +++ b/src/events.rs @@ -0,0 +1,404 @@ +//! Event journal: every action taken by agent-manager is appended as one +//! JSON line per event in `events-YYYYMM.jsonl` files inside the state +//! directory. The journal is the source of truth for stats, sessions and +//! project aggregation (ROADMAP.md, axe 1). +//! +//! Design rules: +//! * append-only, one compact JSON object per line; +//! * monthly rotation (file name derived from the event timestamp); +//! * corrupt lines are skipped on read (never fatal); +//! * no long locks: plain append with a short retry loop for Windows +//! file locks. + +use serde::{Deserialize, Serialize}; +use std::collections::BTreeMap; +use std::io::Write; +use std::path::{Path, PathBuf}; + +/// Kind of event, serialized in snake_case. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum EventKind { + Start, + Stop, + Run, + Restart, + Install, + Update, + Uninstall, + Doctor, + Config, + Repl, + Shell, +} + +/// One journal entry. Optional fields are omitted from the JSON line. +#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] +pub struct Event { + /// RFC 3339 timestamp (UTC), used to pick the monthly file. + pub ts: String, + pub kind: EventKind, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub agent: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub pid: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub cwd: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub project: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub session: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub reason: Option, + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub args: Vec, + /// Variable names only — values (secrets) never enter the journal. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub env_keys: Vec, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub exit_code: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub duration_s: Option, +} + +impl Event { + /// Event stamped with the current UTC time. + pub fn now(kind: EventKind) -> Self { + Self { + ts: crate::installers::now_rfc3339(), + kind, + agent: None, + pid: None, + cwd: None, + project: None, + session: None, + reason: None, + args: Vec::new(), + env_keys: Vec::new(), + exit_code: None, + duration_s: None, + } + } + + pub fn with_agent(mut self, agent: impl Into) -> Self { + self.agent = Some(agent.into()); + self + } + + pub fn with_pid(mut self, pid: u32) -> Self { + self.pid = Some(pid); + self + } + + pub fn with_cwd(mut self, cwd: impl Into) -> Self { + self.cwd = Some(cwd.into()); + self + } + + pub fn with_project(mut self, project: impl Into) -> Self { + self.project = Some(project.into()); + self + } + + pub fn with_session(mut self, session: impl Into) -> Self { + self.session = Some(session.into()); + self + } + + pub fn with_reason(mut self, reason: impl Into) -> Self { + self.reason = Some(reason.into()); + self + } + + pub fn with_args(mut self, args: Vec) -> Self { + self.args = args; + self + } + + pub fn with_env_keys(mut self, env_keys: Vec) -> Self { + self.env_keys = env_keys; + self + } + + pub fn with_exit_code(mut self, code: i32) -> Self { + self.exit_code = Some(code); + self + } + + pub fn with_duration(mut self, secs: u64) -> Self { + self.duration_s = Some(secs); + self + } +} + +/// Sorted list of the variable names passed to a process. Values are never +/// logged — only the keys — so secrets stay out of the journal. +pub fn env_keys(env: &BTreeMap) -> Vec { + let mut keys: Vec = env.keys().cloned().collect(); + keys.sort(); + keys +} + +/// Month of an RFC 3339 timestamp: "2026-08-16T13:49:11Z" -> "202608". +pub fn month_of(ts: &str) -> String { + ts.chars() + .take(7) + .filter(|c| c.is_ascii_digit()) + .collect() +} + +/// Journal file name for a month: "202608" -> "events-202608.jsonl". +pub fn file_name_for_month(month: &str) -> String { + format!("events-{month}.jsonl") +} + +/// Month part of a journal file name, or None when the name does not look +/// like a journal file ("events-202608.jsonl" -> "202608"). +pub fn month_from_file_name(name: &str) -> Option { + let rest = name.strip_prefix("events-")?.strip_suffix(".jsonl")?; + (rest.len() == 6 && rest.chars().all(|c| c.is_ascii_digit())).then(|| rest.to_string()) +} + +/// Sorted journal files of a directory, oldest first. +pub fn journal_files(dir: &Path) -> Vec { + let mut files: Vec = std::fs::read_dir(dir) + .map(|rd| { + rd.flatten() + .map(|e| e.path()) + .filter(|p| p.is_file()) + .filter(|p| { + p.file_name() + .and_then(|n| n.to_str()) + .and_then(month_from_file_name) + .is_some() + }) + .collect() + }) + .unwrap_or_default(); + files.sort(); + files +} + +/// Parse one journal file into events; corrupt lines are skipped silently. +pub fn parse_file(path: &Path) -> Vec { + let Ok(text) = std::fs::read_to_string(path) else { + return Vec::new(); + }; + text.lines() + .filter(|l| !l.trim().is_empty()) + .filter_map(|l| serde_json::from_str::(l).ok()) + .collect() +} + +/// Read events from a directory in chronological order, limited to the most +/// recent `limit` events when limit > 0. +pub fn read_events(dir: &Path, limit: usize) -> Vec { + let mut out: Vec = Vec::new(); + for f in journal_files(dir).iter().rev() { + let mut evs = parse_file(f); + evs.reverse(); // newest first within one file + out.extend(evs); + if limit > 0 && out.len() >= limit { + break; + } + } + if limit > 0 { + out.truncate(limit); + } + out.reverse(); // chronological order + out +} + +/// Append-only journal living in a directory (the state directory). +pub struct EventLog { + dir: PathBuf, + /// Cached (month, path) pair to avoid re-computing the file per event. + current: Option<(String, PathBuf)>, +} + +impl EventLog { + pub fn new(dir: impl Into) -> Self { + Self { + dir: dir.into(), + current: None, + } + } + + /// Directory the journal writes into. + pub fn dir(&self) -> &Path { + &self.dir + } + + /// Path of the journal file for the current month. + pub fn current_file(&mut self) -> PathBuf { + let now = crate::installers::now_rfc3339(); + self.file_for(&month_of(&now)) + } + + fn file_for(&mut self, month: &str) -> PathBuf { + if let Some((m, p)) = &self.current { + if m == month { + return p.clone(); + } + } + let p = self.dir.join(file_name_for_month(month)); + self.current = Some((month.to_string(), p.clone())); + p + } + + /// Append one event as a single JSON line; returns the file written to. + /// Retries briefly when the file is locked (Windows). + pub fn append(&mut self, event: &Event) -> anyhow::Result { + let month = month_of(&event.ts); + let path = self.file_for(&month); + if let Some(parent) = path.parent() { + let _ = std::fs::create_dir_all(parent); + } + let line = serde_json::to_string(event)?; + let mut last_err: Option = None; + for attempt in 0..5 { + match std::fs::OpenOptions::new() + .create(true) + .append(true) + .open(&path) + { + Ok(mut f) => { + f.write_all(line.as_bytes())?; + f.write_all(b"\n")?; + return Ok(path); + } + Err(e) => { + last_err = Some(e); + if attempt < 4 { + std::thread::sleep(std::time::Duration::from_millis(20 * (attempt + 1))); + } + } + } + } + Err(anyhow::anyhow!( + "cannot append to journal {}: {}", + path.display(), + last_err + .map(|e| e.to_string()) + .unwrap_or_else(|| "unknown error".to_string()) + )) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn event_roundtrip() { + let ev = Event::now(EventKind::Start) + .with_agent("claude-code") + .with_pid(42) + .with_session("s1") + .with_env_keys(vec!["OPENAI_KEY".to_string()]) + .with_args(vec!["--dangerously-skip-permissions".to_string()]); + let json = serde_json::to_string(&ev).unwrap(); + assert!(!json.contains('\n')); + let back: Event = serde_json::from_str(&json).unwrap(); + assert_eq!(back, ev); + assert_eq!(back.agent.as_deref(), Some("claude-code")); + assert_eq!(back.pid, Some(42)); + // Absent optional fields are omitted from the line. + assert!(!json.contains("duration_s")); + assert!(!json.contains("reason")); + } + + #[test] + fn env_keys_are_sorted_and_never_values() { + let mut env = BTreeMap::new(); + env.insert("ZEBRA".to_string(), "secret-z".to_string()); + env.insert("ALPHA".to_string(), "secret-a".to_string()); + let keys = env_keys(&env); + assert_eq!(keys, vec!["ALPHA".to_string(), "ZEBRA".to_string()]); + } + + #[test] + fn month_helpers() { + assert_eq!(month_of("2026-08-16T13:49:11Z"), "202608"); + assert_eq!(file_name_for_month("202608"), "events-202608.jsonl"); + assert_eq!( + month_from_file_name("events-202608.jsonl"), + Some("202608".to_string()) + ); + assert_eq!(month_from_file_name("state.json"), None); + assert_eq!(month_from_file_name("events-202608.jsonl.bak"), None); + } + + #[test] + fn append_roundtrip() { + let dir = tempfile::tempdir().unwrap(); + let mut log = EventLog::new(dir.path()); + let ev = Event::now(EventKind::Run) + .with_agent("aider") + .with_exit_code(0) + .with_duration(48); + let path = log.append(&ev).unwrap(); + assert_eq!( + path.file_name().and_then(|n| n.to_str()), + Some(file_name_for_month(&month_of(&ev.ts)).as_str()) + ); + let read = parse_file(&path); + assert_eq!(read.len(), 1); + assert_eq!(read[0], ev); + } + + #[test] + fn corrupt_lines_are_skipped() { + let dir = tempfile::tempdir().unwrap(); + let mut log = EventLog::new(dir.path()); + let ev = Event::now(EventKind::Doctor); + log.append(&ev).unwrap(); + // Corrupt the file with garbage lines around a valid one. + let path = log.current_file(); + let mut text = std::fs::read_to_string(&path).unwrap(); + text.push_str("ceci n'est pas du json +"); + text.push_str(" +"); + std::fs::write(&path, text).unwrap(); + let read = parse_file(&path); + assert_eq!(read.len(), 1); + assert_eq!(read[0].kind, EventKind::Doctor); + } + + #[test] + fn read_events_spans_months_and_respects_limit() { + let dir = tempfile::tempdir().unwrap(); + // Two monthly files, one unrelated file. + std::fs::write( + dir.path().join("events-202607.jsonl"), + r#"{"ts":"2026-07-10T10:00:00Z","kind":"install","agent":"a"}"#, + ) + .unwrap(); + std::fs::write( + dir.path().join("events-202608.jsonl"), + concat!( + r#"{"ts":"2026-08-01T10:00:00Z","kind":"run","agent":"b"}"#, + "\n", + r#"{"ts":"2026-08-02T10:00:00Z","kind":"stop","agent":"b"}"#, + "\n", + ), + ) + .unwrap(); + std::fs::write(dir.path().join("state.json"), "{}").unwrap(); + + let files = journal_files(dir.path()); + assert_eq!(files.len(), 2); + assert!(files[0].to_string_lossy().ends_with("events-202607.jsonl")); + + let all = read_events(dir.path(), 0); + assert_eq!(all.len(), 3); + assert_eq!(all[0].ts, "2026-07-10T10:00:00Z"); + assert_eq!(all[2].ts, "2026-08-02T10:00:00Z"); + + let limited = read_events(dir.path(), 2); + assert_eq!(limited.len(), 2); + assert_eq!(limited[0].ts, "2026-08-01T10:00:00Z"); + assert_eq!(limited[1].ts, "2026-08-02T10:00:00Z"); + } +} diff --git a/src/lib.rs b/src/lib.rs index ab39337..8fc6879 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -19,6 +19,7 @@ pub mod config; pub mod deps; pub mod doctor; pub mod download; +pub mod events; pub mod help; pub mod installers; pub mod nav; diff --git a/tests/events_test.rs b/tests/events_test.rs new file mode 100644 index 0000000..6fb9fc0 --- /dev/null +++ b/tests/events_test.rs @@ -0,0 +1,41 @@ +//! Integration test for the event journal: parses the replayable +//! fixture shipped with the repository and exercises the append path. + +use agent_manager::events::{self, Event, EventKind}; +use std::path::PathBuf; + +fn fixture() -> PathBuf { + PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("tests/fixtures/events-sample.jsonl") +} + +#[test] +fn fixture_parses_valid_lines_and_skips_corrupt_ones() { + let events = events::parse_file(&fixture()); + assert_eq!(events.len(), 2, "the corrupt line must be skipped"); + assert_eq!(events[0].kind, EventKind::Start); + assert_eq!(events[0].agent.as_deref(), Some("claude-code")); + assert_eq!(events[0].pid, Some(1234)); + assert_eq!(events[0].env_keys, vec!["OPENAI_KEY".to_string()]); + assert_eq!(events[1].kind, EventKind::Stop); + assert_eq!(events[1].duration_s, Some(1364)); + assert_eq!(events[1].exit_code, Some(0)); + assert_eq!(events[1].reason.as_deref(), Some("signal")); +} + +#[test] +fn append_then_read_back_through_read_events() { + let dir = tempfile::tempdir().unwrap(); + let mut log = events::EventLog::new(dir.path()); + log.append( + &Event::now(EventKind::Run) + .with_agent("jcode") + .with_cwd("/tmp/demo") + .with_exit_code(0), + ) + .unwrap(); + let all: Vec = events::read_events(dir.path(), 10); + assert_eq!(all.len(), 1); + assert_eq!(all[0].agent.as_deref(), Some("jcode")); + assert_eq!(all[0].cwd.as_deref(), Some("/tmp/demo")); + assert_eq!(all[0].exit_code, Some(0)); +} diff --git a/tests/fixtures/events-sample.jsonl b/tests/fixtures/events-sample.jsonl new file mode 100644 index 0000000..5aff512 --- /dev/null +++ b/tests/fixtures/events-sample.jsonl @@ -0,0 +1,3 @@ +{"ts":"2026-08-15T14:39:26Z","kind":"start","agent":"claude-code","pid":1234,"cwd":"/home/bruno/proj/api","session":"20260815_143926_a1b2c3","env_keys":["OPENAI_KEY"]} +cette ligne est corrompue et doit etre ignoree par parse_file +{"ts":"2026-08-15T15:02:10Z","kind":"stop","agent":"claude-code","pid":1234,"reason":"signal","duration_s":1364,"exit_code":0}