pub mod harness_adapt;
pub mod harness_metrics;
pub mod observability;
pub mod tool_receipts;
pub use observability::{
evaluate_alerts, summarize, summarize_log, Alert, AlertKind, AlertThresholds, MetricsSummary,
};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use serde_json::Value;
use std::collections::HashMap;
use std::fs::{self, OpenOptions};
use std::io::{BufRead, BufReader, BufWriter, Write};
use std::path::{Path, PathBuf};
use std::sync::mpsc;
use std::thread;
use uuid::Uuid;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct EventLogStats {
pub events: usize,
pub spans: usize,
pub approx_event_bytes: usize,
pub approx_span_bytes: usize,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum EventKind {
ProposalReceived,
ActionValidated,
ActionRejected,
ActionExecuting,
ActionSucceeded,
ActionFailed,
ActionSkipped,
ActionRetrying,
ActionDeduplicated,
PolicyViolation,
StateChanged,
StateSnapshot,
StateRollback,
SkillDistilled,
SkillEvolved,
SkillDeprecated,
EvolutionTriggered,
CandidatePromoted,
CandidateRejected,
Consolidated,
ReplanAttempted,
ReplanProposalReceived,
ReplanRejected,
ReplanExhausted,
VoiceFastTurnStarted,
VoiceFastTurnEnded,
VoiceSidecarResolved,
VoiceSidecarFailed,
VoiceSidecarTimedOut,
VoiceTurnCancelled,
VoiceBridgePlayed,
GateAccepted,
GateRejected,
SessionScope,
PermissionDecision,
ApprovalRecorded,
BranchDecision,
AlternativeRejected,
InferenceMetered,
TransactionConflict,
AdmissionGateDecision,
ToolReceiptHallucination,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
#[serde(rename_all = "snake_case")]
pub enum SpanStatus {
Ok,
Error,
Unset,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Span {
pub trace_id: String,
pub span_id: String,
pub parent_span_id: Option<String>,
pub name: String,
pub start_time: DateTime<Utc>,
pub end_time: Option<DateTime<Utc>>,
pub status: SpanStatus,
pub attributes: HashMap<String, Value>,
}
pub mod metric_keys {
pub const DURATION_MS: &str = "duration_ms";
pub const TOKENS_IN: &str = "tokens_in";
pub const TOKENS_OUT: &str = "tokens_out";
pub const COST_USD: &str = "cost_usd";
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
pub struct Metrics {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub duration_ms: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tokens_in: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tokens_out: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub cost_usd: Option<f64>,
}
impl Metrics {
pub fn latency(duration_ms: f64) -> Self {
Self {
duration_ms: Some(duration_ms),
..Default::default()
}
}
pub fn inference(tokens_in: u64, tokens_out: u64, cost_usd: Option<f64>) -> Self {
Self {
duration_ms: None,
tokens_in: Some(tokens_in),
tokens_out: Some(tokens_out),
cost_usd,
}
}
pub fn with_duration(mut self, duration_ms: f64) -> Self {
self.duration_ms = Some(duration_ms);
self
}
fn merge_into(&self, data: &mut HashMap<String, Value>) {
if let Some(d) = self.duration_ms {
data.insert(metric_keys::DURATION_MS.into(), Value::from(d));
}
if let Some(t) = self.tokens_in {
data.insert(metric_keys::TOKENS_IN.into(), Value::from(t));
}
if let Some(t) = self.tokens_out {
data.insert(metric_keys::TOKENS_OUT.into(), Value::from(t));
}
if let Some(c) = self.cost_usd {
data.insert(metric_keys::COST_USD.into(), Value::from(c));
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Event {
pub kind: EventKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub action_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub proposal_id: Option<String>,
#[serde(default)]
pub data: HashMap<String, Value>,
#[serde(default = "Utc::now")]
pub timestamp: DateTime<Utc>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prev_hash: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub hash: Option<String>,
}
impl Event {
pub fn duration_ms(&self) -> Option<f64> {
self.data
.get(metric_keys::DURATION_MS)
.and_then(Value::as_f64)
}
pub fn tokens_in(&self) -> Option<u64> {
self.data
.get(metric_keys::TOKENS_IN)
.and_then(Value::as_u64)
}
pub fn tokens_out(&self) -> Option<u64> {
self.data
.get(metric_keys::TOKENS_OUT)
.and_then(Value::as_u64)
}
pub fn cost_usd(&self) -> Option<f64> {
self.data.get(metric_keys::COST_USD).and_then(Value::as_f64)
}
pub fn metrics(&self) -> Metrics {
Metrics {
duration_ms: self.duration_ms(),
tokens_in: self.tokens_in(),
tokens_out: self.tokens_out(),
cost_usd: self.cost_usd(),
}
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Serialize, Deserialize)]
pub struct MetricsTotals {
pub duration_ms: f64,
pub tokens_in: u64,
pub tokens_out: u64,
pub tokens: u64,
pub cost_usd: f64,
pub metered_events: usize,
}
pub fn metrics_totals_of(events: &[Event]) -> MetricsTotals {
let mut totals = MetricsTotals::default();
for ev in events {
let m = ev.metrics();
let mut metered = false;
if let Some(d) = m.duration_ms {
totals.duration_ms += d;
metered = true;
}
if let Some(t) = m.tokens_in {
totals.tokens_in = totals.tokens_in.saturating_add(t);
metered = true;
}
if let Some(t) = m.tokens_out {
totals.tokens_out = totals.tokens_out.saturating_add(t);
metered = true;
}
if let Some(c) = m.cost_usd {
totals.cost_usd += c;
metered = true;
}
if metered {
totals.metered_events += 1;
}
}
totals.tokens = totals.tokens_in.saturating_add(totals.tokens_out);
totals
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct AgentCost {
pub agent: String,
pub calls: u64,
pub tokens_in: u64,
pub tokens_out: u64,
pub cost_usd: f64,
}
pub fn cost_by_agent_of(events: &[Event]) -> Vec<AgentCost> {
use std::collections::BTreeMap;
let mut map: BTreeMap<String, AgentCost> = BTreeMap::new();
for e in events {
if e.kind != EventKind::InferenceMetered {
continue;
}
let agent = e
.data
.get("agent")
.and_then(|v| v.as_str())
.unwrap_or("unknown")
.to_string();
let entry = map.entry(agent.clone()).or_insert_with(|| AgentCost {
agent,
..Default::default()
});
entry.calls += 1;
entry.tokens_in = entry.tokens_in.saturating_add(e.tokens_in().unwrap_or(0));
entry.tokens_out = entry.tokens_out.saturating_add(e.tokens_out().unwrap_or(0));
entry.cost_usd += e.cost_usd().unwrap_or(0.0);
}
map.into_values().collect()
}
struct JournalWriter {
tx: Option<mpsc::Sender<String>>,
handle: Option<thread::JoinHandle<()>>,
}
impl JournalWriter {
fn spawn(path: PathBuf) -> Self {
let (tx, rx) = mpsc::channel::<String>();
match thread::Builder::new()
.name("car-eventlog-journal".into())
.spawn(move || journal_loop(path, rx))
{
Ok(handle) => Self {
tx: Some(tx),
handle: Some(handle),
},
Err(e) => {
tracing::warn!(error = %e, "car-eventlog: failed to spawn journal writer thread — journaling disabled for this log");
Self {
tx: None,
handle: None,
}
}
}
}
fn send(&self, line: String) {
if let Some(tx) = &self.tx {
let _ = tx.send(line);
}
}
}
impl Drop for JournalWriter {
fn drop(&mut self) {
self.tx.take();
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
fn journal_loop(path: PathBuf, rx: mpsc::Receiver<String>) {
let file = match OpenOptions::new().create(true).append(true).open(&path) {
Ok(file) => file,
Err(e) => {
tracing::warn!(path = %path.display(), error = %e, "car-eventlog: cannot open journal file — events for this log will not be persisted");
while rx.recv().is_ok() {}
return;
}
};
let mut writer = BufWriter::new(file);
while let Ok(line) = rx.recv() {
let _ = writeln!(writer, "{line}");
while let Ok(more) = rx.try_recv() {
let _ = writeln!(writer, "{more}");
}
let _ = writer.flush();
}
let _ = writer.flush();
}
pub struct EventLog {
events: Vec<Event>,
spans: Vec<Span>,
journal: Option<JournalWriter>,
hash_chaining: bool,
last_hash: Option<String>,
retention: Option<RetentionPolicy>,
journal_path: Option<PathBuf>,
journal_lines: usize,
trimmed_events: u64,
cumulative_cost_usd: f64,
}
const JOURNAL_COMPACT_MIN_EXCESS: usize = 1024;
fn event_digest(
prev_hash: &str,
kind: &EventKind,
action_id: Option<&str>,
proposal_id: Option<&str>,
data: &HashMap<String, Value>,
timestamp: &DateTime<Utc>,
) -> String {
use sha2::{Digest, Sha256};
let mut sorted: Vec<(&String, &Value)> = data.iter().collect();
sorted.sort_by(|a, b| a.0.cmp(b.0));
let data_canon: String = sorted
.iter()
.map(|(k, v)| format!("{k}={}", v))
.collect::<Vec<_>>()
.join("\u{1f}");
let kind_str = serde_json::to_string(kind).unwrap_or_default();
let mut hasher = Sha256::new();
hasher.update(prev_hash.as_bytes());
hasher.update(b"\x1e");
hasher.update(kind_str.as_bytes());
hasher.update(b"\x1e");
hasher.update(action_id.unwrap_or("").as_bytes());
hasher.update(b"\x1e");
hasher.update(proposal_id.unwrap_or("").as_bytes());
hasher.update(b"\x1e");
hasher.update(data_canon.as_bytes());
hasher.update(b"\x1e");
hasher.update(timestamp.to_rfc3339().as_bytes());
let digest = hasher.finalize();
digest.iter().map(|b| format!("{b:02x}")).collect()
}
#[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)]
pub struct RetentionPolicy {
#[serde(default)]
pub max_events: Option<usize>,
#[serde(default)]
pub max_age_secs: Option<i64>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct EventQuery {
#[serde(default)]
pub kinds: Vec<EventKind>,
#[serde(default)]
pub action_id: Option<String>,
#[serde(default)]
pub proposal_id: Option<String>,
#[serde(default)]
pub since: Option<DateTime<Utc>>,
#[serde(default)]
pub until: Option<DateTime<Utc>>,
#[serde(default)]
pub data_matches: std::collections::HashMap<String, String>,
#[serde(default)]
pub limit: Option<usize>,
}
fn data_value_matches(v: &Value, want: &str) -> bool {
match v {
Value::String(s) => s == want,
Value::Null => false,
other => *other == want,
}
}
impl EventQuery {
pub fn matches(&self, e: &Event) -> bool {
if !self.kinds.is_empty() && !self.kinds.contains(&e.kind) {
return false;
}
if let Some(aid) = &self.action_id {
if e.action_id.as_deref() != Some(aid.as_str()) {
return false;
}
}
if let Some(pid) = &self.proposal_id {
if e.proposal_id.as_deref() != Some(pid.as_str()) {
return false;
}
}
if let Some(since) = self.since {
if e.timestamp < since {
return false;
}
}
if let Some(until) = self.until {
if e.timestamp >= until {
return false;
}
}
for (k, want) in &self.data_matches {
match e.data.get(k) {
Some(v) if data_value_matches(v, want) => {}
_ => return false,
}
}
true
}
}
impl EventLog {
pub fn new() -> Self {
Self {
events: Vec::new(),
spans: Vec::new(),
journal: None,
hash_chaining: false,
last_hash: None,
retention: None,
journal_path: None,
journal_lines: 0,
trimmed_events: 0,
cumulative_cost_usd: 0.0,
}
}
pub fn with_journal(path: PathBuf) -> Self {
if let Some(parent) = path.parent() {
let _ = fs::create_dir_all(parent);
}
Self {
events: Vec::new(),
spans: Vec::new(),
journal: Some(JournalWriter::spawn(path.clone())),
hash_chaining: false,
last_hash: None,
retention: None,
journal_path: Some(path),
journal_lines: 0,
trimmed_events: 0,
cumulative_cost_usd: 0.0,
}
}
pub fn with_hash_chaining(mut self) -> Self {
self.enable_hash_chaining();
self
}
pub fn enable_hash_chaining(&mut self) {
self.hash_chaining = true;
if self.last_hash.is_none() {
self.last_hash = self.events.last().and_then(|e| e.hash.clone());
}
}
pub fn hash_chaining_enabled(&self) -> bool {
self.hash_chaining
}
pub fn append(
&mut self,
kind: EventKind,
action_id: Option<&str>,
proposal_id: Option<&str>,
data: HashMap<String, Value>,
) -> &Event {
let timestamp = Utc::now();
let (prev_hash, hash) = if self.hash_chaining {
let prev = self.last_hash.clone().unwrap_or_default();
let h = event_digest(&prev, &kind, action_id, proposal_id, &data, ×tamp);
self.last_hash = Some(h.clone());
(Some(prev), Some(h))
} else {
(None, None)
};
let event = Event {
kind,
action_id: action_id.map(|s| s.to_string()),
proposal_id: proposal_id.map(|s| s.to_string()),
data,
timestamp,
prev_hash,
hash,
};
if let Some(journal) = &self.journal {
if let Ok(json) = serde_json::to_string(&event) {
journal.send(json);
self.journal_lines += 1;
}
}
if let Some(c) = event.cost_usd() {
self.cumulative_cost_usd += c;
}
self.events.push(event);
if let Some(max) = self.retention.as_ref().and_then(|p| p.max_events) {
if self.events.len() > max {
let removed = truncate_vec_keep_last(&mut self.events, max);
self.trimmed_events += removed as u64;
self.maybe_compact_journal();
}
}
self.events.last().unwrap()
}
pub fn verify_chain(&self) -> Result<usize, usize> {
let mut prev = String::new();
let mut verified = 0usize;
let mut chain_started = false;
for (i, ev) in self.events.iter().enumerate() {
let Some(stored) = &ev.hash else {
if chain_started {
return Err(i);
}
continue;
};
let recorded_prev = ev.prev_hash.clone().unwrap_or_default();
if chain_started && recorded_prev != prev {
return Err(i);
}
let recomputed = event_digest(
&recorded_prev,
&ev.kind,
ev.action_id.as_deref(),
ev.proposal_id.as_deref(),
&ev.data,
&ev.timestamp,
);
if &recomputed != stored {
return Err(i);
}
prev = stored.clone();
chain_started = true;
verified += 1;
}
Ok(verified)
}
pub fn append_metered(
&mut self,
kind: EventKind,
action_id: Option<&str>,
proposal_id: Option<&str>,
mut data: HashMap<String, Value>,
metrics: Metrics,
) -> &Event {
metrics.merge_into(&mut data);
self.append(kind, action_id, proposal_id, data)
}
pub fn metrics_totals(&self) -> MetricsTotals {
metrics_totals_of(&self.events)
}
pub fn cost_by_agent(&self) -> Vec<AgentCost> {
cost_by_agent_of(&self.events)
}
pub fn events(&self) -> &[Event] {
&self.events
}
pub fn len(&self) -> usize {
self.events.len()
}
pub fn span_len(&self) -> usize {
self.spans.len()
}
pub fn is_empty(&self) -> bool {
self.events.is_empty()
}
pub fn stats(&self) -> EventLogStats {
EventLogStats {
events: self.events.len(),
spans: self.spans.len(),
approx_event_bytes: approx_json_bytes(&self.events),
approx_span_bytes: approx_json_bytes(&self.spans),
}
}
pub fn truncate_events_keep_last(&mut self, keep_last: usize) -> usize {
let removed = truncate_vec_keep_last(&mut self.events, keep_last);
self.trimmed_events += removed as u64;
if removed > 0 {
self.maybe_compact_journal();
}
removed
}
pub fn truncate_spans_keep_last(&mut self, keep_last: usize) -> usize {
truncate_vec_keep_last(&mut self.spans, keep_last)
}
pub fn clear(&mut self) -> EventLogStats {
let removed = self.stats();
self.trimmed_events += removed.events as u64;
self.events.clear();
self.events.shrink_to_fit();
self.spans.clear();
self.spans.shrink_to_fit();
removed
}
pub fn trimmed_events(&self) -> u64 {
self.trimmed_events
}
pub fn cumulative_cost_usd(&self) -> f64 {
self.cumulative_cost_usd
}
pub fn journal_size_bytes(&self) -> Option<u64> {
let path = self.journal_path.as_ref()?;
fs::metadata(path).ok().map(|m| m.len())
}
fn maybe_compact_journal(&mut self) {
if self.journal_path.is_none() {
return;
}
let excess = self.journal_lines.saturating_sub(self.events.len());
if excess >= JOURNAL_COMPACT_MIN_EXCESS && excess.saturating_mul(4) >= self.journal_lines {
self.compact_journal();
}
}
pub fn compact_journal(&mut self) -> bool {
let Some(path) = self.journal_path.clone() else {
return false;
};
self.journal = None;
let tmp = path.with_extension("compact-tmp");
let rewrite = (|| -> std::io::Result<()> {
{
let file = fs::File::create(&tmp)?;
let mut w = BufWriter::new(file);
for ev in &self.events {
let line = serde_json::to_string(ev).map_err(std::io::Error::other)?;
writeln!(w, "{line}")?;
}
w.flush()?;
}
fs::rename(&tmp, &path)
})();
let ok = match rewrite {
Ok(()) => {
self.journal_lines = self.events.len();
true
}
Err(e) => {
let _ = fs::remove_file(&tmp);
tracing::warn!(
path = %path.display(), error = %e,
"car-eventlog: journal compaction failed — journal keeps growing until the next successful compaction"
);
false
}
};
self.journal = Some(JournalWriter::spawn(path));
ok
}
pub fn query(&self, query: &EventQuery) -> Vec<&Event> {
let mut out: Vec<&Event> = self.events.iter().filter(|e| query.matches(e)).collect();
out.reverse(); if let Some(limit) = query.limit.filter(|l| *l > 0) {
out.truncate(limit);
}
out
}
pub fn set_retention(&mut self, policy: Option<RetentionPolicy>) {
self.retention = policy;
}
pub fn retention(&self) -> Option<&RetentionPolicy> {
self.retention.as_ref()
}
pub fn enforce_retention(&mut self, policy: &RetentionPolicy, now: DateTime<Utc>) -> usize {
let before = self.events.len();
if let Some(age) = policy.max_age_secs {
let cutoff = now - chrono::Duration::seconds(age);
self.events.retain(|e| e.timestamp >= cutoff);
}
if let Some(max) = policy.max_events {
truncate_vec_keep_last(&mut self.events, max);
}
let removed = before.saturating_sub(self.events.len());
self.trimmed_events += removed as u64;
if removed > 0 {
self.maybe_compact_journal();
}
removed
}
pub fn filter(&self, kind: Option<&EventKind>, action_id: Option<&str>) -> Vec<&Event> {
self.events
.iter()
.filter(|e| {
if let Some(k) = kind {
if &e.kind != k {
return false;
}
}
if let Some(aid) = action_id {
if e.action_id.as_deref() != Some(aid) {
return false;
}
}
true
})
.collect()
}
pub fn begin_span(
&mut self,
name: &str,
trace_id: &str,
parent_span_id: Option<&str>,
attributes: HashMap<String, Value>,
) -> String {
let span_id = Uuid::new_v4().to_string();
let span = Span {
trace_id: trace_id.to_string(),
span_id: span_id.clone(),
parent_span_id: parent_span_id.map(|s| s.to_string()),
name: name.to_string(),
start_time: Utc::now(),
end_time: None,
status: SpanStatus::Unset,
attributes,
};
self.spans.push(span);
span_id
}
pub fn end_span(&mut self, span_id: &str, status: SpanStatus) {
if let Some(span) = self.spans.iter_mut().find(|s| s.span_id == span_id) {
span.end_time = Some(Utc::now());
span.status = status;
}
}
pub fn spans(&self) -> Vec<Span> {
self.spans.clone()
}
pub fn export_traces(&self) -> String {
let mut traces: HashMap<&str, Vec<&Span>> = HashMap::new();
for span in &self.spans {
traces.entry(span.trace_id.as_str()).or_default().push(span);
}
let resource_spans: Vec<Value> = traces.into_values().map(|spans| {
let scope_spans = spans
.iter()
.map(|s| {
let mut span_obj = serde_json::json!({
"traceId": s.trace_id,
"spanId": s.span_id,
"name": s.name,
"startTimeUnixNano": s.start_time.timestamp_nanos_opt().unwrap_or(0).to_string(),
"status": {
"code": match s.status {
SpanStatus::Ok => 1,
SpanStatus::Error => 2,
SpanStatus::Unset => 0,
}
},
"attributes": s.attributes.iter().map(|(k, v)| {
serde_json::json!({
"key": k,
"value": { "stringValue": v.to_string() }
})
}).collect::<Vec<_>>(),
});
if let Some(ref parent) = s.parent_span_id {
span_obj.as_object_mut().unwrap().insert(
"parentSpanId".to_string(),
Value::from(parent.as_str()),
);
}
if let Some(end) = s.end_time {
span_obj.as_object_mut().unwrap().insert(
"endTimeUnixNano".to_string(),
Value::from(end.timestamp_nanos_opt().unwrap_or(0).to_string()),
);
}
span_obj
})
.collect::<Vec<_>>();
serde_json::json!({
"resource": {
"attributes": [
{ "key": "service.name", "value": { "stringValue": "car-runtime" } }
]
},
"scopeSpans": [{
"scope": { "name": "car-eventlog" },
"spans": scope_spans
}]
})
})
.collect();
serde_json::to_string(&serde_json::json!({
"resourceSpans": resource_spans
}))
.unwrap_or_else(|_| "{}".to_string())
}
pub fn load(path: &Path) -> std::io::Result<Self> {
let file = fs::File::open(path)?;
let reader = BufReader::new(file);
let mut events = Vec::new();
for line in reader.lines() {
let line = line?;
let line = line.trim();
if !line.is_empty() {
if let Ok(event) = serde_json::from_str::<Event>(line) {
events.push(event);
}
}
}
let last_hash = events.last().and_then(|e| e.hash.clone());
let hash_chaining = last_hash.is_some();
let cumulative_cost_usd = events.iter().filter_map(Event::cost_usd).sum();
let journal_lines = events.len();
Ok(Self {
events,
spans: Vec::new(),
journal: Some(JournalWriter::spawn(path.to_path_buf())),
hash_chaining,
last_hash,
retention: None,
journal_path: Some(path.to_path_buf()),
journal_lines,
trimmed_events: 0,
cumulative_cost_usd,
})
}
}
fn approx_json_bytes<T: Serialize>(value: &T) -> usize {
serde_json::to_vec(value)
.map(|bytes| bytes.len())
.unwrap_or(0)
}
fn truncate_vec_keep_last<T>(items: &mut Vec<T>, keep_last: usize) -> usize {
let len = items.len();
if len <= keep_last {
return 0;
}
let removed = len - keep_last;
items.drain(..removed);
items.shrink_to_fit();
removed
}
impl Default for EventLog {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn append_and_read() {
let mut log = EventLog::new();
log.append(
EventKind::ProposalReceived,
None,
Some("p1"),
[("source".to_string(), Value::from("test"))].into(),
);
assert_eq!(log.len(), 1);
assert_eq!(log.events()[0].kind, EventKind::ProposalReceived);
}
#[test]
fn query_filters_by_kind_data_and_time() {
let mut log = EventLog::new();
log.append(
EventKind::PermissionDecision,
Some("a1"),
None,
[
("caller".to_string(), Value::from("alice")),
("tool".to_string(), Value::from("shell")),
]
.into(),
);
log.append(
EventKind::PermissionDecision,
Some("a2"),
None,
[
("caller".to_string(), Value::from("bob")),
("tool".to_string(), Value::from("shell")),
]
.into(),
);
log.append(
EventKind::StateChanged,
Some("a3"),
None,
Default::default(),
);
let q = EventQuery {
kinds: vec![EventKind::PermissionDecision],
..Default::default()
};
assert_eq!(log.query(&q).len(), 2);
let q = EventQuery {
data_matches: [("caller".to_string(), "alice".to_string())].into(),
..Default::default()
};
let hits = log.query(&q);
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].action_id.as_deref(), Some("a1"));
let q = EventQuery {
kinds: vec![EventKind::PermissionDecision],
data_matches: [("tool".to_string(), "shell".to_string())].into(),
limit: Some(1),
..Default::default()
};
let hits = log.query(&q);
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].action_id.as_deref(), Some("a2"));
}
#[test]
fn cost_by_agent_folds_metered_events() {
let mut log = EventLog::new();
log.append_metered(
EventKind::InferenceMetered,
None,
None,
[("agent".to_string(), Value::from("researcher"))].into(),
Metrics {
tokens_in: Some(100),
tokens_out: Some(50),
cost_usd: Some(2.0),
..Default::default()
},
);
log.append_metered(
EventKind::InferenceMetered,
None,
None,
[("agent".to_string(), Value::from("researcher"))].into(),
Metrics {
tokens_in: Some(10),
tokens_out: Some(5),
cost_usd: Some(0.2),
..Default::default()
},
);
log.append_metered(
EventKind::InferenceMetered,
None,
None,
[("agent".to_string(), Value::from("coordinator"))].into(),
Metrics {
cost_usd: Some(0.5),
..Default::default()
},
);
let report = log.cost_by_agent();
assert_eq!(report.len(), 2);
assert_eq!(report[0].agent, "coordinator");
assert_eq!(report[0].cost_usd, 0.5);
assert_eq!(report[1].agent, "researcher");
assert_eq!(report[1].calls, 2);
assert_eq!(report[1].tokens_in, 110);
assert_eq!(report[1].tokens_out, 55);
assert!((report[1].cost_usd - 2.2).abs() < 1e-9);
}
#[test]
fn auto_retention_caps_event_count() {
let mut log = EventLog::new();
log.set_retention(Some(RetentionPolicy {
max_events: Some(3),
max_age_secs: None,
}));
for i in 0..10 {
log.append(
EventKind::StateChanged,
Some(&format!("a{i}")),
None,
Default::default(),
);
}
assert_eq!(log.len(), 3);
assert_eq!(log.events()[0].action_id.as_deref(), Some("a7"));
assert_eq!(log.events()[2].action_id.as_deref(), Some("a9"));
}
#[test]
fn enforce_retention_drops_old_by_age() {
let mut log = EventLog::new();
log.append(
EventKind::StateChanged,
Some("old"),
None,
Default::default(),
);
log.events[0].timestamp = Utc::now() - chrono::Duration::seconds(3600);
log.append(
EventKind::StateChanged,
Some("fresh"),
None,
Default::default(),
);
let removed = log.enforce_retention(
&RetentionPolicy {
max_events: None,
max_age_secs: Some(60),
},
Utc::now(),
);
assert_eq!(removed, 1);
assert_eq!(log.len(), 1);
assert_eq!(log.events()[0].action_id.as_deref(), Some("fresh"));
}
#[test]
fn retention_trims_are_counted() {
let mut log = EventLog::new();
log.set_retention(Some(RetentionPolicy {
max_events: Some(2),
max_age_secs: None,
}));
for i in 0..5 {
log.append(
EventKind::StateChanged,
Some(&format!("a{i}")),
None,
Default::default(),
);
}
assert_eq!(log.trimmed_events(), 3);
assert_eq!(log.truncate_events_keep_last(1), 1);
assert_eq!(log.trimmed_events(), 4);
log.clear();
assert_eq!(log.trimmed_events(), 5);
}
#[test]
fn cumulative_cost_is_monotonic_across_trims_and_reload() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("cost.jsonl");
{
let mut log = EventLog::with_journal(journal.clone());
log.set_retention(Some(RetentionPolicy {
max_events: Some(1),
max_age_secs: None,
}));
for _ in 0..4 {
log.append_metered(
EventKind::InferenceMetered,
None,
None,
Default::default(),
Metrics {
cost_usd: Some(2.5),
..Default::default()
},
);
}
assert_eq!(log.len(), 1);
assert!((log.cumulative_cost_usd() - 10.0).abs() < 1e-9);
}
let reloaded = EventLog::load(&journal).unwrap();
assert!((reloaded.cumulative_cost_usd() - 10.0).abs() < 1e-9);
}
#[test]
fn journal_compaction_rewrites_to_retained_set() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("compact.jsonl");
let keep = 16usize;
let total = keep + JOURNAL_COMPACT_MIN_EXCESS + 8;
{
let mut log = EventLog::with_journal(journal.clone());
log.set_retention(Some(RetentionPolicy {
max_events: Some(keep),
max_age_secs: None,
}));
for i in 0..total {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
None,
HashMap::new(),
);
}
assert_eq!(log.len(), keep);
assert!(log.journal_size_bytes().unwrap_or(0) > 0);
}
let reloaded = EventLog::load(&journal).unwrap();
assert!(
reloaded.len() < total,
"journal must have been compacted (got {} lines)",
reloaded.len()
);
assert_eq!(
reloaded.events().last().unwrap().action_id.as_deref(),
Some(format!("a{}", total - 1).as_str())
);
}
#[test]
fn compact_journal_preserves_hash_chain_of_retained_tail() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("chained.jsonl");
{
let mut log = EventLog::with_journal(journal.clone()).with_hash_chaining();
for i in 0..20 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
None,
HashMap::new(),
);
}
log.truncate_events_keep_last(5);
assert!(log.compact_journal(), "compaction must succeed");
log.append(
EventKind::ActionSucceeded,
Some("post"),
None,
HashMap::new(),
);
}
let reloaded = EventLog::load(&journal).unwrap();
assert_eq!(reloaded.len(), 6);
assert_eq!(reloaded.verify_chain(), Ok(6), "retained tail must verify");
assert_eq!(reloaded.events()[0].action_id.as_deref(), Some("a15"));
assert_eq!(reloaded.events()[5].action_id.as_deref(), Some("post"));
}
#[test]
fn compact_journal_without_journal_is_noop() {
let mut log = EventLog::new();
log.append(EventKind::StateChanged, Some("a"), None, Default::default());
assert!(!log.compact_journal());
assert_eq!(log.journal_size_bytes(), None);
}
#[test]
fn chaining_off_by_default_no_hashes() {
let mut log = EventLog::new();
log.append(
EventKind::ActionSucceeded,
Some("a1"),
Some("p1"),
HashMap::new(),
);
assert!(!log.hash_chaining_enabled());
assert!(log.events()[0].hash.is_none());
assert!(log.events()[0].prev_hash.is_none());
assert_eq!(log.verify_chain(), Ok(0));
}
#[test]
fn hash_chain_verifies_clean_log() {
let mut log = EventLog::new().with_hash_chaining();
for i in 0..5 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
Some("p"),
[("i".to_string(), Value::from(i))].into(),
);
}
assert!(log.events().iter().all(|e| e.hash.is_some()));
assert_eq!(log.verify_chain(), Ok(5));
assert_eq!(log.events()[0].prev_hash.as_deref(), Some(""));
for w in log.events().windows(2) {
assert_eq!(w[1].prev_hash, w[0].hash);
}
}
#[test]
fn tampering_with_data_breaks_chain() {
let mut log = EventLog::new().with_hash_chaining();
for i in 0..4 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
Some("p"),
[("v".to_string(), Value::from(i))].into(),
);
}
assert_eq!(log.verify_chain(), Ok(4));
log.events[2].data.insert("v".to_string(), Value::from(999));
assert_eq!(log.verify_chain(), Err(2));
}
#[test]
fn deleting_an_event_breaks_chain() {
let mut log = EventLog::new().with_hash_chaining();
for i in 0..4 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
Some("p"),
HashMap::new(),
);
}
log.events.remove(1);
assert_eq!(log.verify_chain(), Err(1));
}
#[test]
fn chain_survives_serialize_roundtrip() {
let mut log = EventLog::new().with_hash_chaining();
for i in 0..3 {
log.append(
EventKind::PermissionDecision,
Some(&format!("a{i}")),
Some("p"),
[
("decision".to_string(), Value::from("allow")),
("nested".to_string(), serde_json::json!({"z": 1, "a": 2})),
]
.into(),
);
}
let lines: Vec<String> = log
.events()
.iter()
.map(|e| serde_json::to_string(e).unwrap())
.collect();
let mut rebuilt = EventLog::new();
for line in &lines {
rebuilt.events.push(serde_json::from_str(line).unwrap());
}
assert_eq!(rebuilt.verify_chain(), Ok(3));
}
#[test]
fn chain_survives_journal_load_and_append() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("chain.jsonl");
{
let mut log = EventLog::with_journal(journal.clone());
log.enable_hash_chaining();
log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
log.append(EventKind::ActionSucceeded, Some("a2"), None, HashMap::new());
}
{
let mut log = EventLog::load(&journal).unwrap();
assert!(
log.hash_chaining_enabled(),
"loading a chained tail re-enables chaining"
);
log.append(EventKind::ActionSucceeded, Some("a3"), None, HashMap::new());
assert_eq!(log.verify_chain(), Ok(3), "post-load append stays chained");
}
let reloaded = EventLog::load(&journal).unwrap();
assert_eq!(reloaded.len(), 3);
assert_eq!(reloaded.verify_chain(), Ok(3));
let plain = dir.path().join("plain.jsonl");
{
let mut log = EventLog::with_journal(plain.clone());
log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
}
let loaded = EventLog::load(&plain).unwrap();
assert!(!loaded.hash_chaining_enabled(), "unchained tail stays off");
}
#[test]
fn metered_event_carries_metrics_in_data() {
let mut log = EventLog::new();
log.append_metered(
EventKind::ActionSucceeded,
Some("a1"),
Some("p1"),
[("tool".to_string(), Value::from("search"))].into(),
Metrics::inference(120, 45, Some(0.0012)).with_duration(83.0),
);
let ev = &log.events()[0];
assert_eq!(ev.data.get("tool").unwrap(), "search");
assert_eq!(ev.duration_ms(), Some(83.0));
assert_eq!(ev.tokens_in(), Some(120));
assert_eq!(ev.tokens_out(), Some(45));
assert_eq!(ev.cost_usd(), Some(0.0012));
}
#[test]
fn metrics_totals_sum_across_events() {
let mut log = EventLog::new();
log.append_metered(
EventKind::ActionSucceeded,
Some("a1"),
None,
HashMap::new(),
Metrics::latency(50.0),
);
log.append_metered(
EventKind::ActionSucceeded,
Some("a2"),
None,
HashMap::new(),
Metrics::inference(100, 20, Some(0.5)).with_duration(70.0),
);
log.append(EventKind::ProposalReceived, None, None, HashMap::new());
let t = log.metrics_totals();
assert_eq!(t.duration_ms, 120.0);
assert_eq!(t.tokens_in, 100);
assert_eq!(t.tokens_out, 20);
assert_eq!(t.tokens, 120);
assert_eq!(t.cost_usd, 0.5);
assert_eq!(t.metered_events, 2);
}
#[test]
fn metrics_totals_counts_raw_appended_duration_key() {
let mut log = EventLog::new();
log.append(
EventKind::ActionSucceeded,
Some("a1"),
None,
[(metric_keys::DURATION_MS.to_string(), Value::from(42.0))].into(),
);
let t = log.metrics_totals();
assert_eq!(t.duration_ms, 42.0);
assert_eq!(t.metered_events, 1);
}
#[test]
fn new_telemetry_event_kinds_serialize_snake_case() {
let json = serde_json::to_string(&EventKind::BranchDecision).unwrap();
assert_eq!(json, "\"branch_decision\"");
let json = serde_json::to_string(&EventKind::AlternativeRejected).unwrap();
assert_eq!(json, "\"alternative_rejected\"");
let json = serde_json::to_string(&EventKind::InferenceMetered).unwrap();
assert_eq!(json, "\"inference_metered\"");
}
#[test]
fn filter_by_kind() {
let mut log = EventLog::new();
log.append(
EventKind::ProposalReceived,
None,
Some("p1"),
HashMap::new(),
);
log.append(
EventKind::ActionValidated,
Some("a1"),
Some("p1"),
HashMap::new(),
);
log.append(
EventKind::ActionSucceeded,
Some("a1"),
Some("p1"),
HashMap::new(),
);
let validated = log.filter(Some(&EventKind::ActionValidated), None);
assert_eq!(validated.len(), 1);
}
#[test]
fn filter_by_action_id() {
let mut log = EventLog::new();
log.append(EventKind::ActionValidated, Some("a1"), None, HashMap::new());
log.append(EventKind::ActionValidated, Some("a2"), None, HashMap::new());
let a1_events = log.filter(None, Some("a1"));
assert_eq!(a1_events.len(), 1);
}
#[test]
fn journal_write_and_reload() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("events.jsonl");
{
let mut log = EventLog::with_journal(journal.clone());
log.append(
EventKind::ProposalReceived,
None,
Some("p1"),
HashMap::new(),
);
log.append(
EventKind::ActionSucceeded,
Some("a1"),
Some("p1"),
HashMap::new(),
);
}
assert!(journal.exists());
let reloaded = EventLog::load(&journal).unwrap();
assert_eq!(reloaded.len(), 2);
assert_eq!(reloaded.events()[0].kind, EventKind::ProposalReceived);
assert_eq!(reloaded.events()[1].kind, EventKind::ActionSucceeded);
}
#[test]
fn journal_preserves_order_and_count_under_burst() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("burst.jsonl");
{
let mut log = EventLog::with_journal(journal.clone());
for i in 0..500 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{i}")),
None,
HashMap::new(),
);
}
}
let reloaded = EventLog::load(&journal).unwrap();
assert_eq!(reloaded.len(), 500, "no events lost");
for (i, event) in reloaded.events().iter().enumerate() {
assert_eq!(
event.action_id.as_deref(),
Some(format!("a{i}").as_str()),
"order preserved at {i}"
);
}
}
#[test]
fn unopenable_journal_is_best_effort_not_fatal() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("a-directory");
fs::create_dir(&journal).unwrap();
let mut log = EventLog::with_journal(journal);
log.append(
EventKind::ProposalReceived,
None,
Some("p1"),
HashMap::new(),
);
log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
assert_eq!(
log.len(),
2,
"in-memory log unaffected by an unwritable journal"
);
}
#[test]
fn load_then_append_preserves_existing_and_adds() {
let dir = tempfile::tempdir().unwrap();
let journal = dir.path().join("resume.jsonl");
{
let mut log = EventLog::with_journal(journal.clone());
log.append(
EventKind::ProposalReceived,
None,
Some("p1"),
HashMap::new(),
);
}
{
let mut log = EventLog::load(&journal).unwrap();
assert_eq!(log.len(), 1);
log.append(EventKind::ActionSucceeded, Some("a1"), None, HashMap::new());
}
let reloaded = EventLog::load(&journal).unwrap();
assert_eq!(reloaded.len(), 2, "append-mode preserved the loaded line");
assert_eq!(reloaded.events()[0].kind, EventKind::ProposalReceived);
assert_eq!(reloaded.events()[1].kind, EventKind::ActionSucceeded);
}
#[test]
fn event_kind_serializes_snake_case() {
assert_eq!(
serde_json::to_string(&EventKind::ProposalReceived).unwrap(),
"\"proposal_received\""
);
assert_eq!(
serde_json::to_string(&EventKind::StateSnapshot).unwrap(),
"\"state_snapshot\""
);
}
#[test]
fn stats_truncate_and_clear_release_retained_entries() {
let mut log = EventLog::new();
for idx in 0..5 {
log.append(
EventKind::ActionSucceeded,
Some(&format!("a{idx}")),
Some("p1"),
[("payload".to_string(), Value::from("x".repeat(16)))].into(),
);
log.begin_span("action.tool_call", "trace", None, HashMap::new());
}
let stats = log.stats();
assert_eq!(stats.events, 5);
assert_eq!(stats.spans, 5);
assert!(stats.approx_event_bytes > 0);
assert!(stats.approx_span_bytes > 0);
assert_eq!(log.truncate_events_keep_last(2), 3);
assert_eq!(log.truncate_spans_keep_last(1), 4);
assert_eq!(log.len(), 2);
assert_eq!(log.span_len(), 1);
assert_eq!(log.events()[0].action_id.as_deref(), Some("a3"));
let removed = log.clear();
assert_eq!(removed.events, 2);
assert_eq!(removed.spans, 1);
assert_eq!(log.len(), 0);
assert_eq!(log.span_len(), 0);
}
#[test]
fn span_begin_end_lifecycle() {
let mut log = EventLog::new();
let trace_id = "trace-1".to_string();
let span_id = log.begin_span(
"test.operation",
&trace_id,
None,
[("key".to_string(), Value::from("value"))].into(),
);
let spans = log.spans();
assert_eq!(spans.len(), 1);
assert_eq!(spans[0].name, "test.operation");
assert_eq!(spans[0].trace_id, "trace-1");
assert!(spans[0].parent_span_id.is_none());
assert!(spans[0].end_time.is_none());
assert_eq!(spans[0].status, SpanStatus::Unset);
log.end_span(&span_id, SpanStatus::Ok);
let spans = log.spans();
assert!(spans[0].end_time.is_some());
assert_eq!(spans[0].status, SpanStatus::Ok);
}
#[test]
fn span_parent_child_relationship() {
let mut log = EventLog::new();
let trace_id = "trace-2".to_string();
let parent_id = log.begin_span("parent.op", &trace_id, None, HashMap::new());
let child_id = log.begin_span("child.op", &trace_id, Some(&parent_id), HashMap::new());
let spans = log.spans();
assert_eq!(spans.len(), 2);
let child = spans.iter().find(|s| s.span_id == child_id).unwrap();
assert_eq!(child.parent_span_id.as_deref(), Some(parent_id.as_str()));
assert_eq!(child.trace_id, trace_id);
let parent = spans.iter().find(|s| s.span_id == parent_id).unwrap();
assert!(parent.parent_span_id.is_none());
}
#[test]
fn export_traces_produces_valid_json() {
let mut log = EventLog::new();
let trace_id = "trace-3".to_string();
let root = log.begin_span(
"proposal.execute",
&trace_id,
None,
[("proposal_id".to_string(), Value::from("p1"))].into(),
);
let child = log.begin_span(
"action.tool_call",
&trace_id,
Some(&root),
[("tool".to_string(), Value::from("read_file"))].into(),
);
log.end_span(&child, SpanStatus::Ok);
log.end_span(&root, SpanStatus::Ok);
let json_str = log.export_traces();
let parsed: Value =
serde_json::from_str(&json_str).expect("export_traces must produce valid JSON");
let resource_spans = parsed["resourceSpans"].as_array().unwrap();
assert_eq!(resource_spans.len(), 1);
let scope_spans = &resource_spans[0]["scopeSpans"][0]["spans"];
let spans_arr = scope_spans.as_array().unwrap();
assert_eq!(spans_arr.len(), 2);
for span in spans_arr {
assert!(span.get("traceId").is_some());
assert!(span.get("spanId").is_some());
assert!(span.get("name").is_some());
assert!(span.get("startTimeUnixNano").is_some());
assert!(span.get("endTimeUnixNano").is_some());
assert!(span.get("status").is_some());
}
let child_span = spans_arr
.iter()
.find(|s| s["name"] == "action.tool_call")
.unwrap();
assert!(child_span.get("parentSpanId").is_some());
}
#[test]
fn span_status_set_on_error() {
let mut log = EventLog::new();
let trace_id = "trace-4".to_string();
let span_id = log.begin_span("failing.op", &trace_id, None, HashMap::new());
log.end_span(&span_id, SpanStatus::Error);
let spans = log.spans();
assert_eq!(spans[0].status, SpanStatus::Error);
assert!(spans[0].end_time.is_some());
}
}