415 lines
13 KiB
Rust
415 lines
13 KiB
Rust
//! 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: Option<String>) -> Self {
|
|
self.cwd = cwd;
|
|
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 one event as a single JSON line to the journal file for the
|
|
/// event's month (created under `dir`). Returns the file written to.
|
|
/// Retries briefly when the file is locked (Windows).
|
|
pub fn append(dir: &Path, event: &Event) -> anyhow::Result<PathBuf> {
|
|
let path = dir.join(file_name_for_month(&month_of(&event.ts)));
|
|
append_to(&path, event)?;
|
|
Ok(path)
|
|
}
|
|
|
|
/// Append one event to a specific journal file.
|
|
pub fn append_to(path: &Path, event: &Event) -> anyhow::Result<()> {
|
|
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(());
|
|
}
|
|
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())
|
|
))
|
|
}
|
|
|
|
/// 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.
|
|
pub fn append(&mut self, event: &Event) -> anyhow::Result<PathBuf> {
|
|
let month = month_of(&event.ts);
|
|
let path = self.file_for(&month);
|
|
append_to(&path, event)?;
|
|
Ok(path)
|
|
}
|
|
}
|
|
|
|
#[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");
|
|
}
|
|
}
|