| 1 | //! `codewhale metrics` — reads the audit log and session/task stores and prints |
| 2 | //! a human-readable usage rollup. |
| 3 | //! |
| 4 | //! Data sources, all resolved through the shared Codewhale state resolvers so |
| 5 | //! the reader lands on the same files the writers use: |
| 6 | //! - `~/.codewhale/audit.log` — one JSON line per event (approvals, credentials) |
| 7 | //! - `~/.codewhale/sessions/` — saved session JSON files (tool call history) |
| 8 | //! - `~/.codewhale/tasks/runtime/events/` — runtime thread JSONL event streams |
| 9 | //! - `~/.codewhale/sessions/<id>/runtime/events/` — session-scoped runtime |
| 10 | //! stores (`default_runtime_store_root` in `crates/tui/src/runtime_threads.rs`) |
| 11 | //! |
| 12 | //! `CODEWHALE_RUNTIME_DIR` / `DEEPSEEK_RUNTIME_DIR` is an *exclusive* store root |
| 13 | //! for every Runtime store in the writing process, so when it is set the reader |
| 14 | //! uses it alone. Mixing it with the default roots would count one call twice. |
| 15 | //! |
| 16 | //! Default-root audit history includes retained rotations and legacy receipts, |
| 17 | //! excluding records copied across roots. An explicit `CODEWHALE_HOME` never |
| 18 | //! reads outside that root. |
| 19 | //! |
| 20 | //! The three sources overlap and are deliberately not de-duplicated against |
| 21 | //! each other: an approval receipt, a saved-session transcript entry, and a |
| 22 | //! runtime `item.*` receipt each describe one tool call from a different |
| 23 | //! vantage point, and collapsing them would assert an identity the data does |
| 24 | //! not carry. Cross-reference with `--json` when an exact count matters. |
| 25 | //! |
| 26 | //! There is no fourth source. The opt-in tool audit file behind |
| 27 | //! `CODEWHALE_TOOL_AUDIT_LOG` / `DEEPSEEK_TOOL_AUDIT_LOG` (`emit_tool_audit`) |
| 28 | //! is a *different* file from `~/.codewhale/audit.log` and is not discovered |
| 29 | //! here. Because that variable can be pointed at the audit log itself, and |
| 30 | //! because `emit_tool_audit` puts `tool_name` and `success` at the JSON top |
| 31 | //! level rather than under `details`, the audit reader accepts both shapes. |
| 32 | |
| 33 | use std::collections::{HashMap, HashSet}; |
| 34 | use std::path::{Path, PathBuf}; |
| 35 | |
| 36 | use anyhow::Result; |
| 37 | use chrono::{DateTime, Duration, Utc}; |
| 38 | use serde_json::Value; |
| 39 | use sha2::{Digest, Sha256}; |
| 40 | |
| 41 | // ────────────────────────────────────────────────────────────────────────────── |
| 42 | // Public entry-point |
| 43 | // ────────────────────────────────────────────────────────────────────────────── |
| 44 | |
| 45 | /// Arguments accepted by `codewhale metrics`. |
| 46 | #[derive(Debug, Default)] |
| 47 | pub struct MetricsArgs { |
| 48 | /// Emit machine-readable JSON instead of human text. |
| 49 | pub json: bool, |
| 50 | /// Restrict to events newer than this cutoff (inclusive). |
| 51 | pub since: Option<DateTime<Utc>>, |
| 52 | } |
| 53 | |
| 54 | pub fn run(args: MetricsArgs) -> Result<()> { |
| 55 | // `resolve_state_dir` is the shared read-path resolver already used by |
| 56 | // `doctor` and the session store; the runtime thread store hangs its event |
| 57 | // streams off `<tasks>/runtime`. Resolving the home is fallible, and a |
| 58 | // rollup of zeros is indistinguishable from real emptiness, so a home we |
| 59 | // cannot resolve is an error rather than a silent all-zero report. |
| 60 | let audit_roots = resolve_audit_roots()?; |
| 61 | let sessions = codewhale_config::resolve_state_dir("sessions")?; |
| 62 | let tasks = codewhale_config::resolve_state_dir("tasks")?; |
| 63 | let runtime_events = runtime_event_dirs(&tasks, &sessions); |
| 64 | |
| 65 | // Collect data from every source; treat missing files as empty. |
| 66 | let mut rollup = Rollup::default(); |
| 67 | read_audit_history(&audit_roots, args.since, &mut rollup); |
| 68 | read_session_files(&sessions, args.since, &mut rollup); |
| 69 | read_runtime_events(&runtime_events, args.since, &mut rollup); |
| 70 | |
| 71 | if args.json { |
| 72 | print_json(&rollup)?; |
| 73 | } else { |
| 74 | print_human(&rollup, args.since); |
| 75 | } |
| 76 | |
| 77 | Ok(()) |
| 78 | } |
| 79 | |
| 80 | // ────────────────────────────────────────────────────────────────────────────── |
| 81 | // Duration-string parser ("7d", "24h", "30m", "2h", "now-2h", "2h30m") |
| 82 | // ────────────────────────────────────────────────────────────────────────────── |
| 83 | |
| 84 | /// Parse a loose humantime-ish duration string into an absolute `DateTime<Utc>` |
| 85 | /// cutoff (i.e. `Utc::now() - duration`). |
| 86 | /// |
| 87 | /// Accepted forms: |
| 88 | /// - `7d` / `24h` / `30m` / `90s` |
| 89 | /// - `2h30m`, `1d12h` |
| 90 | /// - `now-2h` (leading `now-` is stripped before parsing) |
| 91 | pub fn parse_since(s: &str) -> Result<DateTime<Utc>> { |
| 92 | let s = s.trim().to_ascii_lowercase(); |
| 93 | let s = s.strip_prefix("now-").unwrap_or(&s); |
| 94 | let secs = parse_duration_secs(s)?; |
| 95 | let delta = Duration::try_seconds(secs) |
| 96 | .ok_or_else(|| anyhow::anyhow!("duration {s:?} is too large"))?; |
| 97 | Utc::now() |
| 98 | .checked_sub_signed(delta) |
| 99 | .ok_or_else(|| anyhow::anyhow!("duration {s:?} reaches before the earliest supported time")) |
| 100 | } |
| 101 | |
| 102 | fn parse_duration_secs(s: &str) -> Result<i64> { |
| 103 | // Walk through the string accumulating numbers and consuming unit suffixes. |
| 104 | let mut total: i64 = 0; |
| 105 | let mut num_buf = String::new(); |
| 106 | |
| 107 | for ch in s.chars() { |
| 108 | match ch { |
| 109 | '0'..='9' => num_buf.push(ch), |
| 110 | 'd' | 'h' | 'm' | 's' => { |
| 111 | if num_buf.is_empty() { |
| 112 | anyhow::bail!("unit {ch:?} in duration {s:?} has no number before it"); |
| 113 | } |
| 114 | let n: i64 = num_buf |
| 115 | .parse() |
| 116 | .map_err(|_| anyhow::anyhow!("duration component {num_buf:?} is too large"))?; |
| 117 | num_buf.clear(); |
| 118 | let factor = match ch { |
| 119 | 'd' => 86_400, |
| 120 | 'h' => 3_600, |
| 121 | 'm' => 60, |
| 122 | 's' => 1, |
| 123 | _ => unreachable!(), |
| 124 | }; |
| 125 | total = n |
| 126 | .checked_mul(factor) |
| 127 | .and_then(|secs| total.checked_add(secs)) |
| 128 | .ok_or_else(|| anyhow::anyhow!("duration {s:?} is too large"))?; |
| 129 | } |
| 130 | _ => anyhow::bail!("unrecognised character {ch:?} in duration {s:?}"), |
| 131 | } |
| 132 | } |
| 133 | |
| 134 | if !num_buf.is_empty() { |
| 135 | // Trailing bare number — treat as seconds. |
| 136 | let n: i64 = num_buf |
| 137 | .parse() |
| 138 | .map_err(|_| anyhow::anyhow!("duration component {num_buf:?} is too large"))?; |
| 139 | total = total |
| 140 | .checked_add(n) |
| 141 | .ok_or_else(|| anyhow::anyhow!("duration {s:?} is too large"))?; |
| 142 | } |
| 143 | |
| 144 | if total == 0 { |
| 145 | anyhow::bail!("duration {s:?} resolved to zero seconds"); |
| 146 | } |
| 147 | |
| 148 | Ok(total) |
| 149 | } |
| 150 | |
| 151 | // ────────────────────────────────────────────────────────────────────────────── |
| 152 | // Rollup data model |
| 153 | // ────────────────────────────────────────────────────────────────────────────── |
| 154 | |
| 155 | /// Per-tool aggregated counters. |
| 156 | #[derive(Debug, Default, serde::Serialize)] |
| 157 | pub struct ToolStats { |
| 158 | pub calls: u64, |
| 159 | /// Calls that were auto-approved (no prompt required). |
| 160 | pub auto_approved: u64, |
| 161 | /// Calls that required a manual prompt. |
| 162 | pub prompted: u64, |
| 163 | /// Total elapsed ms (from events that carry this field). |
| 164 | pub total_elapsed_ms: u64, |
| 165 | /// Number of elapsed_ms samples included in `total_elapsed_ms`. |
| 166 | pub elapsed_samples: u64, |
| 167 | /// Successful calls (where we have result data). |
| 168 | pub successes: u64, |
| 169 | /// Failed calls. |
| 170 | pub failures: u64, |
| 171 | /// Calls an approval receipt blocked before they ran. A denial is neither |
| 172 | /// a success nor a failure, so it stays out of `success_rate_pct`. |
| 173 | pub denied: u64, |
| 174 | /// Durable receipts whose outcome could not be read. Never folded into |
| 175 | /// `failures`: an unrecorded outcome is unknown, not a failure. |
| 176 | pub outcome_unknown: u64, |
| 177 | /// Terminal receipts with no usable `started_at`/`ended_at` pair. Counted |
| 178 | /// rather than contributing a 0 ms sample, which would understate a call |
| 179 | /// that was actually slow. |
| 180 | pub elapsed_unavailable: u64, |
| 181 | } |
| 182 | |
| 183 | impl ToolStats { |
| 184 | fn success_rate_pct(&self) -> Option<f64> { |
| 185 | let judged = self.successes + self.failures; |
| 186 | if judged == 0 { |
| 187 | None |
| 188 | } else { |
| 189 | Some(self.successes as f64 / judged as f64 * 100.0) |
| 190 | } |
| 191 | } |
| 192 | |
| 193 | fn avg_elapsed_ms(&self) -> Option<u64> { |
| 194 | self.total_elapsed_ms.checked_div(self.elapsed_samples) |
| 195 | } |
| 196 | } |
| 197 | |
| 198 | /// Compaction event stats. |
| 199 | #[derive(Debug, Default, serde::Serialize)] |
| 200 | pub struct CompactionStats { |
| 201 | pub events: u64, |
| 202 | pub refusals: HashMap<String, u64>, |
| 203 | pub triggers: HashMap<String, u64>, |
| 204 | pub paths: HashMap<String, u64>, |
| 205 | pub summarizer_usage_samples: u64, |
| 206 | pub summarizer_input_tokens: u64, |
| 207 | pub summarizer_output_tokens: u64, |
| 208 | #[serde(skip)] |
| 209 | receipt_ids: HashSet<String>, |
| 210 | /// Sum of `reduction_ratio` from events that carry it (0.0–1.0 each). |
| 211 | pub ratio_sum: f64, |
| 212 | pub ratio_samples: u64, |
| 213 | } |
| 214 | |
| 215 | impl CompactionStats { |
| 216 | fn avg_reduction_pct(&self) -> Option<f64> { |
| 217 | if self.ratio_samples == 0 { |
| 218 | None |
| 219 | } else { |
| 220 | Some(self.ratio_sum / self.ratio_samples as f64 * 100.0) |
| 221 | } |
| 222 | } |
| 223 | } |
| 224 | |
| 225 | /// Sub-agent lifecycle receipt counts; these are not unique worker totals. |
| 226 | #[derive(Debug, Default, serde::Serialize)] |
| 227 | pub struct AgentStats { |
| 228 | pub spawns: u64, |
| 229 | pub successes: u64, |
| 230 | pub failures: u64, |
| 231 | pub cancelled: u64, |
| 232 | pub interrupted: u64, |
| 233 | pub budget_exhausted: u64, |
| 234 | /// Terminal receipts with missing, malformed, or unrecognized outcomes. |
| 235 | pub unknown_outcomes: u64, |
| 236 | /// Completions carrying a usage receipt. Token sums cover exactly these; |
| 237 | /// a completion without usage is a missing receipt, never zero tokens. |
| 238 | pub usage_receipts: u64, |
| 239 | pub input_tokens: u64, |
| 240 | pub output_tokens: u64, |
| 241 | pub total_tokens: u64, |
| 242 | /// Summed priced subtotal in microdollars, JSON consumers only. |
| 243 | pub cost_microusd: u64, |
| 244 | } |
| 245 | |
| 246 | impl AgentStats { |
| 247 | fn record_completion(&mut self, event: &Value) { |
| 248 | // Runtime's worker_status owns the outcome. A completed status item |
| 249 | // means its receipt settled, not that the worker succeeded. Preserve |
| 250 | // explicit unknown values instead of falling back to a legacy boolean. |
| 251 | let status = event |
| 252 | .pointer("/details/worker_status") |
| 253 | .or_else(|| event.pointer("/payload/worker_status")) |
| 254 | .or_else(|| event.pointer("/details/status")) |
| 255 | .or_else(|| event.pointer("/payload/status")); |
| 256 | let count = if let Some(status) = status { |
| 257 | match status.as_str() { |
| 258 | Some("completed") => &mut self.successes, |
| 259 | Some("failed") => &mut self.failures, |
| 260 | Some("cancelled") => &mut self.cancelled, |
| 261 | Some("interrupted") => &mut self.interrupted, |
| 262 | Some("budget_exhausted") => &mut self.budget_exhausted, |
| 263 | _ => &mut self.unknown_outcomes, |
| 264 | } |
| 265 | } else { |
| 266 | match event |
| 267 | .pointer("/details/success") |
| 268 | .or_else(|| event.pointer("/payload/success")) |
| 269 | .and_then(Value::as_bool) |
| 270 | { |
| 271 | Some(true) => &mut self.successes, |
| 272 | Some(false) => &mut self.failures, |
| 273 | None => &mut self.unknown_outcomes, |
| 274 | } |
| 275 | }; |
| 276 | *count = count.saturating_add(1); |
| 277 | // Child cost visibility (#6315): the runtime persists the completion |
| 278 | // receipt's usage on the agent.completed payload. |
| 279 | let usage = event |
| 280 | .pointer("/payload/usage") |
| 281 | .or_else(|| event.pointer("/details/usage")); |
| 282 | if let Some(usage) = usage |
| 283 | && usage.is_object() |
| 284 | { |
| 285 | self.usage_receipts = self.usage_receipts.saturating_add(1); |
| 286 | for (field, sum) in [ |
| 287 | ("input_tokens", &mut self.input_tokens), |
| 288 | ("output_tokens", &mut self.output_tokens), |
| 289 | ("total_tokens", &mut self.total_tokens), |
| 290 | ("cost_microusd", &mut self.cost_microusd), |
| 291 | ] { |
| 292 | if let Some(n) = usage.get(field).and_then(Value::as_u64) { |
| 293 | *sum = (*sum).saturating_add(n); |
| 294 | } |
| 295 | } |
| 296 | } |
| 297 | } |
| 298 | |
| 299 | fn summary(&self) -> String { |
| 300 | let outcomes = [ |
| 301 | (self.successes, "completed"), |
| 302 | (self.failures, "failed"), |
| 303 | (self.cancelled, "cancelled"), |
| 304 | (self.interrupted, "interrupted"), |
| 305 | (self.budget_exhausted, "budget exhausted"), |
| 306 | (self.unknown_outcomes, "outcome unconfirmed"), |
| 307 | ] |
| 308 | .into_iter() |
| 309 | .filter(|(count, _)| *count > 0) |
| 310 | .map(|(count, label)| format!("{} {label}", fmt_num(count))) |
| 311 | .collect::<Vec<_>>(); |
| 312 | if self.spawns == 0 && outcomes.is_empty() { |
| 313 | return "Sub-agents: (no data)".to_string(); |
| 314 | } |
| 315 | let mut summary = format!("Sub-agents: {} spawn receipts", fmt_num(self.spawns)); |
| 316 | if !outcomes.is_empty() { |
| 317 | summary.push_str(&format!("; outcomes: {}", outcomes.join(", "))); |
| 318 | } |
| 319 | if self.usage_receipts > 0 { |
| 320 | summary.push_str(&format!( |
| 321 | "; tokens: {} in/{} out/{} total ({} {})", |
| 322 | fmt_num(self.input_tokens), |
| 323 | fmt_num(self.output_tokens), |
| 324 | fmt_num(self.total_tokens), |
| 325 | fmt_num(self.usage_receipts), |
| 326 | if self.usage_receipts == 1 { |
| 327 | "receipt" |
| 328 | } else { |
| 329 | "receipts" |
| 330 | }, |
| 331 | )); |
| 332 | } |
| 333 | summary |
| 334 | } |
| 335 | } |
| 336 | |
| 337 | /// Capacity-controller / rate-limit intervention stats. |
| 338 | #[derive(Debug, Default, serde::Serialize)] |
| 339 | pub struct CapacityStats { |
| 340 | pub total: u64, |
| 341 | pub by_category: HashMap<String, u64>, |
| 342 | } |
| 343 | |
| 344 | /// Credential / session event stats (from audit log). |
| 345 | #[derive(Debug, Default, serde::Serialize)] |
| 346 | pub struct CredentialStats { |
| 347 | pub saves: u64, |
| 348 | pub clears: u64, |
| 349 | } |
| 350 | |
| 351 | /// Runtime receipts for model-client dispatch and provider-reported usage. |
| 352 | /// |
| 353 | /// These are deliberately not billing records: the terminal diagnostics count |
| 354 | /// parent model-client calls, while `turn.usage` exists only when a provider |
| 355 | /// supplied usage for one call. Client-internal HTTP retries and invoices are |
| 356 | /// outside both receipts. |
| 357 | #[derive(Debug, Default, serde::Serialize)] |
| 358 | pub struct RuntimeRequestStats { |
| 359 | /// Distinct durable `turn.completed` receipts with a usable `(thread, turn)` identity. |
| 360 | pub terminal_turn_receipts: u64, |
| 361 | /// Terminal receipts carrying the optional request diagnostics projection. |
| 362 | pub diagnostics_turn_receipts: u64, |
| 363 | /// Terminal receipts from older or partial logs with no diagnostics projection. |
| 364 | pub diagnostics_unavailable_turn_receipts: u64, |
| 365 | /// Present-but-incomplete diagnostics are unknown rather than zero. |
| 366 | pub diagnostics_incomplete_turn_receipts: u64, |
| 367 | /// Terminal receipts omitted because their identity could not be verified. |
| 368 | pub terminal_receipts_without_identity: u64, |
| 369 | /// Repeated terminal snapshots for one `(thread, turn)` omitted from the rollup. |
| 370 | pub duplicate_terminal_receipts_skipped: u64, |
| 371 | /// Parent model-client calls recorded by terminal diagnostics, not HTTP retries or invoices. |
| 372 | pub model_requests_started: u64, |
| 373 | pub transparent_stream_retries: u64, |
| 374 | pub stream_resumes: u64, |
| 375 | /// Distinct `turn.usage` receipts with a verified runtime event identity. |
| 376 | pub provider_usage_receipts: u64, |
| 377 | /// `turn.usage` records that could not be identified, so their values are unknown. |
| 378 | pub provider_usage_receipts_without_identity: u64, |
| 379 | /// Repeated runtime event identities omitted from provider usage totals. |
| 380 | pub duplicate_provider_usage_receipts_skipped: u64, |
| 381 | /// `turn.usage` records missing either required token total are not treated as zero. |
| 382 | pub provider_usage_receipts_incomplete: u64, |
| 383 | /// Provider-reported per-request input tokens only; terminal cumulative snapshots are excluded. |
| 384 | pub provider_reported_input_tokens: u64, |
| 385 | /// Provider-reported per-request output tokens only; terminal cumulative snapshots are excluded. |
| 386 | pub provider_reported_output_tokens: u64, |
| 387 | } |
| 388 | |
| 389 | /// Top-level rollup. |
| 390 | #[derive(Debug, Default, serde::Serialize)] |
| 391 | pub struct Rollup { |
| 392 | /// UTC timestamp of the earliest event we've seen. |
| 393 | pub earliest_ts: Option<DateTime<Utc>>, |
| 394 | /// UTC timestamp of the latest event we've seen. |
| 395 | pub latest_ts: Option<DateTime<Utc>>, |
| 396 | /// Per-tool stats keyed by tool name. |
| 397 | pub tools: HashMap<String, ToolStats>, |
| 398 | pub compaction: CompactionStats, |
| 399 | pub agents: AgentStats, |
| 400 | pub capacity: CapacityStats, |
| 401 | pub credentials: CredentialStats, |
| 402 | pub runtime_requests: RuntimeRequestStats, |
| 403 | /// Total lines read across all sources. |
| 404 | pub total_lines: u64, |
| 405 | /// Lines successfully parsed. |
| 406 | pub parsed_lines: u64, |
| 407 | } |
| 408 | |
| 409 | #[derive(Default)] |
| 410 | struct RuntimeEventDedup { |
| 411 | terminal_turns: HashSet<(String, String)>, |
| 412 | event_records: HashSet<(String, u64)>, |
| 413 | } |
| 414 | |
| 415 | impl Rollup { |
| 416 | fn touch_ts(&mut self, ts: &DateTime<Utc>) { |
| 417 | match self.earliest_ts { |
| 418 | None => self.earliest_ts = Some(*ts), |
| 419 | Some(ref cur) if ts < cur => self.earliest_ts = Some(*ts), |
| 420 | _ => {} |
| 421 | } |
| 422 | match self.latest_ts { |
| 423 | None => self.latest_ts = Some(*ts), |
| 424 | Some(ref cur) if ts > cur => self.latest_ts = Some(*ts), |
| 425 | _ => {} |
| 426 | } |
| 427 | } |
| 428 | |
| 429 | fn tool_mut(&mut self, name: &str) -> &mut ToolStats { |
| 430 | self.tools.entry(name.to_string()).or_default() |
| 431 | } |
| 432 | |
| 433 | fn total_tool_calls(&self) -> u64 { |
| 434 | self.tools.values().map(|t| t.calls).sum() |
| 435 | } |
| 436 | } |
| 437 | |
| 438 | // ────────────────────────────────────────────────────────────────────────────── |
| 439 | // Source readers |
| 440 | // ────────────────────────────────────────────────────────────────────────────── |
| 441 | |
| 442 | /// Read both retained generations from each root. A copied legacy record is |
| 443 | /// counted once across roots, while repeated records within one root retain |
| 444 | /// their multiplicity. No source log is rewritten or removed. |
| 445 | fn read_audit_history(roots: &[PathBuf], since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 446 | let mut earlier_roots = HashMap::new(); |
| 447 | for root in roots { |
| 448 | let mut root_counts = HashMap::new(); |
| 449 | for name in ["audit.log.1", "audit.log"] { |
| 450 | read_audit_log( |
| 451 | &root.join(name), |
| 452 | since, |
| 453 | rollup, |
| 454 | &earlier_roots, |
| 455 | &mut root_counts, |
| 456 | ); |
| 457 | } |
| 458 | for (record, count) in root_counts { |
| 459 | let prior = earlier_roots.entry(record).or_insert(0); |
| 460 | *prior = (*prior).max(count); |
| 461 | } |
| 462 | } |
| 463 | } |
| 464 | |
| 465 | /// Read one JSON event per line, excluding copies already seen in other roots. |
| 466 | fn read_audit_log( |
| 467 | path: &Path, |
| 468 | since: Option<DateTime<Utc>>, |
| 469 | rollup: &mut Rollup, |
| 470 | earlier_roots: &HashMap<[u8; 32], u64>, |
| 471 | root_counts: &mut HashMap<[u8; 32], u64>, |
| 472 | ) { |
| 473 | let file = match std::fs::File::open(path) { |
| 474 | Ok(file) => file, |
| 475 | Err(e) if e.kind() == std::io::ErrorKind::NotFound => return, |
| 476 | Err(e) => { |
| 477 | tracing::trace!( |
| 478 | "metrics: could not read audit log {}: {}", |
| 479 | path.display(), |
| 480 | e |
| 481 | ); |
| 482 | return; |
| 483 | } |
| 484 | }; |
| 485 | |
| 486 | // Streamed line by line: an append-only log is unbounded, and holding it |
| 487 | // whole only to aggregate one line at a time made allocation scale with it. |
| 488 | for raw_line in std::io::BufRead::lines(std::io::BufReader::new(file)) { |
| 489 | let Ok(raw_line) = raw_line.inspect_err(|e| { |
| 490 | tracing::trace!("metrics: stopped reading audit log {}: {e}", path.display()); |
| 491 | }) else { |
| 492 | break; |
| 493 | }; |
| 494 | rollup.total_lines += 1; |
| 495 | let line = raw_line.trim(); |
| 496 | if line.is_empty() { |
| 497 | continue; |
| 498 | } |
| 499 | |
| 500 | let v: Value = match serde_json::from_str(line) { |
| 501 | Ok(v) => v, |
| 502 | Err(e) => { |
| 503 | tracing::trace!("metrics: skipping malformed audit line: {e}"); |
| 504 | continue; |
| 505 | } |
| 506 | }; |
| 507 | |
| 508 | // Copy migration preserves the complete event, including its timestamp. |
| 509 | // Count occurrences so two identical legitimate records in one source |
| 510 | // are not collapsed into one merely because another root also exists. |
| 511 | let fingerprint: [u8; 32] = Sha256::digest(v.to_string().as_bytes()).into(); |
| 512 | let count = root_counts.entry(fingerprint).or_insert(0); |
| 513 | *count += 1; |
| 514 | if *count <= earlier_roots.get(&fingerprint).copied().unwrap_or(0) { |
| 515 | continue; |
| 516 | } |
| 517 | |
| 518 | // Parse timestamp — field is "ts" in audit log. |
| 519 | let ts = parse_ts_field(&v, "ts"); |
| 520 | |
| 521 | if let Some(cutoff) = since { |
| 522 | match ts { |
| 523 | Some(t) if t < cutoff => continue, |
| 524 | _ => {} |
| 525 | } |
| 526 | } |
| 527 | |
| 528 | rollup.parsed_lines += 1; |
| 529 | if let Some(t) = &ts { |
| 530 | rollup.touch_ts(t); |
| 531 | } |
| 532 | |
| 533 | let event = v.get("event").and_then(|e| e.as_str()).unwrap_or(""); |
| 534 | |
| 535 | match event { |
| 536 | // `log_sensitive_event` emits `auto_approve_session` |
| 537 | // (`crates/tui/src/tui/ui/event_loop.rs`). The bare `auto_approve` |
| 538 | // name only ever appears in older logs; keep it as an alias. |
| 539 | "tool.approval.auto_approve" | "tool.approval.auto_approve_session" => { |
| 540 | let tool_name = audit_tool_name(&v); |
| 541 | let stats = rollup.tool_mut(tool_name); |
| 542 | stats.calls += 1; |
| 543 | stats.auto_approved += 1; |
| 544 | } |
| 545 | // Every denial name written by `log_sensitive_event` and |
| 546 | // `auto_deny_session_approval`. The call never ran, so it is |
| 547 | // counted as its own class rather than as a failed execution. |
| 548 | "tool.approval.auto_deny" |
| 549 | | "tool.approval.auto_deny_session" |
| 550 | | "tool.approval.auto_deny_auto_review" |
| 551 | | "tool.approval.auto_deny_full_access_policy" => { |
| 552 | let tool_name = audit_tool_name(&v); |
| 553 | let stats = rollup.tool_mut(tool_name); |
| 554 | stats.calls += 1; |
| 555 | stats.denied += 1; |
| 556 | } |
| 557 | "tool.approval.prompted" => { |
| 558 | let tool_name = audit_tool_name(&v); |
| 559 | let stats = rollup.tool_mut(tool_name); |
| 560 | stats.calls += 1; |
| 561 | stats.prompted += 1; |
| 562 | } |
| 563 | "tool.completed" | "tool.result" => { |
| 564 | let tool_name = audit_tool_name(&v); |
| 565 | let stats = rollup.tool_mut(tool_name); |
| 566 | stats.calls += 1; |
| 567 | |
| 568 | // Optional elapsed_ms |
| 569 | if let Some(ms) = v |
| 570 | .pointer("/details/elapsed_ms") |
| 571 | .or_else(|| v.pointer("/payload/elapsed_ms")) |
| 572 | .or_else(|| v.get("elapsed_ms")) |
| 573 | .and_then(|v| v.as_u64()) |
| 574 | { |
| 575 | stats.total_elapsed_ms += ms; |
| 576 | stats.elapsed_samples += 1; |
| 577 | } |
| 578 | |
| 579 | // Success / failure. An absent outcome is unknown, not a |
| 580 | // success — the previous default silently graded every |
| 581 | // outcome-free receipt as passing. |
| 582 | match v |
| 583 | .pointer("/details/success") |
| 584 | .or_else(|| v.pointer("/payload/success")) |
| 585 | .or_else(|| v.get("success")) |
| 586 | .and_then(|b| b.as_bool()) |
| 587 | { |
| 588 | Some(true) => stats.successes += 1, |
| 589 | Some(false) => stats.failures += 1, |
| 590 | None => stats.outcome_unknown += 1, |
| 591 | } |
| 592 | } |
| 593 | "compaction.refused" => { |
| 594 | let reason = v |
| 595 | .pointer("/details/reason") |
| 596 | .and_then(Value::as_str) |
| 597 | .unwrap_or("unknown"); |
| 598 | *rollup |
| 599 | .compaction |
| 600 | .refusals |
| 601 | .entry(reason.to_string()) |
| 602 | .or_default() += 1; |
| 603 | } |
| 604 | "compaction.completed" | "context.compaction" => { |
| 605 | if let Some(id) = compaction_receipt_identity(&v, false) |
| 606 | && !rollup.compaction.receipt_ids.insert(id) |
| 607 | { |
| 608 | continue; |
| 609 | } |
| 610 | for (field, counts) in [ |
| 611 | ("trigger", &mut rollup.compaction.triggers), |
| 612 | ("path", &mut rollup.compaction.paths), |
| 613 | ] { |
| 614 | if let Some(value) = v |
| 615 | .pointer(&format!("/details/{field}")) |
| 616 | .and_then(Value::as_str) |
| 617 | { |
| 618 | *counts.entry(value.to_string()).or_default() += 1; |
| 619 | } |
| 620 | } |
| 621 | if let (Some(input), Some(output)) = ( |
| 622 | v.pointer("/details/summarizer_usage/input_tokens") |
| 623 | .and_then(Value::as_u64), |
| 624 | v.pointer("/details/summarizer_usage/output_tokens") |
| 625 | .and_then(Value::as_u64), |
| 626 | ) { |
| 627 | rollup.compaction.summarizer_usage_samples += 1; |
| 628 | rollup.compaction.summarizer_input_tokens += input; |
| 629 | rollup.compaction.summarizer_output_tokens += output; |
| 630 | } |
| 631 | rollup.compaction.events += 1; |
| 632 | if let Some(ratio) = v |
| 633 | .pointer("/details/reduction_ratio") |
| 634 | .or_else(|| v.pointer("/payload/reduction_ratio")) |
| 635 | .and_then(|r| r.as_f64()) |
| 636 | { |
| 637 | rollup.compaction.ratio_sum += ratio; |
| 638 | rollup.compaction.ratio_samples += 1; |
| 639 | } |
| 640 | } |
| 641 | "agent.spawn" | "agent.spawned" | "subagent.spawned" => { |
| 642 | rollup.agents.spawns += 1; |
| 643 | } |
| 644 | "agent.completed" | "subagent.completed" => { |
| 645 | rollup.agents.record_completion(&v); |
| 646 | } |
| 647 | e if e.starts_with("capacity.") => { |
| 648 | rollup.capacity.total += 1; |
| 649 | let category = v |
| 650 | .pointer("/details/category") |
| 651 | .or_else(|| v.pointer("/payload/category")) |
| 652 | .and_then(|c| c.as_str()) |
| 653 | .unwrap_or(e.trim_start_matches("capacity.")); |
| 654 | *rollup |
| 655 | .capacity |
| 656 | .by_category |
| 657 | .entry(category.to_string()) |
| 658 | .or_insert(0) += 1; |
| 659 | } |
| 660 | "credential.save" => { |
| 661 | rollup.credentials.saves += 1; |
| 662 | } |
| 663 | "credential.clear" => { |
| 664 | rollup.credentials.clears += 1; |
| 665 | } |
| 666 | _ => { |
| 667 | // Unknown event — tracked in parsed_lines but otherwise ignored. |
| 668 | } |
| 669 | } |
| 670 | } |
| 671 | } |
| 672 | |
| 673 | /// Read session JSON files under `sessions/` (one per session). |
| 674 | /// These carry tool call history with optional elapsed_ms and result data. |
| 675 | fn read_session_files(sessions_dir: &Path, since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 676 | let rd = match std::fs::read_dir(sessions_dir) { |
| 677 | Ok(rd) => rd, |
| 678 | Err(e) if e.kind() == std::io::ErrorKind::NotFound => return, |
| 679 | Err(e) => { |
| 680 | tracing::trace!( |
| 681 | "metrics: could not list sessions dir {}: {}", |
| 682 | sessions_dir.display(), |
| 683 | e |
| 684 | ); |
| 685 | return; |
| 686 | } |
| 687 | }; |
| 688 | |
| 689 | for entry in rd.flatten() { |
| 690 | let path = entry.path(); |
| 691 | // Only look at .json files directly in sessions/; skip sub-dirs. |
| 692 | if path.is_dir() || path.extension().map(|e| e != "json").unwrap_or(true) { |
| 693 | continue; |
| 694 | } |
| 695 | read_session_file(&path, since, rollup); |
| 696 | } |
| 697 | } |
| 698 | |
| 699 | fn read_session_file(path: &Path, since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 700 | let content = match std::fs::read_to_string(path) { |
| 701 | Ok(c) => c, |
| 702 | Err(e) => { |
| 703 | tracing::trace!( |
| 704 | "metrics: could not read session file {}: {}", |
| 705 | path.display(), |
| 706 | e |
| 707 | ); |
| 708 | return; |
| 709 | } |
| 710 | }; |
| 711 | |
| 712 | rollup.total_lines += 1; |
| 713 | |
| 714 | let v: Value = match serde_json::from_str(&content) { |
| 715 | Ok(v) => v, |
| 716 | Err(e) => { |
| 717 | tracing::trace!( |
| 718 | "metrics: skipping malformed session file {}: {}", |
| 719 | path.display(), |
| 720 | e |
| 721 | ); |
| 722 | return; |
| 723 | } |
| 724 | }; |
| 725 | |
| 726 | rollup.parsed_lines += 1; |
| 727 | |
| 728 | // Session-level timestamp filter (check metadata.created_at or updated_at). |
| 729 | let session_ts = v |
| 730 | .pointer("/metadata/updated_at") |
| 731 | .or_else(|| v.pointer("/metadata/created_at")) |
| 732 | .and_then(|t| t.as_str()) |
| 733 | .and_then(|s| s.parse::<DateTime<Utc>>().ok()); |
| 734 | |
| 735 | if let Some(cutoff) = since |
| 736 | && let Some(ts) = &session_ts |
| 737 | && *ts < cutoff |
| 738 | { |
| 739 | return; |
| 740 | } |
| 741 | |
| 742 | if let Some(ts) = session_ts { |
| 743 | rollup.touch_ts(&ts); |
| 744 | } |
| 745 | |
| 746 | // Walk messages looking for tool_use calls with associated results. |
| 747 | let messages = match v.get("messages").and_then(|m| m.as_array()) { |
| 748 | Some(m) => m, |
| 749 | None => return, |
| 750 | }; |
| 751 | |
| 752 | // Build a map from tool_use_id → (tool_name, elapsed_ms_option, started_at_option). |
| 753 | let mut pending: HashMap<String, (String, Option<u64>)> = HashMap::new(); |
| 754 | |
| 755 | for msg in messages { |
| 756 | let role = msg.get("role").and_then(|r| r.as_str()).unwrap_or(""); |
| 757 | let content_arr = match msg.get("content").and_then(|c| c.as_array()) { |
| 758 | Some(c) => c, |
| 759 | None => continue, |
| 760 | }; |
| 761 | |
| 762 | for block in content_arr { |
| 763 | let block_type = block.get("type").and_then(|t| t.as_str()).unwrap_or(""); |
| 764 | match (role, block_type) { |
| 765 | ("assistant", "tool_use") => { |
| 766 | let id = block.get("id").and_then(|i| i.as_str()).unwrap_or(""); |
| 767 | let name = block |
| 768 | .get("name") |
| 769 | .and_then(|n| n.as_str()) |
| 770 | .unwrap_or("unknown"); |
| 771 | let elapsed_ms = block.get("elapsed_ms").and_then(|e| e.as_u64()); |
| 772 | if !id.is_empty() { |
| 773 | pending.insert(id.to_string(), (name.to_string(), elapsed_ms)); |
| 774 | } |
| 775 | } |
| 776 | ("user", "tool_result") => { |
| 777 | let id = block |
| 778 | .get("tool_use_id") |
| 779 | .and_then(|i| i.as_str()) |
| 780 | .unwrap_or(""); |
| 781 | if let Some((name, elapsed_ms)) = pending.remove(id) { |
| 782 | let stats = rollup.tool_mut(&name); |
| 783 | // Only count if not already counted via audit log (we don't de-dup, so |
| 784 | // session files may double-count approvals; that's acceptable — users who |
| 785 | // want precise counts should use --json and cross-reference). |
| 786 | stats.calls += 1; |
| 787 | if let Some(ms) = elapsed_ms { |
| 788 | stats.total_elapsed_ms += ms; |
| 789 | stats.elapsed_samples += 1; |
| 790 | } |
| 791 | // Tool result success: absence of "is_error": true |
| 792 | let is_error = block |
| 793 | .get("is_error") |
| 794 | .and_then(|e| e.as_bool()) |
| 795 | .unwrap_or(false); |
| 796 | if is_error { |
| 797 | stats.failures += 1; |
| 798 | } else { |
| 799 | stats.successes += 1; |
| 800 | } |
| 801 | } |
| 802 | } |
| 803 | _ => {} |
| 804 | } |
| 805 | } |
| 806 | } |
| 807 | |
| 808 | // Walk messages for compaction events embedded as special user messages. |
| 809 | for msg in messages { |
| 810 | if let Some(compaction) = msg |
| 811 | .get("compaction") |
| 812 | .or_else(|| msg.pointer("/metadata/compaction")) |
| 813 | { |
| 814 | rollup.compaction.events += 1; |
| 815 | if let Some(ratio) = compaction.get("reduction_ratio").and_then(|r| r.as_f64()) { |
| 816 | rollup.compaction.ratio_sum += ratio; |
| 817 | rollup.compaction.ratio_samples += 1; |
| 818 | } |
| 819 | } |
| 820 | } |
| 821 | } |
| 822 | |
| 823 | /// Every Runtime event directory this install can have written to. |
| 824 | /// |
| 825 | /// Mirrors `default_runtime_store_root` / `runtime_dir_override` in |
| 826 | /// `crates/tui/src/runtime_threads.rs`: the task-scoped store lives at |
| 827 | /// `<tasks>/runtime`, a session-scoped store at `<sessions>/<id>/runtime`, and |
| 828 | /// an explicit `CODEWHALE_RUNTIME_DIR` replaces both. Missing directories are |
| 829 | /// simply empty; the walk stays inside the resolved state roots. |
| 830 | fn runtime_event_dirs(tasks: &Path, sessions: &Path) -> Vec<PathBuf> { |
| 831 | if let Some(root) = runtime_dir_override() { |
| 832 | return vec![root.join("events")]; |
| 833 | } |
| 834 | let mut dirs = vec![tasks.join("runtime").join("events")]; |
| 835 | let Ok(rd) = std::fs::read_dir(sessions) else { |
| 836 | return dirs; |
| 837 | }; |
| 838 | let mut session_dirs: Vec<PathBuf> = rd |
| 839 | .flatten() |
| 840 | // `DirEntry::file_type` does not follow symlinks, so a link planted in |
| 841 | // `sessions/` cannot walk the reader into another root. |
| 842 | .filter(|entry| entry.file_type().is_ok_and(|ty| ty.is_dir())) |
| 843 | .map(|entry| entry.path().join("runtime").join("events")) |
| 844 | .collect(); |
| 845 | session_dirs.sort(); |
| 846 | dirs.append(&mut session_dirs); |
| 847 | dirs |
| 848 | } |
| 849 | |
| 850 | /// The writer's exclusive store-root override (`runtime_dir_override`). |
| 851 | fn runtime_dir_override() -> Option<PathBuf> { |
| 852 | std::env::var("CODEWHALE_RUNTIME_DIR") |
| 853 | .or_else(|_| std::env::var("DEEPSEEK_RUNTIME_DIR")) |
| 854 | .ok() |
| 855 | .filter(|dir| !dir.trim().is_empty()) |
| 856 | .map(PathBuf::from) |
| 857 | } |
| 858 | |
| 859 | /// Read every runtime event root under one de-duplication scope, so the same |
| 860 | /// `(thread_id, seq)` receipt is counted once however many roots list it. |
| 861 | fn read_runtime_events(events_dirs: &[PathBuf], since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 862 | let mut dedup = RuntimeEventDedup::default(); |
| 863 | for dir in events_dirs { |
| 864 | read_runtime_events_dir(dir, since, rollup, &mut dedup); |
| 865 | } |
| 866 | } |
| 867 | |
| 868 | /// Read JSONL event streams from one runtime events directory. |
| 869 | fn read_runtime_events_dir( |
| 870 | events_dir: &Path, |
| 871 | since: Option<DateTime<Utc>>, |
| 872 | rollup: &mut Rollup, |
| 873 | dedup: &mut RuntimeEventDedup, |
| 874 | ) { |
| 875 | let rd = match std::fs::read_dir(events_dir) { |
| 876 | Ok(rd) => rd, |
| 877 | Err(e) if e.kind() == std::io::ErrorKind::NotFound => return, |
| 878 | Err(e) => { |
| 879 | tracing::trace!( |
| 880 | "metrics: could not list events dir {}: {}", |
| 881 | events_dir.display(), |
| 882 | e |
| 883 | ); |
| 884 | return; |
| 885 | } |
| 886 | }; |
| 887 | |
| 888 | for entry in rd.flatten() { |
| 889 | let path = entry.path(); |
| 890 | if path.extension().map(|e| e != "jsonl").unwrap_or(true) { |
| 891 | continue; |
| 892 | } |
| 893 | read_events_jsonl(&path, since, rollup, dedup); |
| 894 | } |
| 895 | } |
| 896 | |
| 897 | fn read_events_jsonl( |
| 898 | path: &Path, |
| 899 | since: Option<DateTime<Utc>>, |
| 900 | rollup: &mut Rollup, |
| 901 | dedup: &mut RuntimeEventDedup, |
| 902 | ) { |
| 903 | let file = match std::fs::File::open(path) { |
| 904 | Ok(file) => file, |
| 905 | Err(e) => { |
| 906 | tracing::trace!( |
| 907 | "metrics: could not read events file {}: {}", |
| 908 | path.display(), |
| 909 | e |
| 910 | ); |
| 911 | return; |
| 912 | } |
| 913 | }; |
| 914 | |
| 915 | // Streamed like the audit log: allocation follows one line, not the file. |
| 916 | for raw_line in std::io::BufRead::lines(std::io::BufReader::new(file)) { |
| 917 | let Ok(raw_line) = raw_line.inspect_err(|e| { |
| 918 | tracing::trace!( |
| 919 | "metrics: stopped reading events file {}: {e}", |
| 920 | path.display() |
| 921 | ); |
| 922 | }) else { |
| 923 | break; |
| 924 | }; |
| 925 | rollup.total_lines += 1; |
| 926 | let line = raw_line.trim(); |
| 927 | if line.is_empty() { |
| 928 | continue; |
| 929 | } |
| 930 | |
| 931 | let v: Value = match serde_json::from_str(line) { |
| 932 | Ok(v) => v, |
| 933 | Err(e) => { |
| 934 | tracing::trace!("metrics: skipping malformed event line: {e}"); |
| 935 | continue; |
| 936 | } |
| 937 | }; |
| 938 | |
| 939 | let ts = parse_ts_field(&v, "timestamp"); |
| 940 | |
| 941 | if let Some(cutoff) = since { |
| 942 | match ts { |
| 943 | Some(t) if t < cutoff => continue, |
| 944 | _ => {} |
| 945 | } |
| 946 | } |
| 947 | |
| 948 | rollup.parsed_lines += 1; |
| 949 | if let Some(t) = &ts { |
| 950 | rollup.touch_ts(t); |
| 951 | } |
| 952 | |
| 953 | let event = v.get("event").and_then(|e| e.as_str()).unwrap_or(""); |
| 954 | |
| 955 | match event { |
| 956 | "turn.completed" => record_terminal_request_diagnostics(&v, rollup, dedup), |
| 957 | "turn.usage" => record_provider_usage_receipt(&v, rollup, dedup), |
| 958 | // Tool and compaction receipts are durable *item* records. The |
| 959 | // `tool.started` / `tool.completed` names this reader used to |
| 960 | // match are synthesized for HTTP clients by |
| 961 | // `map_compat_stream_event` (`crates/tui/src/runtime_api.rs`) and |
| 962 | // are never persisted, so those arms counted nothing. |
| 963 | "item.started" | "item.completed" | "item.failed" => { |
| 964 | record_runtime_item_receipt(event, &v, rollup, dedup); |
| 965 | } |
| 966 | "agent.spawned" | "subagent.spawned" => { |
| 967 | rollup.agents.spawns += 1; |
| 968 | } |
| 969 | "agent.completed" | "subagent.completed" => { |
| 970 | rollup.agents.record_completion(&v); |
| 971 | } |
| 972 | e if e.starts_with("capacity.") => { |
| 973 | rollup.capacity.total += 1; |
| 974 | let category = v |
| 975 | .pointer("/payload/category") |
| 976 | .and_then(|c| c.as_str()) |
| 977 | .unwrap_or(e.trim_start_matches("capacity.")); |
| 978 | *rollup |
| 979 | .capacity |
| 980 | .by_category |
| 981 | .entry(category.to_string()) |
| 982 | .or_insert(0) += 1; |
| 983 | } |
| 984 | _ => {} |
| 985 | } |
| 986 | } |
| 987 | } |
| 988 | |
| 989 | fn runtime_event_identity(v: &Value) -> Option<(String, u64)> { |
| 990 | Some(( |
| 991 | v.get("thread_id")?.as_str()?.to_string(), |
| 992 | v.get("seq")?.as_u64()?, |
| 993 | )) |
| 994 | } |
| 995 | |
| 996 | fn terminal_turn_identity(v: &Value) -> Option<(String, String)> { |
| 997 | Some(( |
| 998 | v.get("thread_id")?.as_str()?.to_string(), |
| 999 | v.get("turn_id")?.as_str()?.to_string(), |
| 1000 | )) |
| 1001 | } |
| 1002 | |
| 1003 | fn record_terminal_request_diagnostics( |
| 1004 | v: &Value, |
| 1005 | rollup: &mut Rollup, |
| 1006 | dedup: &mut RuntimeEventDedup, |
| 1007 | ) { |
| 1008 | let Some(identity) = terminal_turn_identity(v) else { |
| 1009 | rollup.runtime_requests.terminal_receipts_without_identity = rollup |
| 1010 | .runtime_requests |
| 1011 | .terminal_receipts_without_identity |
| 1012 | .saturating_add(1); |
| 1013 | return; |
| 1014 | }; |
| 1015 | if !dedup.terminal_turns.insert(identity) { |
| 1016 | rollup.runtime_requests.duplicate_terminal_receipts_skipped = rollup |
| 1017 | .runtime_requests |
| 1018 | .duplicate_terminal_receipts_skipped |
| 1019 | .saturating_add(1); |
| 1020 | return; |
| 1021 | } |
| 1022 | |
| 1023 | let stats = &mut rollup.runtime_requests; |
| 1024 | stats.terminal_turn_receipts = stats.terminal_turn_receipts.saturating_add(1); |
| 1025 | let Some(diagnostics) = v.pointer("/payload/turn/modelRequestDiagnostics") else { |
| 1026 | stats.diagnostics_unavailable_turn_receipts = stats |
| 1027 | .diagnostics_unavailable_turn_receipts |
| 1028 | .saturating_add(1); |
| 1029 | return; |
| 1030 | }; |
| 1031 | let Some(model_requests_started) = diagnostics |
| 1032 | .get("modelRequestsStarted") |
| 1033 | .and_then(Value::as_u64) |
| 1034 | else { |
| 1035 | stats.diagnostics_incomplete_turn_receipts = |
| 1036 | stats.diagnostics_incomplete_turn_receipts.saturating_add(1); |
| 1037 | return; |
| 1038 | }; |
| 1039 | let Some(transparent_stream_retries) = diagnostics |
| 1040 | .get("transparentStreamRetries") |
| 1041 | .and_then(Value::as_u64) |
| 1042 | else { |
| 1043 | stats.diagnostics_incomplete_turn_receipts = |
| 1044 | stats.diagnostics_incomplete_turn_receipts.saturating_add(1); |
| 1045 | return; |
| 1046 | }; |
| 1047 | let Some(stream_resumes) = diagnostics.get("streamResumes").and_then(Value::as_u64) else { |
| 1048 | stats.diagnostics_incomplete_turn_receipts = |
| 1049 | stats.diagnostics_incomplete_turn_receipts.saturating_add(1); |
| 1050 | return; |
| 1051 | }; |
| 1052 | |
| 1053 | stats.diagnostics_turn_receipts = stats.diagnostics_turn_receipts.saturating_add(1); |
| 1054 | stats.model_requests_started = stats |
| 1055 | .model_requests_started |
| 1056 | .saturating_add(model_requests_started); |
| 1057 | stats.transparent_stream_retries = stats |
| 1058 | .transparent_stream_retries |
| 1059 | .saturating_add(transparent_stream_retries); |
| 1060 | stats.stream_resumes = stats.stream_resumes.saturating_add(stream_resumes); |
| 1061 | } |
| 1062 | |
| 1063 | /// Fold one durable `item.*` receipt into the tool or compaction rollup. |
| 1064 | /// |
| 1065 | /// Reads only the fields a rollup needs — `kind`, `metadata.tool_name`, |
| 1066 | /// `metadata.is_error`, the two timestamps, and the compaction message counts. |
| 1067 | /// `item.detail` and `metadata.tool_input` carry tool output and arguments and |
| 1068 | /// are deliberately never read here. |
| 1069 | fn record_runtime_item_receipt( |
| 1070 | event: &str, |
| 1071 | v: &Value, |
| 1072 | rollup: &mut Rollup, |
| 1073 | dedup: &mut RuntimeEventDedup, |
| 1074 | ) { |
| 1075 | let Some(item) = v.pointer("/payload/item") else { |
| 1076 | return; |
| 1077 | }; |
| 1078 | let kind = item.get("kind").and_then(Value::as_str).unwrap_or_default(); |
| 1079 | let is_tool = matches!(kind, "tool_call" | "file_change" | "command_execution"); |
| 1080 | if !is_tool && kind != "context_compaction" { |
| 1081 | return; |
| 1082 | } |
| 1083 | // One receipt, counted once. A record with no verifiable runtime identity |
| 1084 | // cannot be de-duplicated, so it is counted without one rather than |
| 1085 | // dropped — `thread_id` and `seq` are required fields of every record the |
| 1086 | // store writes, so this only affects hand-edited logs. |
| 1087 | if let Some(identity) = runtime_event_identity(v) |
| 1088 | && !dedup.event_records.insert(identity) |
| 1089 | { |
| 1090 | return; |
| 1091 | } |
| 1092 | |
| 1093 | if !is_tool { |
| 1094 | record_compaction_item_receipt(event, v, rollup); |
| 1095 | return; |
| 1096 | } |
| 1097 | |
| 1098 | // `tool_name` is copied forward into the completion metadata, but the |
| 1099 | // redaction and error branches rewrite or leave that object alone, so fall |
| 1100 | // back to the started record's `tool` projection before giving up. An |
| 1101 | // unresolvable name is bucketed, never dropped. |
| 1102 | let tool_name = item |
| 1103 | .pointer("/metadata/tool_name") |
| 1104 | .or_else(|| v.pointer("/payload/tool/name")) |
| 1105 | .and_then(Value::as_str) |
| 1106 | .unwrap_or("unknown"); |
| 1107 | let elapsed_ms = item_elapsed_ms(item); |
| 1108 | let stats = rollup.tool_mut(tool_name); |
| 1109 | match event { |
| 1110 | "item.started" => stats.calls += 1, |
| 1111 | "item.completed" => match item.pointer("/metadata/is_error").and_then(Value::as_bool) { |
| 1112 | Some(false) => stats.successes += 1, |
| 1113 | Some(true) => stats.failures += 1, |
| 1114 | None => stats.outcome_unknown += 1, |
| 1115 | }, |
| 1116 | // `item.failed` is the engine's own error branch: the tool call did |
| 1117 | // run and did not succeed. |
| 1118 | _ => stats.failures += 1, |
| 1119 | } |
| 1120 | if event != "item.started" { |
| 1121 | match elapsed_ms { |
| 1122 | Some(ms) => { |
| 1123 | stats.total_elapsed_ms = stats.total_elapsed_ms.saturating_add(ms); |
| 1124 | stats.elapsed_samples += 1; |
| 1125 | } |
| 1126 | None => stats.elapsed_unavailable += 1, |
| 1127 | } |
| 1128 | } |
| 1129 | } |
| 1130 | |
| 1131 | fn compaction_receipt_identity(v: &Value, runtime: bool) -> Option<String> { |
| 1132 | let (session, id) = if runtime { |
| 1133 | ( |
| 1134 | v.get("thread_id")?, |
| 1135 | v.pointer("/payload/item/metadata/compaction_id")?, |
| 1136 | ) |
| 1137 | } else { |
| 1138 | ( |
| 1139 | v.pointer("/details/thread_id") |
| 1140 | .filter(|id| id.is_string()) |
| 1141 | .or_else(|| v.pointer("/details/session_id"))?, |
| 1142 | v.pointer("/details/compaction_id")?, |
| 1143 | ) |
| 1144 | }; |
| 1145 | Some(format!("{}:{}", session.as_str()?, id.as_str()?)) |
| 1146 | } |
| 1147 | |
| 1148 | /// A compaction is only counted when it completed. The size reduction comes |
| 1149 | /// from the two persisted message counts or it stays unknown — a compaction |
| 1150 | /// with no counts must never average in as a 0% reduction. |
| 1151 | fn record_compaction_item_receipt(event: &str, v: &Value, rollup: &mut Rollup) { |
| 1152 | if event != "item.completed" { |
| 1153 | return; |
| 1154 | } |
| 1155 | if let Some(id) = compaction_receipt_identity(v, true) |
| 1156 | && !rollup.compaction.receipt_ids.insert(id) |
| 1157 | { |
| 1158 | return; |
| 1159 | } |
| 1160 | rollup.compaction.events += 1; |
| 1161 | let before = v |
| 1162 | .pointer("/payload/messages_before") |
| 1163 | .and_then(Value::as_u64); |
| 1164 | let after = v.pointer("/payload/messages_after").and_then(Value::as_u64); |
| 1165 | let (Some(before), Some(after)) = (before, after) else { |
| 1166 | return; |
| 1167 | }; |
| 1168 | if before == 0 { |
| 1169 | return; |
| 1170 | } |
| 1171 | rollup.compaction.ratio_sum += 1.0 - (after as f64 / before as f64); |
| 1172 | rollup.compaction.ratio_samples += 1; |
| 1173 | } |
| 1174 | |
| 1175 | /// Exact interval between two persisted item timestamps, or `None`. |
| 1176 | /// |
| 1177 | /// Both fields are optional on the record, and a reversed pair is not a |
| 1178 | /// measurement. Neither case may contribute a 0 ms sample, and neither is |
| 1179 | /// evidence about a provider charge. |
| 1180 | fn item_elapsed_ms(item: &Value) -> Option<u64> { |
| 1181 | let started = parse_ts_field(item, "started_at")?; |
| 1182 | let ended = parse_ts_field(item, "ended_at")?; |
| 1183 | u64::try_from((ended - started).num_milliseconds()).ok() |
| 1184 | } |
| 1185 | |
| 1186 | fn record_provider_usage_receipt(v: &Value, rollup: &mut Rollup, dedup: &mut RuntimeEventDedup) { |
| 1187 | let Some(identity) = runtime_event_identity(v) else { |
| 1188 | rollup |
| 1189 | .runtime_requests |
| 1190 | .provider_usage_receipts_without_identity = rollup |
| 1191 | .runtime_requests |
| 1192 | .provider_usage_receipts_without_identity |
| 1193 | .saturating_add(1); |
| 1194 | return; |
| 1195 | }; |
| 1196 | if !dedup.event_records.insert(identity) { |
| 1197 | rollup |
| 1198 | .runtime_requests |
| 1199 | .duplicate_provider_usage_receipts_skipped = rollup |
| 1200 | .runtime_requests |
| 1201 | .duplicate_provider_usage_receipts_skipped |
| 1202 | .saturating_add(1); |
| 1203 | return; |
| 1204 | } |
| 1205 | |
| 1206 | let stats = &mut rollup.runtime_requests; |
| 1207 | let Some(input_tokens) = v |
| 1208 | .pointer("/payload/usage/input_tokens") |
| 1209 | .and_then(Value::as_u64) |
| 1210 | else { |
| 1211 | stats.provider_usage_receipts_incomplete = |
| 1212 | stats.provider_usage_receipts_incomplete.saturating_add(1); |
| 1213 | return; |
| 1214 | }; |
| 1215 | let Some(output_tokens) = v |
| 1216 | .pointer("/payload/usage/output_tokens") |
| 1217 | .and_then(Value::as_u64) |
| 1218 | else { |
| 1219 | stats.provider_usage_receipts_incomplete = |
| 1220 | stats.provider_usage_receipts_incomplete.saturating_add(1); |
| 1221 | return; |
| 1222 | }; |
| 1223 | stats.provider_usage_receipts = stats.provider_usage_receipts.saturating_add(1); |
| 1224 | stats.provider_reported_input_tokens = stats |
| 1225 | .provider_reported_input_tokens |
| 1226 | .saturating_add(input_tokens); |
| 1227 | stats.provider_reported_output_tokens = stats |
| 1228 | .provider_reported_output_tokens |
| 1229 | .saturating_add(output_tokens); |
| 1230 | } |
| 1231 | |
| 1232 | // ────────────────────────────────────────────────────────────────────────────── |
| 1233 | // Output formatters |
| 1234 | // ────────────────────────────────────────────────────────────────────────────── |
| 1235 | |
| 1236 | fn print_json(rollup: &Rollup) -> Result<()> { |
| 1237 | println!("{}", serde_json::to_string_pretty(rollup)?); |
| 1238 | Ok(()) |
| 1239 | } |
| 1240 | |
| 1241 | fn print_human(rollup: &Rollup, since: Option<DateTime<Utc>>) { |
| 1242 | // Period header. When a --since cutoff yields nothing, print it: a bare |
| 1243 | // `--since 7` is seven seconds, and the empty window must not read the |
| 1244 | // same as a genuinely idle period (#6315). |
| 1245 | match (rollup.earliest_ts, rollup.latest_ts) { |
| 1246 | (Some(start), Some(end)) => { |
| 1247 | let days = (end - start).num_days(); |
| 1248 | println!( |
| 1249 | "Period: {} → {} ({} days)", |
| 1250 | start.format("%Y-%m-%d"), |
| 1251 | end.format("%Y-%m-%d"), |
| 1252 | days |
| 1253 | ); |
| 1254 | } |
| 1255 | (Some(start), None) | (None, Some(start)) => { |
| 1256 | println!("Period: {} → (unknown)", start.format("%Y-%m-%d")); |
| 1257 | } |
| 1258 | (None, None) => match since { |
| 1259 | Some(cutoff) => println!( |
| 1260 | "Period: (no data since {})", |
| 1261 | cutoff.format("%Y-%m-%d %H:%M UTC") |
| 1262 | ), |
| 1263 | None => println!("Period: (no data)"), |
| 1264 | }, |
| 1265 | } |
| 1266 | |
| 1267 | // ── Tools ────────────────────────────────────────────────────────────── |
| 1268 | let total_calls = rollup.total_tool_calls(); |
| 1269 | if total_calls > 0 { |
| 1270 | // Overall success rate from session-file data (where we have result info). |
| 1271 | let total_ok: u64 = rollup.tools.values().map(|t| t.successes).sum(); |
| 1272 | let total_judged: u64 = rollup |
| 1273 | .tools |
| 1274 | .values() |
| 1275 | .map(|t| t.successes + t.failures) |
| 1276 | .sum(); |
| 1277 | let mut overall_rate = if total_judged > 0 { |
| 1278 | format!( |
| 1279 | "{:.1}% success", |
| 1280 | total_ok as f64 / total_judged as f64 * 100.0 |
| 1281 | ) |
| 1282 | } else { |
| 1283 | // Only approval events — show prompt breakdown. |
| 1284 | let auto: u64 = rollup.tools.values().map(|t| t.auto_approved).sum(); |
| 1285 | let prompted: u64 = rollup.tools.values().map(|t| t.prompted).sum(); |
| 1286 | format!("{auto} auto-approved, {prompted} prompted") |
| 1287 | }; |
| 1288 | // Denied and outcome-unknown calls are excluded from the rate above by |
| 1289 | // construction, so they are named rather than silently dropped. |
| 1290 | let total_denied: u64 = rollup.tools.values().map(|t| t.denied).sum(); |
| 1291 | if total_denied > 0 { |
| 1292 | overall_rate.push_str(&format!(", {} denied", fmt_num(total_denied))); |
| 1293 | } |
| 1294 | let total_unknown: u64 = rollup.tools.values().map(|t| t.outcome_unknown).sum(); |
| 1295 | if total_unknown > 0 { |
| 1296 | overall_rate.push_str(&format!(", {} outcome unknown", fmt_num(total_unknown))); |
| 1297 | } |
| 1298 | |
| 1299 | println!( |
| 1300 | "Tools: {:>6} calls ({})", |
| 1301 | fmt_num(total_calls), |
| 1302 | overall_rate |
| 1303 | ); |
| 1304 | |
| 1305 | // Sort tools by call count descending, top 15. |
| 1306 | let mut tools: Vec<(&String, &ToolStats)> = rollup.tools.iter().collect(); |
| 1307 | tools.sort_by_key(|b| std::cmp::Reverse(b.1.calls)); |
| 1308 | for (name, stats) in tools.iter().take(15) { |
| 1309 | let rate_str = match stats.success_rate_pct() { |
| 1310 | Some(pct) => format!("{pct:5.1}%"), |
| 1311 | None if stats.denied > 0 => { |
| 1312 | // Nothing ran, so an approval breakdown would read as if |
| 1313 | // it had. |
| 1314 | format!("{} denied", fmt_num(stats.denied)) |
| 1315 | } |
| 1316 | None if stats.outcome_unknown > 0 => { |
| 1317 | format!("{} unknown", fmt_num(stats.outcome_unknown)) |
| 1318 | } |
| 1319 | None => { |
| 1320 | // Only approval data available — show auto/prompted breakdown. |
| 1321 | let a = stats.auto_approved; |
| 1322 | let p = stats.prompted; |
| 1323 | if p == 0 { |
| 1324 | format!("auto×{a} ") |
| 1325 | } else { |
| 1326 | format!("auto×{a}/prompted×{p}") |
| 1327 | } |
| 1328 | } |
| 1329 | }; |
| 1330 | let avg_str = match stats.avg_elapsed_ms() { |
| 1331 | Some(ms) => format!(" avg {ms}ms"), |
| 1332 | None => String::new(), |
| 1333 | }; |
| 1334 | println!( |
| 1335 | " {name:<22} {:>6} {rate_str}{avg_str}", |
| 1336 | fmt_num(stats.calls) |
| 1337 | ); |
| 1338 | } |
| 1339 | if tools.len() > 15 { |
| 1340 | println!(" … and {} more tools", tools.len() - 15); |
| 1341 | } |
| 1342 | } else { |
| 1343 | println!("Tools: (no data)"); |
| 1344 | } |
| 1345 | |
| 1346 | // ── Compaction ───────────────────────────────────────────────────────── |
| 1347 | let compaction_refusals: u64 = rollup.compaction.refusals.values().sum(); |
| 1348 | if rollup.compaction.events > 0 || compaction_refusals > 0 { |
| 1349 | let avg_str = match rollup.compaction.avg_reduction_pct() { |
| 1350 | Some(pct) => format!(", avg {pct:.0}% size reduction"), |
| 1351 | // No message counts were recorded. Saying nothing here reads as |
| 1352 | // "no reduction"; say that it is unknown. |
| 1353 | None => ", size reduction unknown".to_string(), |
| 1354 | }; |
| 1355 | let usage = if rollup.compaction.summarizer_usage_samples > 0 { |
| 1356 | format!( |
| 1357 | "{} input / {} output tokens across {} measured passes", |
| 1358 | fmt_num(rollup.compaction.summarizer_input_tokens), |
| 1359 | fmt_num(rollup.compaction.summarizer_output_tokens), |
| 1360 | fmt_num(rollup.compaction.summarizer_usage_samples) |
| 1361 | ) |
| 1362 | } else { |
| 1363 | "usage unavailable".to_string() |
| 1364 | }; |
| 1365 | println!( |
| 1366 | "Compaction: {} completed, {} refused{}; summarizer {}", |
| 1367 | fmt_num(rollup.compaction.events), |
| 1368 | fmt_num(compaction_refusals), |
| 1369 | avg_str, |
| 1370 | usage |
| 1371 | ); |
| 1372 | } else { |
| 1373 | println!("Compaction: (no data)"); |
| 1374 | } |
| 1375 | |
| 1376 | // ── Sub-agents ───────────────────────────────────────────────────────── |
| 1377 | println!("{}", rollup.agents.summary()); |
| 1378 | |
| 1379 | // ── Capacity interventions ───────────────────────────────────────────── |
| 1380 | if rollup.capacity.total > 0 { |
| 1381 | let cat_str: String = { |
| 1382 | let mut cats: Vec<(&String, &u64)> = rollup.capacity.by_category.iter().collect(); |
| 1383 | cats.sort_by(|a, b| b.1.cmp(a.1)); |
| 1384 | cats.iter() |
| 1385 | .map(|(k, v)| format!("{v} {k}")) |
| 1386 | .collect::<Vec<_>>() |
| 1387 | .join(", ") |
| 1388 | }; |
| 1389 | println!( |
| 1390 | "Capacity interventions: {} ({})", |
| 1391 | fmt_num(rollup.capacity.total), |
| 1392 | cat_str |
| 1393 | ); |
| 1394 | } else { |
| 1395 | println!("Capacity interventions: (no data)"); |
| 1396 | } |
| 1397 | |
| 1398 | // ── Runtime request and provider-usage receipts ─────────────────────── |
| 1399 | let runtime = &rollup.runtime_requests; |
| 1400 | if runtime.terminal_turn_receipts == 0 { |
| 1401 | println!("Runtime requests: (no terminal receipts; model-client counts unknown)"); |
| 1402 | } else if runtime.diagnostics_turn_receipts == 0 { |
| 1403 | println!( |
| 1404 | "Runtime requests: (diagnostics unavailable for all {} terminal receipts)", |
| 1405 | fmt_num(runtime.terminal_turn_receipts) |
| 1406 | ); |
| 1407 | } else { |
| 1408 | println!( |
| 1409 | "Runtime requests: {} model-client calls, {} stream resumes, {} transparent retries (diagnostics for {}/{} terminal receipts; status events excluded)", |
| 1410 | fmt_num(runtime.model_requests_started), |
| 1411 | fmt_num(runtime.stream_resumes), |
| 1412 | fmt_num(runtime.transparent_stream_retries), |
| 1413 | fmt_num(runtime.diagnostics_turn_receipts), |
| 1414 | fmt_num(runtime.terminal_turn_receipts), |
| 1415 | ); |
| 1416 | } |
| 1417 | if runtime.provider_usage_receipts == 0 { |
| 1418 | println!("Provider usage receipts: (none recorded; this is not zero usage)"); |
| 1419 | } else { |
| 1420 | println!( |
| 1421 | "Provider usage receipts: {} records, {} input tokens, {} output tokens", |
| 1422 | fmt_num(runtime.provider_usage_receipts), |
| 1423 | fmt_num(runtime.provider_reported_input_tokens), |
| 1424 | fmt_num(runtime.provider_reported_output_tokens), |
| 1425 | ); |
| 1426 | } |
| 1427 | if runtime.diagnostics_unavailable_turn_receipts > 0 |
| 1428 | || runtime.diagnostics_incomplete_turn_receipts > 0 |
| 1429 | || runtime.terminal_receipts_without_identity > 0 |
| 1430 | || runtime.provider_usage_receipts_without_identity > 0 |
| 1431 | || runtime.provider_usage_receipts_incomplete > 0 |
| 1432 | || runtime.duplicate_terminal_receipts_skipped > 0 |
| 1433 | || runtime.duplicate_provider_usage_receipts_skipped > 0 |
| 1434 | { |
| 1435 | println!( |
| 1436 | "Runtime receipt coverage: {} diagnostics unavailable, {} diagnostics incomplete, {} terminal receipts without identity, {} usage receipts without identity, {} usage receipts incomplete, {} duplicate terminal receipts skipped, {} duplicate usage receipts skipped", |
| 1437 | fmt_num(runtime.diagnostics_unavailable_turn_receipts), |
| 1438 | fmt_num(runtime.diagnostics_incomplete_turn_receipts), |
| 1439 | fmt_num(runtime.terminal_receipts_without_identity), |
| 1440 | fmt_num(runtime.provider_usage_receipts_without_identity), |
| 1441 | fmt_num(runtime.provider_usage_receipts_incomplete), |
| 1442 | fmt_num(runtime.duplicate_terminal_receipts_skipped), |
| 1443 | fmt_num(runtime.duplicate_provider_usage_receipts_skipped), |
| 1444 | ); |
| 1445 | } |
| 1446 | |
| 1447 | // ── Credentials ──────────────────────────────────────────────────────── |
| 1448 | if rollup.credentials.saves > 0 || rollup.credentials.clears > 0 { |
| 1449 | println!( |
| 1450 | "Credentials: {} saves, {} clears", |
| 1451 | rollup.credentials.saves, rollup.credentials.clears |
| 1452 | ); |
| 1453 | } |
| 1454 | } |
| 1455 | |
| 1456 | // ────────────────────────────────────────────────────────────────────────────── |
| 1457 | // Helpers |
| 1458 | // ────────────────────────────────────────────────────────────────────────────── |
| 1459 | |
| 1460 | /// An explicit home is an isolation boundary. Default installs can have |
| 1461 | /// distinct audit histories in both roots, even after a copied migration. |
| 1462 | fn resolve_audit_roots() -> Result<Vec<PathBuf>> { |
| 1463 | let primary = codewhale_config::codewhale_home()?; |
| 1464 | let mut roots = vec![primary]; |
| 1465 | if !codewhale_config::codewhale_home_is_explicit() { |
| 1466 | let legacy = codewhale_config::legacy_deepseek_home()?; |
| 1467 | if !roots.contains(&legacy) { |
| 1468 | roots.push(legacy); |
| 1469 | } |
| 1470 | } |
| 1471 | Ok(roots) |
| 1472 | } |
| 1473 | |
| 1474 | /// Resolve the tool a durable audit record is about. |
| 1475 | /// |
| 1476 | /// `~/.codewhale/audit.log` nests its payload under `details`; the opt-in |
| 1477 | /// `CODEWHALE_TOOL_AUDIT_LOG` file writes `tool_name` at the top level. |
| 1478 | fn audit_tool_name(v: &Value) -> &str { |
| 1479 | v.pointer("/details/tool_name") |
| 1480 | .or_else(|| v.pointer("/payload/tool_name")) |
| 1481 | .or_else(|| v.get("tool_name")) |
| 1482 | .and_then(Value::as_str) |
| 1483 | .unwrap_or("unknown") |
| 1484 | } |
| 1485 | |
| 1486 | /// Parse a timestamp from a JSON value field (tries RFC3339). |
| 1487 | fn parse_ts_field(v: &Value, field: &str) -> Option<DateTime<Utc>> { |
| 1488 | v.get(field)?.as_str()?.parse::<DateTime<Utc>>().ok() |
| 1489 | } |
| 1490 | |
| 1491 | /// Format a number with thousands separators. |
| 1492 | fn fmt_num(n: u64) -> String { |
| 1493 | let s = n.to_string(); |
| 1494 | let mut result = String::with_capacity(s.len() + s.len() / 3); |
| 1495 | for (i, ch) in s.chars().rev().enumerate() { |
| 1496 | if i > 0 && i % 3 == 0 { |
| 1497 | result.push(','); |
| 1498 | } |
| 1499 | result.push(ch); |
| 1500 | } |
| 1501 | result.chars().rev().collect() |
| 1502 | } |
| 1503 | |
| 1504 | // ────────────────────────────────────────────────────────────────────────────── |
| 1505 | // Tests |
| 1506 | // ────────────────────────────────────────────────────────────────────────────── |
| 1507 | |
| 1508 | #[cfg(test)] |
| 1509 | mod tests { |
| 1510 | use super::*; |
| 1511 | |
| 1512 | #[test] |
| 1513 | fn compaction_audit_counts_usage_refusals_and_deduplicates_runtime_receipt() { |
| 1514 | let tmp = tempfile::NamedTempFile::new().unwrap(); |
| 1515 | let completed = serde_json::json!({"ts": "2026-09-19T12:00:00Z", "event": "compaction.completed", "details": { |
| 1516 | "session_id": "engine-session-a", "thread_id": "thread-a", "compaction_id": "pass-a", "trigger": "manual", "path": "summary", |
| 1517 | "reduction_ratio": 0.75, "summarizer_usage": {"input_tokens": 120, "output_tokens": 15} |
| 1518 | }}); |
| 1519 | let refused = serde_json::json!({"ts": "2026-09-19T12:01:00Z", "event": "compaction.refused", "details": {"reason": "retained_floor"}}); |
| 1520 | std::fs::write(tmp.path(), format!("{completed}\n{completed}\n{refused}\n")).unwrap(); |
| 1521 | let mut rollup = Rollup::default(); |
| 1522 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 1523 | record_compaction_item_receipt( |
| 1524 | "item.completed", |
| 1525 | &serde_json::json!({ |
| 1526 | "thread_id": "thread-a", "payload": {"item": {"metadata": {"compaction_id": "pass-a"}}, "messages_before": 4, "messages_after": 1} |
| 1527 | }), |
| 1528 | &mut rollup, |
| 1529 | ); |
| 1530 | assert_eq!(rollup.compaction.events, 1); |
| 1531 | assert_eq!(rollup.compaction.ratio_samples, 1); |
| 1532 | assert_eq!(rollup.compaction.avg_reduction_pct(), Some(75.0)); |
| 1533 | assert_eq!(rollup.compaction.refusals["retained_floor"], 1); |
| 1534 | assert_eq!(rollup.compaction.triggers["manual"], 1); |
| 1535 | assert_eq!(rollup.compaction.paths["summary"], 1); |
| 1536 | assert_eq!(rollup.compaction.summarizer_input_tokens, 120); |
| 1537 | assert_eq!(rollup.compaction.summarizer_output_tokens, 15); |
| 1538 | } |
| 1539 | |
| 1540 | fn read_audit_test_log(path: &Path, since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 1541 | super::read_audit_log(path, since, rollup, &HashMap::new(), &mut HashMap::new()); |
| 1542 | } |
| 1543 | |
| 1544 | fn read_runtime_test_log(path: &Path, since: Option<DateTime<Utc>>, rollup: &mut Rollup) { |
| 1545 | super::read_events_jsonl(path, since, rollup, &mut RuntimeEventDedup::default()); |
| 1546 | } |
| 1547 | |
| 1548 | fn runtime_event( |
| 1549 | seq: u64, |
| 1550 | timestamp: &str, |
| 1551 | thread_id: &str, |
| 1552 | turn_id: Option<&str>, |
| 1553 | event: &str, |
| 1554 | payload: Value, |
| 1555 | ) -> Value { |
| 1556 | serde_json::json!({ |
| 1557 | "schema_version": 4, |
| 1558 | "seq": seq, |
| 1559 | "timestamp": timestamp, |
| 1560 | "thread_id": thread_id, |
| 1561 | "turn_id": turn_id, |
| 1562 | "event": event, |
| 1563 | "payload": payload, |
| 1564 | }) |
| 1565 | } |
| 1566 | |
| 1567 | fn write_runtime_events(events: &[Value]) -> tempfile::NamedTempFile { |
| 1568 | use std::io::Write; |
| 1569 | |
| 1570 | let mut tmp = tempfile::NamedTempFile::new().unwrap(); |
| 1571 | for event in events { |
| 1572 | writeln!(tmp, "{event}").unwrap(); |
| 1573 | } |
| 1574 | tmp |
| 1575 | } |
| 1576 | |
| 1577 | #[test] |
| 1578 | fn runtime_worker_completion_uses_owner_outcome_not_completed_receipt_status() { |
| 1579 | let statuses = [ |
| 1580 | serde_json::json!("completed"), |
| 1581 | serde_json::json!("failed"), |
| 1582 | serde_json::json!("cancelled"), |
| 1583 | serde_json::json!("interrupted"), |
| 1584 | serde_json::json!("budget_exhausted"), |
| 1585 | Value::Null, |
| 1586 | ]; |
| 1587 | let events: Vec<_> = statuses |
| 1588 | .into_iter() |
| 1589 | .enumerate() |
| 1590 | .map(|(seq, worker_status)| { |
| 1591 | runtime_event( |
| 1592 | seq as u64, |
| 1593 | "2026-09-08T10:00:00Z", |
| 1594 | "thread-a", |
| 1595 | Some("turn-a"), |
| 1596 | "agent.completed", |
| 1597 | serde_json::json!({ |
| 1598 | "item": { "kind": "status", "status": "completed" }, |
| 1599 | "agent_id": format!("worker-{seq}"), |
| 1600 | "worker_status": worker_status, |
| 1601 | "parent_run_id": "run-a", |
| 1602 | "spawn_depth": 1, |
| 1603 | "continuable": false, |
| 1604 | }), |
| 1605 | ) |
| 1606 | }) |
| 1607 | .collect(); |
| 1608 | let tmp = write_runtime_events(&events); |
| 1609 | let mut rollup = Rollup::default(); |
| 1610 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 1611 | let agents = &rollup.agents; |
| 1612 | assert_eq!(agents.successes, 1, "a settled item is not worker success"); |
| 1613 | assert_eq!(agents.failures, 1); |
| 1614 | assert_eq!(agents.cancelled, 1); |
| 1615 | assert_eq!(agents.interrupted, 1); |
| 1616 | assert_eq!(agents.budget_exhausted, 1); |
| 1617 | assert_eq!(agents.unknown_outcomes, 1); |
| 1618 | assert_eq!(agents.spawns, 0); |
| 1619 | let summary = agents.summary(); |
| 1620 | assert!(summary.contains("1 failed")); |
| 1621 | assert!(summary.contains("1 outcome unconfirmed")); |
| 1622 | assert!( |
| 1623 | !summary.contains("no data"), |
| 1624 | "terminal-only windows have data" |
| 1625 | ); |
| 1626 | assert!( |
| 1627 | !summary.contains("%"), |
| 1628 | "partial receipts are not a success rate" |
| 1629 | ); |
| 1630 | } |
| 1631 | |
| 1632 | #[test] |
| 1633 | fn runtime_worker_usage_sums_tokens_across_completion_receipts() { |
| 1634 | // #6315: completions with usage receipts sum into the rollup; a |
| 1635 | // completion without usage is a missing receipt, never zero tokens. |
| 1636 | let events = vec![ |
| 1637 | runtime_event( |
| 1638 | 0, |
| 1639 | "2026-09-08T10:00:00Z", |
| 1640 | "thread-a", |
| 1641 | Some("turn-a"), |
| 1642 | "agent.completed", |
| 1643 | serde_json::json!({ |
| 1644 | "agent_id": "worker-0", |
| 1645 | "worker_status": "completed", |
| 1646 | "usage": { |
| 1647 | "status": "completed", |
| 1648 | "input_tokens": 800, |
| 1649 | "output_tokens": 200, |
| 1650 | "total_tokens": 1000, |
| 1651 | "cost_microusd": 50, |
| 1652 | }, |
| 1653 | }), |
| 1654 | ), |
| 1655 | runtime_event( |
| 1656 | 1, |
| 1657 | "2026-09-08T10:01:00Z", |
| 1658 | "thread-a", |
| 1659 | Some("turn-a"), |
| 1660 | "agent.completed", |
| 1661 | serde_json::json!({ |
| 1662 | "agent_id": "worker-1", |
| 1663 | "worker_status": "completed", |
| 1664 | }), |
| 1665 | ), |
| 1666 | ]; |
| 1667 | let tmp = write_runtime_events(&events); |
| 1668 | let mut rollup = Rollup::default(); |
| 1669 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 1670 | let agents = &rollup.agents; |
| 1671 | assert_eq!(agents.successes, 2); |
| 1672 | assert_eq!(agents.usage_receipts, 1); |
| 1673 | assert_eq!(agents.input_tokens, 800); |
| 1674 | assert_eq!(agents.output_tokens, 200); |
| 1675 | assert_eq!(agents.total_tokens, 1000); |
| 1676 | assert_eq!(agents.cost_microusd, 50); |
| 1677 | let summary = agents.summary(); |
| 1678 | assert!(summary.contains("800 in/200 out/1,000 total (1 receipt)")); |
| 1679 | } |
| 1680 | |
| 1681 | #[test] |
| 1682 | fn runtime_legacy_worker_receipts_require_explicit_success_evidence() { |
| 1683 | let payloads = [ |
| 1684 | serde_json::json!({ "success": true }), |
| 1685 | serde_json::json!({ "success": false }), |
| 1686 | serde_json::json!({}), |
| 1687 | serde_json::json!({ "success": "true" }), |
| 1688 | serde_json::json!({ "worker_status": "failed", "success": true }), |
| 1689 | serde_json::json!({ "worker_status": null, "success": true }), |
| 1690 | serde_json::json!({ "worker_status": "running", "success": true }), |
| 1691 | serde_json::json!({ "worker_status": { "completed": true }, "success": true }), |
| 1692 | serde_json::json!({ "worker_status": "future_outcome", "success": true }), |
| 1693 | ]; |
| 1694 | let events: Vec<_> = payloads |
| 1695 | .into_iter() |
| 1696 | .enumerate() |
| 1697 | .map(|(seq, payload)| { |
| 1698 | runtime_event( |
| 1699 | seq as u64, |
| 1700 | "2026-09-08T10:00:00Z", |
| 1701 | "thread-a", |
| 1702 | Some("turn-a"), |
| 1703 | "agent.completed", |
| 1704 | payload, |
| 1705 | ) |
| 1706 | }) |
| 1707 | .collect(); |
| 1708 | let tmp = write_runtime_events(&events); |
| 1709 | let mut rollup = Rollup::default(); |
| 1710 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 1711 | assert_eq!(rollup.agents.successes, 1); |
| 1712 | assert_eq!(rollup.agents.failures, 2); |
| 1713 | assert_eq!(rollup.agents.unknown_outcomes, 6); |
| 1714 | } |
| 1715 | |
| 1716 | #[test] |
| 1717 | fn audit_worker_receipts_share_typed_and_legacy_outcome_rules() { |
| 1718 | let events = [ |
| 1719 | serde_json::json!({ "event": "agent.completed", "details": { "worker_status": "failed", "success": true } }), |
| 1720 | serde_json::json!({ "event": "subagent.completed", "payload": { "status": "cancelled", "success": true } }), |
| 1721 | serde_json::json!({ "event": "subagent.completed", "details": { "status": "completed" } }), |
| 1722 | serde_json::json!({ "event": "agent.completed", "details": { "success": false } }), |
| 1723 | serde_json::json!({ "event": "agent.completed", "payload": { "success": true } }), |
| 1724 | serde_json::json!({ "event": "agent.completed", "details": { "worker_status": null, "success": true } }), |
| 1725 | serde_json::json!({ "event": "agent.completed", "details": {} }), |
| 1726 | ]; |
| 1727 | let tmp = write_runtime_events(&events); |
| 1728 | let mut rollup = Rollup::default(); |
| 1729 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 1730 | assert_eq!(rollup.agents.successes, 2); |
| 1731 | assert_eq!(rollup.agents.failures, 2); |
| 1732 | assert_eq!(rollup.agents.cancelled, 1); |
| 1733 | assert_eq!(rollup.agents.unknown_outcomes, 2); |
| 1734 | let json = serde_json::to_value(&rollup).unwrap(); |
| 1735 | assert_eq!(json["agents"]["unknown_outcomes"], 2); |
| 1736 | assert_eq!(json["agents"]["cancelled"], 1); |
| 1737 | } |
| 1738 | |
| 1739 | // ── Duration parser ── |
| 1740 | |
| 1741 | #[test] |
| 1742 | fn parse_since_7d() { |
| 1743 | let cutoff = parse_since("7d").unwrap(); |
| 1744 | let expected = Utc::now() - Duration::days(7); |
| 1745 | // Allow ±2s for test execution time. |
| 1746 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1747 | } |
| 1748 | |
| 1749 | #[test] |
| 1750 | fn parse_since_24h() { |
| 1751 | let cutoff = parse_since("24h").unwrap(); |
| 1752 | let expected = Utc::now() - Duration::hours(24); |
| 1753 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1754 | } |
| 1755 | |
| 1756 | #[test] |
| 1757 | fn parse_since_30m() { |
| 1758 | let cutoff = parse_since("30m").unwrap(); |
| 1759 | let expected = Utc::now() - Duration::minutes(30); |
| 1760 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1761 | } |
| 1762 | |
| 1763 | #[test] |
| 1764 | fn parse_since_now_prefix() { |
| 1765 | // "now-2h" should strip "now-" and parse "2h". |
| 1766 | let cutoff = parse_since("now-2h").unwrap(); |
| 1767 | let expected = Utc::now() - Duration::hours(2); |
| 1768 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1769 | } |
| 1770 | |
| 1771 | #[test] |
| 1772 | fn parse_since_compound() { |
| 1773 | let cutoff = parse_since("2h30m").unwrap(); |
| 1774 | let expected = Utc::now() - Duration::seconds(2 * 3600 + 30 * 60); |
| 1775 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1776 | } |
| 1777 | |
| 1778 | #[test] |
| 1779 | fn parse_since_compound_days_hours() { |
| 1780 | let cutoff = parse_since("1d12h").unwrap(); |
| 1781 | let expected = Utc::now() - Duration::seconds(36 * 3600); |
| 1782 | assert!((cutoff - expected).num_seconds().abs() < 2); |
| 1783 | } |
| 1784 | |
| 1785 | #[test] |
| 1786 | fn parse_since_error_on_invalid() { |
| 1787 | assert!(parse_since("xyz").is_err()); |
| 1788 | assert!(parse_since("").is_err()); |
| 1789 | } |
| 1790 | |
| 1791 | #[test] |
| 1792 | fn parse_since_rejects_bare_unit() { |
| 1793 | let err = parse_since("d").unwrap_err().to_string(); |
| 1794 | assert!(err.contains("no number"), "{err}"); |
| 1795 | } |
| 1796 | |
| 1797 | #[test] |
| 1798 | fn parse_since_rejects_overflow_without_panicking() { |
| 1799 | // n * factor overflows i64. |
| 1800 | let err = parse_since("106751991167301d").unwrap_err().to_string(); |
| 1801 | assert!(err.contains("too large"), "{err}"); |
| 1802 | // Sum of components overflows i64. |
| 1803 | assert!(parse_since("9223372036854775807s1s").is_err()); |
| 1804 | // Fits in i64 seconds but exceeds TimeDelta's range. |
| 1805 | assert!(parse_since("9223372036854775807").is_err()); |
| 1806 | // Valid TimeDelta, but before the earliest representable DateTime. |
| 1807 | assert!(parse_since("100000000000d").is_err()); |
| 1808 | // Component too large to parse as i64. |
| 1809 | assert!(parse_since("99999999999999999999h").is_err()); |
| 1810 | } |
| 1811 | |
| 1812 | // ── fmt_num ── |
| 1813 | |
| 1814 | #[test] |
| 1815 | fn fmt_num_zero() { |
| 1816 | assert_eq!(fmt_num(0), "0"); |
| 1817 | } |
| 1818 | |
| 1819 | #[test] |
| 1820 | fn fmt_num_thousands() { |
| 1821 | assert_eq!(fmt_num(1_000), "1,000"); |
| 1822 | assert_eq!(fmt_num(12_453), "12,453"); |
| 1823 | assert_eq!(fmt_num(1_000_000), "1,000,000"); |
| 1824 | } |
| 1825 | |
| 1826 | // ── Rollup from audit log ── |
| 1827 | |
| 1828 | fn make_audit_line(event: &str, tool: &str, ts: &str) -> String { |
| 1829 | format!( |
| 1830 | r#"{{"details":{{"mode":"YOLO","session_id":null,"tool_name":"{tool}"}},"event":"{event}","ts":"{ts}"}}"# |
| 1831 | ) |
| 1832 | } |
| 1833 | |
| 1834 | #[test] |
| 1835 | fn audit_log_empty_file() { |
| 1836 | let mut rollup = Rollup::default(); |
| 1837 | // Non-existent path — should not panic, rollup stays empty. |
| 1838 | read_audit_test_log(Path::new("/nonexistent/audit.log"), None, &mut rollup); |
| 1839 | assert_eq!(rollup.total_lines, 0); |
| 1840 | } |
| 1841 | |
| 1842 | #[test] |
| 1843 | fn audit_log_parses_auto_approve() { |
| 1844 | use std::io::Write; |
| 1845 | let mut tmp = tempfile::NamedTempFile::new().unwrap(); |
| 1846 | let line1 = make_audit_line( |
| 1847 | "tool.approval.auto_approve", |
| 1848 | "exec_shell", |
| 1849 | "2026-04-01T10:00:00+00:00", |
| 1850 | ); |
| 1851 | let line2 = make_audit_line( |
| 1852 | "tool.approval.auto_approve", |
| 1853 | "read_file", |
| 1854 | "2026-04-02T10:00:00+00:00", |
| 1855 | ); |
| 1856 | writeln!(tmp, "{line1}").unwrap(); |
| 1857 | writeln!(tmp, "{line2}").unwrap(); |
| 1858 | |
| 1859 | let mut rollup = Rollup::default(); |
| 1860 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 1861 | |
| 1862 | assert_eq!(rollup.parsed_lines, 2); |
| 1863 | assert_eq!(rollup.tools["exec_shell"].calls, 1); |
| 1864 | assert_eq!(rollup.tools["exec_shell"].auto_approved, 1); |
| 1865 | assert_eq!(rollup.tools["read_file"].calls, 1); |
| 1866 | } |
| 1867 | |
| 1868 | #[test] |
| 1869 | fn audit_log_skips_malformed_lines() { |
| 1870 | use std::io::Write; |
| 1871 | let mut tmp = tempfile::NamedTempFile::new().unwrap(); |
| 1872 | writeln!(tmp, "not json at all").unwrap(); |
| 1873 | writeln!( |
| 1874 | tmp, |
| 1875 | r#"{{"event":"credential.save","ts":"2026-04-01T10:00:00+00:00"}}"# |
| 1876 | ) |
| 1877 | .unwrap(); |
| 1878 | |
| 1879 | let mut rollup = Rollup::default(); |
| 1880 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 1881 | |
| 1882 | // 2 lines total, 1 malformed skipped, 1 parsed. |
| 1883 | assert_eq!(rollup.total_lines, 2); |
| 1884 | assert_eq!(rollup.parsed_lines, 1); |
| 1885 | assert_eq!(rollup.credentials.saves, 1); |
| 1886 | } |
| 1887 | |
| 1888 | #[test] |
| 1889 | fn audit_log_since_filter() { |
| 1890 | use std::io::Write; |
| 1891 | let mut tmp = tempfile::NamedTempFile::new().unwrap(); |
| 1892 | let line_old = make_audit_line( |
| 1893 | "tool.approval.auto_approve", |
| 1894 | "exec_shell", |
| 1895 | "2025-01-01T00:00:00+00:00", |
| 1896 | ); |
| 1897 | let line_new = make_audit_line( |
| 1898 | "tool.approval.auto_approve", |
| 1899 | "read_file", |
| 1900 | "2026-04-01T00:00:00+00:00", |
| 1901 | ); |
| 1902 | writeln!(tmp, "{line_old}").unwrap(); |
| 1903 | writeln!(tmp, "{line_new}").unwrap(); |
| 1904 | |
| 1905 | let cutoff: DateTime<Utc> = "2026-01-01T00:00:00Z".parse().unwrap(); |
| 1906 | let mut rollup = Rollup::default(); |
| 1907 | read_audit_test_log(tmp.path(), Some(cutoff), &mut rollup); |
| 1908 | |
| 1909 | // Only the newer line should be counted. |
| 1910 | assert_eq!(rollup.parsed_lines, 1); |
| 1911 | assert!(!rollup.tools.contains_key("exec_shell")); |
| 1912 | assert_eq!(rollup.tools["read_file"].calls, 1); |
| 1913 | } |
| 1914 | |
| 1915 | #[test] |
| 1916 | fn total_tool_calls_sums_across_tools() { |
| 1917 | let mut rollup = Rollup::default(); |
| 1918 | rollup.tool_mut("read_file").calls = 4_012; |
| 1919 | rollup.tool_mut("exec_shell").calls = 1_118; |
| 1920 | assert_eq!(rollup.total_tool_calls(), 5_130); |
| 1921 | } |
| 1922 | |
| 1923 | // ── Runtime request and provider-usage receipts ── |
| 1924 | |
| 1925 | #[test] |
| 1926 | fn runtime_receipts_separate_terminal_requests_from_per_request_usage() { |
| 1927 | let timestamp = "2026-09-08T10:00:00Z"; |
| 1928 | let terminal = runtime_event( |
| 1929 | 2, |
| 1930 | timestamp, |
| 1931 | "thread-a", |
| 1932 | Some("turn-a"), |
| 1933 | "turn.completed", |
| 1934 | serde_json::json!({ |
| 1935 | "turn": { |
| 1936 | "usage": { "input_tokens": 10_000, "output_tokens": 9_000 }, |
| 1937 | "modelRequestDiagnostics": { |
| 1938 | "modelRequestsStarted": 2, |
| 1939 | "transparentStreamRetries": 1, |
| 1940 | "streamResumes": 1, |
| 1941 | }, |
| 1942 | }, |
| 1943 | }), |
| 1944 | ); |
| 1945 | let usage_one = runtime_event( |
| 1946 | 3, |
| 1947 | timestamp, |
| 1948 | "thread-a", |
| 1949 | Some("turn-a"), |
| 1950 | "turn.usage", |
| 1951 | serde_json::json!({ "usage": { "input_tokens": 7, "output_tokens": 2 } }), |
| 1952 | ); |
| 1953 | let usage_two = runtime_event( |
| 1954 | 4, |
| 1955 | timestamp, |
| 1956 | "thread-a", |
| 1957 | Some("turn-a"), |
| 1958 | "turn.usage", |
| 1959 | serde_json::json!({ "usage": { "input_tokens": 11, "output_tokens": 3 } }), |
| 1960 | ); |
| 1961 | let duplicate_terminal = runtime_event( |
| 1962 | 5, |
| 1963 | timestamp, |
| 1964 | "thread-a", |
| 1965 | Some("turn-a"), |
| 1966 | "turn.completed", |
| 1967 | serde_json::json!({ |
| 1968 | "turn": { |
| 1969 | "modelRequestDiagnostics": { |
| 1970 | "modelRequestsStarted": 99, |
| 1971 | "transparentStreamRetries": 99, |
| 1972 | "streamResumes": 99, |
| 1973 | }, |
| 1974 | }, |
| 1975 | }), |
| 1976 | ); |
| 1977 | let legacy_terminal = runtime_event( |
| 1978 | 6, |
| 1979 | timestamp, |
| 1980 | "thread-a", |
| 1981 | Some("turn-b"), |
| 1982 | "turn.completed", |
| 1983 | serde_json::json!({ "turn": { "usage": { "input_tokens": 50, "output_tokens": 5 } } }), |
| 1984 | ); |
| 1985 | let status = runtime_event( |
| 1986 | 7, |
| 1987 | timestamp, |
| 1988 | "thread-a", |
| 1989 | Some("turn-a"), |
| 1990 | "item.completed", |
| 1991 | serde_json::json!({ "item": { "kind": "status" } }), |
| 1992 | ); |
| 1993 | let tmp = write_runtime_events(&[ |
| 1994 | terminal, |
| 1995 | usage_one, |
| 1996 | usage_two, |
| 1997 | duplicate_terminal, |
| 1998 | legacy_terminal, |
| 1999 | status, |
| 2000 | ]); |
| 2001 | let mut rollup = Rollup::default(); |
| 2002 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2003 | |
| 2004 | let runtime = &rollup.runtime_requests; |
| 2005 | assert_eq!(runtime.terminal_turn_receipts, 2); |
| 2006 | assert_eq!(runtime.diagnostics_turn_receipts, 1); |
| 2007 | assert_eq!(runtime.diagnostics_unavailable_turn_receipts, 1); |
| 2008 | assert_eq!(runtime.duplicate_terminal_receipts_skipped, 1); |
| 2009 | assert_eq!(runtime.model_requests_started, 2); |
| 2010 | assert_eq!(runtime.transparent_stream_retries, 1); |
| 2011 | assert_eq!(runtime.stream_resumes, 1); |
| 2012 | assert_eq!(runtime.provider_usage_receipts, 2); |
| 2013 | assert_eq!(runtime.provider_reported_input_tokens, 18); |
| 2014 | assert_eq!(runtime.provider_reported_output_tokens, 5); |
| 2015 | assert_ne!(runtime.provider_reported_input_tokens, 10_018); |
| 2016 | assert_eq!(runtime.model_requests_started, 2, "status is not a request"); |
| 2017 | } |
| 2018 | |
| 2019 | #[test] |
| 2020 | fn runtime_usage_receipts_deduplicate_by_runtime_event_identity() { |
| 2021 | let usage = runtime_event( |
| 2022 | 20, |
| 2023 | "2026-09-08T10:00:00Z", |
| 2024 | "thread-a", |
| 2025 | Some("turn-a"), |
| 2026 | "turn.usage", |
| 2027 | serde_json::json!({ "usage": { "input_tokens": 7, "output_tokens": 2 } }), |
| 2028 | ); |
| 2029 | let tmp = write_runtime_events(&[usage.clone(), usage]); |
| 2030 | let mut rollup = Rollup::default(); |
| 2031 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2032 | |
| 2033 | let runtime = &rollup.runtime_requests; |
| 2034 | assert_eq!(runtime.provider_usage_receipts, 1); |
| 2035 | assert_eq!(runtime.provider_reported_input_tokens, 7); |
| 2036 | assert_eq!(runtime.provider_reported_output_tokens, 2); |
| 2037 | assert_eq!(runtime.duplicate_provider_usage_receipts_skipped, 1); |
| 2038 | } |
| 2039 | |
| 2040 | #[test] |
| 2041 | fn runtime_receipt_coverage_marks_unidentified_or_incomplete_old_records_unknown() { |
| 2042 | let terminal_without_identity = serde_json::json!({ |
| 2043 | "timestamp": "2026-09-08T10:00:00Z", |
| 2044 | "event": "turn.completed", |
| 2045 | "payload": { |
| 2046 | "turn": { |
| 2047 | "modelRequestDiagnostics": { |
| 2048 | "modelRequestsStarted": 3, |
| 2049 | "transparentStreamRetries": 1, |
| 2050 | "streamResumes": 2, |
| 2051 | }, |
| 2052 | }, |
| 2053 | }, |
| 2054 | }); |
| 2055 | let usage_without_identity = serde_json::json!({ |
| 2056 | "timestamp": "2026-09-08T10:00:00Z", |
| 2057 | "event": "turn.usage", |
| 2058 | "payload": { "usage": { "input_tokens": 9, "output_tokens": 4 } }, |
| 2059 | }); |
| 2060 | let incomplete_diagnostics = runtime_event( |
| 2061 | 30, |
| 2062 | "2026-09-08T10:00:00Z", |
| 2063 | "thread-a", |
| 2064 | Some("turn-b"), |
| 2065 | "turn.completed", |
| 2066 | serde_json::json!({ |
| 2067 | "turn": { "modelRequestDiagnostics": { "modelRequestsStarted": 3 } }, |
| 2068 | }), |
| 2069 | ); |
| 2070 | let incomplete_usage = runtime_event( |
| 2071 | 31, |
| 2072 | "2026-09-08T10:00:00Z", |
| 2073 | "thread-a", |
| 2074 | Some("turn-b"), |
| 2075 | "turn.usage", |
| 2076 | serde_json::json!({ "usage": { "input_tokens": 9 } }), |
| 2077 | ); |
| 2078 | let tmp = write_runtime_events(&[ |
| 2079 | terminal_without_identity, |
| 2080 | usage_without_identity, |
| 2081 | incomplete_diagnostics, |
| 2082 | incomplete_usage, |
| 2083 | ]); |
| 2084 | let mut rollup = Rollup::default(); |
| 2085 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2086 | |
| 2087 | let runtime = &rollup.runtime_requests; |
| 2088 | assert_eq!(runtime.terminal_receipts_without_identity, 1); |
| 2089 | assert_eq!(runtime.terminal_turn_receipts, 1); |
| 2090 | assert_eq!(runtime.diagnostics_incomplete_turn_receipts, 1); |
| 2091 | assert_eq!(runtime.model_requests_started, 0); |
| 2092 | assert_eq!(runtime.provider_usage_receipts_without_identity, 1); |
| 2093 | assert_eq!(runtime.provider_usage_receipts_incomplete, 1); |
| 2094 | assert_eq!(runtime.provider_usage_receipts, 0); |
| 2095 | assert_eq!(runtime.provider_reported_input_tokens, 0); |
| 2096 | } |
| 2097 | |
| 2098 | #[test] |
| 2099 | fn runtime_receipts_respect_since_cutoff_without_crossing_snapshot_boundaries() { |
| 2100 | let old_terminal = runtime_event( |
| 2101 | 40, |
| 2102 | "2026-09-01T10:00:00Z", |
| 2103 | "thread-a", |
| 2104 | Some("turn-old"), |
| 2105 | "turn.completed", |
| 2106 | serde_json::json!({ |
| 2107 | "turn": { "modelRequestDiagnostics": { |
| 2108 | "modelRequestsStarted": 4, |
| 2109 | "transparentStreamRetries": 1, |
| 2110 | "streamResumes": 2, |
| 2111 | } }, |
| 2112 | }), |
| 2113 | ); |
| 2114 | let old_usage = runtime_event( |
| 2115 | 41, |
| 2116 | "2026-09-01T10:00:00Z", |
| 2117 | "thread-a", |
| 2118 | Some("turn-old"), |
| 2119 | "turn.usage", |
| 2120 | serde_json::json!({ "usage": { "input_tokens": 40, "output_tokens": 4 } }), |
| 2121 | ); |
| 2122 | let new_terminal = runtime_event( |
| 2123 | 42, |
| 2124 | "2026-09-08T10:00:00Z", |
| 2125 | "thread-a", |
| 2126 | Some("turn-new"), |
| 2127 | "turn.completed", |
| 2128 | serde_json::json!({ |
| 2129 | "turn": { "modelRequestDiagnostics": { |
| 2130 | "modelRequestsStarted": 1, |
| 2131 | "transparentStreamRetries": 0, |
| 2132 | "streamResumes": 0, |
| 2133 | } }, |
| 2134 | }), |
| 2135 | ); |
| 2136 | let new_usage = runtime_event( |
| 2137 | 43, |
| 2138 | "2026-09-08T10:00:00Z", |
| 2139 | "thread-a", |
| 2140 | Some("turn-new"), |
| 2141 | "turn.usage", |
| 2142 | serde_json::json!({ "usage": { "input_tokens": 10, "output_tokens": 1 } }), |
| 2143 | ); |
| 2144 | let tmp = write_runtime_events(&[old_terminal, old_usage, new_terminal, new_usage]); |
| 2145 | let mut rollup = Rollup::default(); |
| 2146 | read_runtime_test_log( |
| 2147 | tmp.path(), |
| 2148 | Some("2026-09-08T00:00:00Z".parse().unwrap()), |
| 2149 | &mut rollup, |
| 2150 | ); |
| 2151 | |
| 2152 | let runtime = &rollup.runtime_requests; |
| 2153 | assert_eq!(runtime.terminal_turn_receipts, 1); |
| 2154 | assert_eq!(runtime.model_requests_started, 1); |
| 2155 | assert_eq!(runtime.provider_usage_receipts, 1); |
| 2156 | assert_eq!(runtime.provider_reported_input_tokens, 10); |
| 2157 | assert_eq!(runtime.provider_reported_output_tokens, 1); |
| 2158 | } |
| 2159 | |
| 2160 | // ── Durable runtime item receipts ── |
| 2161 | // |
| 2162 | // These pin *which* runtime event names carry tool and compaction data. |
| 2163 | // Before the fix this reader matched `tool.started` / `tool.completed` / |
| 2164 | // `tool.failed` and `compaction.completed`, none of which the Runtime |
| 2165 | // store has ever written, so every per-tool counter was structurally 0. |
| 2166 | |
| 2167 | fn tool_item(kind: &str, tool_name: &str, extra: Value) -> Value { |
| 2168 | let mut item = serde_json::json!({ |
| 2169 | "schema_version": 4, |
| 2170 | "id": "item_abc", |
| 2171 | "turn_id": "turn-a", |
| 2172 | "kind": kind, |
| 2173 | "status": "completed", |
| 2174 | "summary": "exec_shell: ok", |
| 2175 | "metadata": { "tool_use_id": "call-1", "tool_name": tool_name }, |
| 2176 | "started_at": "2026-09-08T10:00:00Z", |
| 2177 | }); |
| 2178 | merge_json(&mut item, extra); |
| 2179 | item |
| 2180 | } |
| 2181 | |
| 2182 | fn merge_json(target: &mut Value, extra: Value) { |
| 2183 | let Value::Object(extra) = extra else { return }; |
| 2184 | let Some(target) = target.as_object_mut() else { |
| 2185 | return; |
| 2186 | }; |
| 2187 | for (key, value) in extra { |
| 2188 | let nested = matches!(value, Value::Object(_)) |
| 2189 | && matches!(target.get(&key), Some(Value::Object(_))); |
| 2190 | if nested { |
| 2191 | merge_json(target.get_mut(&key).expect("checked above"), value); |
| 2192 | } else { |
| 2193 | target.insert(key, value); |
| 2194 | } |
| 2195 | } |
| 2196 | } |
| 2197 | |
| 2198 | #[test] |
| 2199 | fn runtime_tool_receipts_come_from_durable_item_events() { |
| 2200 | let started = runtime_event( |
| 2201 | 60, |
| 2202 | "2026-09-08T10:00:00Z", |
| 2203 | "thread-a", |
| 2204 | Some("turn-a"), |
| 2205 | "item.started", |
| 2206 | serde_json::json!({ |
| 2207 | "item": tool_item("tool_call", "exec_shell", serde_json::json!({ |
| 2208 | "status": "in_progress", |
| 2209 | "metadata": { "tool_input": "{}" }, |
| 2210 | })), |
| 2211 | "tool": { "id": "call-1", "name": "exec_shell", "input": {} }, |
| 2212 | }), |
| 2213 | ); |
| 2214 | let completed = runtime_event( |
| 2215 | 61, |
| 2216 | "2026-09-08T10:00:02Z", |
| 2217 | "thread-a", |
| 2218 | Some("turn-a"), |
| 2219 | "item.completed", |
| 2220 | serde_json::json!({ |
| 2221 | "item": tool_item("tool_call", "exec_shell", serde_json::json!({ |
| 2222 | "ended_at": "2026-09-08T10:00:02Z", |
| 2223 | "metadata": { "is_error": false }, |
| 2224 | })), |
| 2225 | }), |
| 2226 | ); |
| 2227 | let tmp = write_runtime_events(&[started, completed]); |
| 2228 | let mut rollup = Rollup::default(); |
| 2229 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2230 | |
| 2231 | let stats = &rollup.tools["exec_shell"]; |
| 2232 | assert_eq!(stats.calls, 1); |
| 2233 | assert_eq!(stats.successes, 1); |
| 2234 | assert_eq!(stats.failures, 0); |
| 2235 | assert_eq!(stats.outcome_unknown, 0); |
| 2236 | assert_eq!(stats.elapsed_samples, 1); |
| 2237 | assert_eq!(stats.total_elapsed_ms, 2_000); |
| 2238 | assert_eq!(stats.elapsed_unavailable, 0); |
| 2239 | } |
| 2240 | |
| 2241 | #[test] |
| 2242 | fn runtime_file_change_and_command_execution_items_count_as_tools() { |
| 2243 | // `tool_kind_for_name` splits one tool call across three item kinds; |
| 2244 | // dropping two of them would hide every shell and edit receipt. |
| 2245 | let events: Vec<_> = [ |
| 2246 | ("file_change", "apply_patch"), |
| 2247 | ("command_execution", "exec_shell"), |
| 2248 | ] |
| 2249 | .into_iter() |
| 2250 | .enumerate() |
| 2251 | .map(|(i, (kind, name))| { |
| 2252 | runtime_event( |
| 2253 | 70 + i as u64, |
| 2254 | "2026-09-08T10:00:00Z", |
| 2255 | "thread-a", |
| 2256 | Some("turn-a"), |
| 2257 | "item.completed", |
| 2258 | serde_json::json!({ |
| 2259 | "item": tool_item(kind, name, serde_json::json!({ |
| 2260 | "ended_at": "2026-09-08T10:00:01Z", |
| 2261 | "metadata": { "is_error": false }, |
| 2262 | })), |
| 2263 | }), |
| 2264 | ) |
| 2265 | }) |
| 2266 | .collect(); |
| 2267 | let tmp = write_runtime_events(&events); |
| 2268 | let mut rollup = Rollup::default(); |
| 2269 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2270 | |
| 2271 | assert_eq!(rollup.tools["apply_patch"].successes, 1); |
| 2272 | assert_eq!(rollup.tools["exec_shell"].successes, 1); |
| 2273 | assert_eq!(rollup.tools["exec_shell"].total_elapsed_ms, 1_000); |
| 2274 | } |
| 2275 | |
| 2276 | #[test] |
| 2277 | fn runtime_tool_outcome_without_is_error_is_unknown_not_success() { |
| 2278 | let completed = runtime_event( |
| 2279 | 80, |
| 2280 | "2026-09-08T10:00:00Z", |
| 2281 | "thread-a", |
| 2282 | Some("turn-a"), |
| 2283 | "item.completed", |
| 2284 | serde_json::json!({ |
| 2285 | "item": tool_item("tool_call", "exec_shell", serde_json::json!({})), |
| 2286 | }), |
| 2287 | ); |
| 2288 | let tmp = write_runtime_events(&[completed]); |
| 2289 | let mut rollup = Rollup::default(); |
| 2290 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2291 | |
| 2292 | let stats = &rollup.tools["exec_shell"]; |
| 2293 | assert_eq!(stats.successes, 0); |
| 2294 | assert_eq!(stats.failures, 0, "unknown is never folded into failures"); |
| 2295 | assert_eq!(stats.outcome_unknown, 1); |
| 2296 | assert_eq!(stats.success_rate_pct(), None); |
| 2297 | assert_eq!( |
| 2298 | stats.elapsed_unavailable, 1, |
| 2299 | "a missing ended_at is not a 0 ms call" |
| 2300 | ); |
| 2301 | assert_eq!(stats.elapsed_samples, 0); |
| 2302 | assert_eq!(stats.avg_elapsed_ms(), None); |
| 2303 | } |
| 2304 | |
| 2305 | #[test] |
| 2306 | fn sse_only_tool_event_names_are_not_durable_receipts() { |
| 2307 | // `tool.started` / `tool.completed` / `tool.failed` are synthesized by |
| 2308 | // `map_compat_stream_event` for HTTP clients and never persisted. |
| 2309 | let events: Vec<_> = ["tool.started", "tool.completed", "tool.failed"] |
| 2310 | .into_iter() |
| 2311 | .enumerate() |
| 2312 | .map(|(i, event)| { |
| 2313 | runtime_event( |
| 2314 | 90 + i as u64, |
| 2315 | "2026-09-08T10:00:00Z", |
| 2316 | "thread-a", |
| 2317 | Some("turn-a"), |
| 2318 | event, |
| 2319 | serde_json::json!({ "tool_name": "exec_shell", "elapsed_ms": 5 }), |
| 2320 | ) |
| 2321 | }) |
| 2322 | .collect(); |
| 2323 | let tmp = write_runtime_events(&events); |
| 2324 | let mut rollup = Rollup::default(); |
| 2325 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2326 | |
| 2327 | assert_eq!(rollup.total_tool_calls(), 0); |
| 2328 | assert!(rollup.tools.is_empty()); |
| 2329 | } |
| 2330 | |
| 2331 | #[test] |
| 2332 | fn duplicate_item_receipts_are_counted_once() { |
| 2333 | let completed = runtime_event( |
| 2334 | 100, |
| 2335 | "2026-09-08T10:00:00Z", |
| 2336 | "thread-a", |
| 2337 | Some("turn-a"), |
| 2338 | "item.completed", |
| 2339 | serde_json::json!({ |
| 2340 | "item": tool_item("tool_call", "exec_shell", serde_json::json!({ |
| 2341 | "ended_at": "2026-09-08T10:00:01Z", |
| 2342 | "metadata": { "is_error": false }, |
| 2343 | })), |
| 2344 | }), |
| 2345 | ); |
| 2346 | let tmp = write_runtime_events(&[completed.clone(), completed]); |
| 2347 | let mut rollup = Rollup::default(); |
| 2348 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2349 | |
| 2350 | let stats = &rollup.tools["exec_shell"]; |
| 2351 | assert_eq!(stats.successes, 1); |
| 2352 | assert_eq!(stats.elapsed_samples, 1); |
| 2353 | assert_eq!(stats.total_elapsed_ms, 1_000); |
| 2354 | } |
| 2355 | |
| 2356 | #[test] |
| 2357 | fn compaction_reduction_is_computed_from_message_counts() { |
| 2358 | let completed = runtime_event( |
| 2359 | 110, |
| 2360 | "2026-09-08T10:00:00Z", |
| 2361 | "thread-a", |
| 2362 | Some("turn-a"), |
| 2363 | "item.completed", |
| 2364 | serde_json::json!({ |
| 2365 | "item": { "kind": "context_compaction", "status": "completed" }, |
| 2366 | "auto": true, |
| 2367 | "messages_before": 40, |
| 2368 | "messages_after": 10, |
| 2369 | }), |
| 2370 | ); |
| 2371 | let tmp = write_runtime_events(&[completed]); |
| 2372 | let mut rollup = Rollup::default(); |
| 2373 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2374 | |
| 2375 | assert_eq!(rollup.compaction.events, 1); |
| 2376 | assert_eq!(rollup.compaction.ratio_samples, 1); |
| 2377 | assert_eq!(rollup.compaction.avg_reduction_pct(), Some(75.0)); |
| 2378 | } |
| 2379 | |
| 2380 | #[test] |
| 2381 | fn compaction_without_message_counts_stays_unknown() { |
| 2382 | let completed = runtime_event( |
| 2383 | 120, |
| 2384 | "2026-09-08T10:00:00Z", |
| 2385 | "thread-a", |
| 2386 | Some("turn-a"), |
| 2387 | "item.completed", |
| 2388 | serde_json::json!({ |
| 2389 | "item": { "kind": "context_compaction", "status": "completed" }, |
| 2390 | "auto": true, |
| 2391 | }), |
| 2392 | ); |
| 2393 | let tmp = write_runtime_events(&[completed]); |
| 2394 | let mut rollup = Rollup::default(); |
| 2395 | read_runtime_test_log(tmp.path(), None, &mut rollup); |
| 2396 | |
| 2397 | assert_eq!(rollup.compaction.events, 1); |
| 2398 | assert_eq!(rollup.compaction.ratio_samples, 0); |
| 2399 | assert_eq!( |
| 2400 | rollup.compaction.avg_reduction_pct(), |
| 2401 | None, |
| 2402 | "no counts is unknown, never a 0% reduction" |
| 2403 | ); |
| 2404 | } |
| 2405 | |
| 2406 | // ── Approval receipts ── |
| 2407 | |
| 2408 | #[test] |
| 2409 | fn session_auto_approvals_are_counted_under_their_emitted_name() { |
| 2410 | let events = [ |
| 2411 | serde_json::json!({ |
| 2412 | "ts": "2026-09-08T10:00:00Z", |
| 2413 | "event": "tool.approval.auto_approve_session", |
| 2414 | "details": { "tool_name": "exec_shell" }, |
| 2415 | }), |
| 2416 | serde_json::json!({ |
| 2417 | "ts": "2026-09-08T10:00:01Z", |
| 2418 | "event": "tool.approval.auto_approve", |
| 2419 | "details": { "tool_name": "exec_shell" }, |
| 2420 | }), |
| 2421 | ]; |
| 2422 | let tmp = write_runtime_events(&events); |
| 2423 | let mut rollup = Rollup::default(); |
| 2424 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 2425 | |
| 2426 | let stats = &rollup.tools["exec_shell"]; |
| 2427 | assert_eq!(stats.auto_approved, 2, "emitted name plus legacy alias"); |
| 2428 | assert_eq!(stats.calls, 2); |
| 2429 | } |
| 2430 | |
| 2431 | #[test] |
| 2432 | fn tool_denials_are_a_class_of_their_own() { |
| 2433 | let events: Vec<_> = [ |
| 2434 | "tool.approval.auto_deny", |
| 2435 | "tool.approval.auto_deny_session", |
| 2436 | "tool.approval.auto_deny_auto_review", |
| 2437 | "tool.approval.auto_deny_full_access_policy", |
| 2438 | ] |
| 2439 | .into_iter() |
| 2440 | .map(|event| { |
| 2441 | serde_json::json!({ |
| 2442 | "ts": "2026-09-08T10:00:00Z", |
| 2443 | "event": event, |
| 2444 | "details": { "tool_name": "exec_shell" }, |
| 2445 | }) |
| 2446 | }) |
| 2447 | .collect(); |
| 2448 | let tmp = write_runtime_events(&events); |
| 2449 | let mut rollup = Rollup::default(); |
| 2450 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 2451 | |
| 2452 | let stats = &rollup.tools["exec_shell"]; |
| 2453 | assert_eq!(stats.denied, 4); |
| 2454 | assert_eq!(stats.calls, 4); |
| 2455 | assert_eq!(stats.successes, 0); |
| 2456 | assert_eq!(stats.failures, 0); |
| 2457 | assert_eq!( |
| 2458 | stats.success_rate_pct(), |
| 2459 | None, |
| 2460 | "a blocked call is not a judged outcome" |
| 2461 | ); |
| 2462 | } |
| 2463 | |
| 2464 | #[test] |
| 2465 | fn opt_in_tool_audit_records_resolve_their_top_level_tool_name() { |
| 2466 | // `emit_tool_audit` writes `tool_name` at the top level, not under |
| 2467 | // `details`, so pointing `--since` at that file used to bucket every |
| 2468 | // record as "unknown" and grade an absent outcome as a success. |
| 2469 | let events = [ |
| 2470 | serde_json::json!({ |
| 2471 | "event": "tool.result", |
| 2472 | "tool_id": "call-1", |
| 2473 | "tool_name": "exec_shell", |
| 2474 | "success": false, |
| 2475 | }), |
| 2476 | serde_json::json!({ |
| 2477 | "event": "tool.result", |
| 2478 | "tool_id": "call-2", |
| 2479 | "tool_name": "exec_shell", |
| 2480 | }), |
| 2481 | ]; |
| 2482 | let tmp = write_runtime_events(&events); |
| 2483 | let mut rollup = Rollup::default(); |
| 2484 | read_audit_test_log(tmp.path(), None, &mut rollup); |
| 2485 | |
| 2486 | let stats = &rollup.tools["exec_shell"]; |
| 2487 | assert_eq!(stats.calls, 2); |
| 2488 | assert_eq!(stats.failures, 1); |
| 2489 | assert_eq!(stats.successes, 0); |
| 2490 | assert_eq!(stats.outcome_unknown, 1); |
| 2491 | assert!(!rollup.tools.contains_key("unknown")); |
| 2492 | } |
| 2493 | |
| 2494 | #[test] |
| 2495 | fn rollup_json_exposes_the_new_tool_classes() { |
| 2496 | let mut rollup = Rollup::default(); |
| 2497 | let stats = rollup.tool_mut("exec_shell"); |
| 2498 | stats.denied = 2; |
| 2499 | stats.outcome_unknown = 1; |
| 2500 | stats.elapsed_unavailable = 3; |
| 2501 | let json = serde_json::to_value(&rollup).unwrap(); |
| 2502 | assert_eq!(json["tools"]["exec_shell"]["denied"], 2); |
| 2503 | assert_eq!(json["tools"]["exec_shell"]["outcome_unknown"], 1); |
| 2504 | assert_eq!(json["tools"]["exec_shell"]["elapsed_unavailable"], 3); |
| 2505 | // Existing keys keep their names and positions for JSON consumers. |
| 2506 | assert_eq!(json["tools"]["exec_shell"]["calls"], 0); |
| 2507 | assert_eq!(json["tools"]["exec_shell"]["successes"], 0); |
| 2508 | assert_eq!(json["tools"]["exec_shell"]["failures"], 0); |
| 2509 | } |
| 2510 | |
| 2511 | // ── Runtime store roots ── |
| 2512 | |
| 2513 | #[test] |
| 2514 | fn every_runtime_store_root_is_read_including_session_scoped_stores() { |
| 2515 | let dir = tempfile::TempDir::new().unwrap(); |
| 2516 | let tasks = dir.path().join("tasks"); |
| 2517 | let sessions = dir.path().join("sessions"); |
| 2518 | std::fs::create_dir_all(sessions.join("sess-1").join("runtime").join("events")).unwrap(); |
| 2519 | std::fs::create_dir_all(&tasks).unwrap(); |
| 2520 | std::fs::write(sessions.join("loose.json"), "{}").unwrap(); |
| 2521 | |
| 2522 | let _lock = crate::tests::env_lock(); |
| 2523 | let _override = crate::tests::ScopedEnvVar::remove("CODEWHALE_RUNTIME_DIR"); |
| 2524 | let _legacy = crate::tests::ScopedEnvVar::remove("DEEPSEEK_RUNTIME_DIR"); |
| 2525 | assert_eq!( |
| 2526 | runtime_event_dirs(&tasks, &sessions), |
| 2527 | vec![ |
| 2528 | tasks.join("runtime").join("events"), |
| 2529 | sessions.join("sess-1").join("runtime").join("events"), |
| 2530 | ], |
| 2531 | "a loose session file is not a store root" |
| 2532 | ); |
| 2533 | } |
| 2534 | |
| 2535 | #[test] |
| 2536 | fn an_explicit_runtime_dir_override_is_the_only_root_read() { |
| 2537 | let dir = tempfile::TempDir::new().unwrap(); |
| 2538 | let override_dir = dir.path().join("elsewhere"); |
| 2539 | let _lock = crate::tests::env_lock(); |
| 2540 | let _override = crate::tests::ScopedEnvVar::set( |
| 2541 | "CODEWHALE_RUNTIME_DIR", |
| 2542 | &override_dir.to_string_lossy(), |
| 2543 | ); |
| 2544 | assert_eq!( |
| 2545 | runtime_event_dirs(&dir.path().join("tasks"), &dir.path().join("sessions")), |
| 2546 | vec![override_dir.join("events")], |
| 2547 | "mixing an override with the default roots would double count" |
| 2548 | ); |
| 2549 | } |
| 2550 | |
| 2551 | // ── State-root resolution ── |
| 2552 | // |
| 2553 | // These pin *which* files the rollup reads. Before the fix the reader |
| 2554 | // resolved `$HOME/.deepseek`, which nothing has written since the v0.8.44 |
| 2555 | // rename, so `codewhale metrics` printed an all-zero rollup as truth. |
| 2556 | |
| 2557 | /// Isolate the ambient home so the resolver sees a clean, empty install. |
| 2558 | /// |
| 2559 | /// The returned guards must stay bound for the life of the test: dropping |
| 2560 | /// them restores the previous environment. Destructure the tuple so the |
| 2561 | /// bindings drop in reverse order — the environment is restored *before* |
| 2562 | /// the lock is released, or a concurrent env-mutating test sees a torn HOME. |
| 2563 | fn isolated_home() -> ( |
| 2564 | tempfile::TempDir, |
| 2565 | std::sync::MutexGuard<'static, ()>, |
| 2566 | Vec<crate::tests::ScopedEnvVar>, |
| 2567 | ) { |
| 2568 | let guard = crate::tests::env_lock(); |
| 2569 | let home = tempfile::TempDir::new().expect("tempdir"); |
| 2570 | let vars = vec![ |
| 2571 | crate::tests::ScopedEnvVar::set("HOME", &home.path().to_string_lossy()), |
| 2572 | crate::tests::ScopedEnvVar::set("USERPROFILE", &home.path().to_string_lossy()), |
| 2573 | crate::tests::ScopedEnvVar::remove("CODEWHALE_HOME"), |
| 2574 | crate::tests::ScopedEnvVar::remove("DEEPSEEK_HOME"), |
| 2575 | ]; |
| 2576 | (home, guard, vars) |
| 2577 | } |
| 2578 | |
| 2579 | #[test] |
| 2580 | fn default_audit_history_includes_both_roots_without_requiring_existing_files() { |
| 2581 | let (home, _lock, _env) = isolated_home(); |
| 2582 | assert_eq!( |
| 2583 | resolve_audit_roots().expect("resolves"), |
| 2584 | vec![ |
| 2585 | home.path().join(".codewhale"), |
| 2586 | home.path().join(".deepseek") |
| 2587 | ], |
| 2588 | ); |
| 2589 | } |
| 2590 | |
| 2591 | #[test] |
| 2592 | fn copied_audit_history_keeps_unique_legacy_and_rotated_records() { |
| 2593 | let dir = tempfile::TempDir::new().expect("tempdir"); |
| 2594 | let primary = dir.path().join("primary"); |
| 2595 | let legacy = dir.path().join("legacy"); |
| 2596 | std::fs::create_dir_all(&primary).unwrap(); |
| 2597 | std::fs::create_dir_all(&legacy).unwrap(); |
| 2598 | let shared = r#"{"ts":"2026-09-01T00:00:00Z","event":"credential.save","details":{}}"#; |
| 2599 | let old = r#"{"ts":"2026-08-01T00:00:00Z","event":"credential.clear","details":{}}"#; |
| 2600 | let new = r#"{"ts":"2026-09-02T00:00:00Z","event":"credential.save","details":{}}"#; |
| 2601 | std::fs::write(primary.join("audit.log.1"), format!("{shared}\n{shared}\n")).unwrap(); |
| 2602 | std::fs::write(primary.join("audit.log"), format!("{new}\nmalformed\n")).unwrap(); |
| 2603 | std::fs::write( |
| 2604 | legacy.join("audit.log"), |
| 2605 | format!("{shared}\n{shared}\n{shared}\n"), |
| 2606 | ) |
| 2607 | .unwrap(); |
| 2608 | std::fs::write(legacy.join("audit.log.1"), format!("{old}\n")).unwrap(); |
| 2609 | let roots = [primary, legacy]; |
| 2610 | let before: Vec<_> = roots |
| 2611 | .iter() |
| 2612 | .flat_map(|root| { |
| 2613 | ["audit.log.1", "audit.log"].map(|name| { |
| 2614 | let path = root.join(name); |
| 2615 | (path.clone(), std::fs::read(path).unwrap()) |
| 2616 | }) |
| 2617 | }) |
| 2618 | .collect(); |
| 2619 | let mut rollup = Rollup::default(); |
| 2620 | read_audit_history(&roots, None, &mut rollup); |
| 2621 | assert_eq!( |
| 2622 | rollup.credentials.saves, 4, |
| 2623 | "maximum occurrence count across copied roots" |
| 2624 | ); |
| 2625 | assert_eq!( |
| 2626 | rollup.credentials.clears, 1, |
| 2627 | "unique old rotation is retained" |
| 2628 | ); |
| 2629 | assert_eq!(rollup.parsed_lines, 5); |
| 2630 | for (path, bytes) in before { |
| 2631 | assert_eq!( |
| 2632 | std::fs::read(path).unwrap(), |
| 2633 | bytes, |
| 2634 | "source history is read-only" |
| 2635 | ); |
| 2636 | } |
| 2637 | let mut recent = Rollup::default(); |
| 2638 | read_audit_history( |
| 2639 | &roots, |
| 2640 | Some("2026-09-01T00:00:00Z".parse().unwrap()), |
| 2641 | &mut recent, |
| 2642 | ); |
| 2643 | assert_eq!(recent.credentials.saves, 4); |
| 2644 | assert_eq!(recent.credentials.clears, 0); |
| 2645 | } |
| 2646 | |
| 2647 | #[test] |
| 2648 | fn an_explicit_codewhale_home_is_an_audit_isolation_boundary() { |
| 2649 | let (home, _lock, _env) = isolated_home(); |
| 2650 | let legacy = home.path().join(".deepseek"); |
| 2651 | std::fs::create_dir_all(&legacy).unwrap(); |
| 2652 | std::fs::write(legacy.join("audit.log"), r#"{"event":"credential.save"}"#).unwrap(); |
| 2653 | let explicit = tempfile::TempDir::new().unwrap(); |
| 2654 | let _pin = |
| 2655 | crate::tests::ScopedEnvVar::set("CODEWHALE_HOME", &explicit.path().to_string_lossy()); |
| 2656 | let roots = resolve_audit_roots().unwrap(); |
| 2657 | assert_eq!(roots, vec![explicit.path().to_path_buf()]); |
| 2658 | let mut rollup = Rollup::default(); |
| 2659 | read_audit_history(&roots, None, &mut rollup); |
| 2660 | assert_eq!(rollup.parsed_lines, 0); |
| 2661 | } |
| 2662 | |
| 2663 | #[test] |
| 2664 | fn the_legacy_deepseek_home_variable_is_no_longer_honoured() { |
| 2665 | let (home, _lock, _env) = isolated_home(); |
| 2666 | let stale = tempfile::TempDir::new().unwrap(); |
| 2667 | let _stale = |
| 2668 | crate::tests::ScopedEnvVar::set("DEEPSEEK_HOME", &stale.path().to_string_lossy()); |
| 2669 | assert_eq!( |
| 2670 | resolve_audit_roots().unwrap(), |
| 2671 | vec![ |
| 2672 | home.path().join(".codewhale"), |
| 2673 | home.path().join(".deepseek") |
| 2674 | ], |
| 2675 | ); |
| 2676 | } |
| 2677 | } |
| 2678 |