Files
agent-manager/src/events.rs
T

488 lines
16 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,
/// Personal annotations: favorite, note, tags (issue #39).
Annotate,
/// Local model runtime action: prune of unused models (issue #73).
Model,
/// Remote catalog action: update or add of an external catalog (issue #60).
Catalog,
/// Pre-update backup created (issue #64).
Backup,
/// Rollback of a previous backup (issue #64).
Rollback,
/// Derived index pruned (sessions retention, issue #53).
Prune,
/// Session usage recorded: tokens + estimated cost (issue #49).
Cost,
/// Threshold alert from 'am monitor' (issue #50).
Alert,
/// One benchmark run of 'am lab' (issue #62).
Lab,
/// LLM provider registry change: add, remove, default (issue #88).
Provider,
}
impl EventKind {
pub fn as_str(&self) -> &'static str {
match self {
EventKind::Start => "start",
EventKind::Stop => "stop",
EventKind::Run => "run",
EventKind::Restart => "restart",
EventKind::Install => "install",
EventKind::Update => "update",
EventKind::Uninstall => "uninstall",
EventKind::Doctor => "doctor",
EventKind::Config => "config",
EventKind::Repl => "repl",
EventKind::Shell => "shell",
EventKind::Annotate => "annotate",
EventKind::Model => "model",
EventKind::Catalog => "catalog",
EventKind::Backup => "backup",
EventKind::Rollback => "rollback",
EventKind::Prune => "prune",
EventKind::Cost => "cost",
EventKind::Alert => "alert",
EventKind::Lab => "lab",
EventKind::Provider => "provider",
}
}
}
/// 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>,
/// Input tokens billed for the session (issue #49).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tokens_in: Option<u64>,
/// Output tokens billed for the session (issue #49).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tokens_out: Option<u64>,
/// Estimated cost in USD (issue #49).
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cost_usd: Option<f64>,
}
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,
tokens_in: None,
tokens_out: None,
cost_usd: 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: Option<String>) -> Self {
self.project = project;
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, d: u64) -> Self {
self.duration_s = Some(d);
self
}
/// Record token usage (issue #49).
pub fn with_usage(mut self, tokens_in: u64, tokens_out: u64) -> Self {
self.tokens_in = Some(tokens_in);
self.tokens_out = Some(tokens_out);
self
}
/// Record the estimated cost in USD (issue #49).
pub fn with_cost(mut self, usd: f64) -> Self {
self.cost_usd = Some(usd);
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");
}
}