events : journal d'evenements JSONL (issue #2) #17
+404
@@ -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<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub pid: Option<u32>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub cwd: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub project: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub session: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub reason: Option<String>,
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub args: Vec<String>,
|
||||
/// Variable names only — values (secrets) never enter the journal.
|
||||
#[serde(default, skip_serializing_if = "Vec::is_empty")]
|
||||
pub env_keys: Vec<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub exit_code: Option<i32>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub duration_s: Option<u64>,
|
||||
}
|
||||
|
||||
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<String>) -> 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<String>) -> Self {
|
||||
self.cwd = Some(cwd.into());
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_project(mut self, project: impl Into<String>) -> Self {
|
||||
self.project = Some(project.into());
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_session(mut self, session: impl Into<String>) -> Self {
|
||||
self.session = Some(session.into());
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_reason(mut self, reason: impl Into<String>) -> Self {
|
||||
self.reason = Some(reason.into());
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_args(mut self, args: Vec<String>) -> Self {
|
||||
self.args = args;
|
||||
self
|
||||
}
|
||||
|
||||
pub fn with_env_keys(mut self, env_keys: Vec<String>) -> 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<String, String>) -> Vec<String> {
|
||||
let mut keys: Vec<String> = 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<String> {
|
||||
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<PathBuf> {
|
||||
let mut files: Vec<PathBuf> = 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<Event> {
|
||||
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::<Event>(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<Event> {
|
||||
let mut out: Vec<Event> = 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<PathBuf>) -> 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<PathBuf> {
|
||||
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<std::io::Error> = 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");
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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<Event> = 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));
|
||||
}
|
||||
Vendored
+3
@@ -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}
|
||||
Reference in New Issue
Block a user