1use std::{
4 fs::File,
5 io::{BufRead, BufReader, Cursor},
6 path::Path,
7};
8
9use anyhow::{Context, Result};
10use serde_json::Value;
11
12mod metrics;
13mod model;
14mod stream;
15
16use metrics::{LatencyAccumulator, LifecycleTiming, UsageAccounting};
17pub use model::{HarnessTraceSummary, LatencyStatistics, TokenUsage};
18
19const GENERIC_ERROR_CATEGORY: &str = "error";
20const MAX_DISTINCT_TOOL_LABELS: usize = 256;
21const MAX_TRACE_LINE_BYTES: usize = 1_048_576;
22const OTHER_TOOL_LABEL: &str = "other_tool";
23
24pub fn analyze_jsonl(input: &str) -> Result<HarnessTraceSummary> {
26 analyze_jsonl_reader(Cursor::new(input.as_bytes()))
27}
28
29pub fn analyze_jsonl_reader<R: BufRead>(reader: R) -> Result<HarnessTraceSummary> {
31 stream::analyze_jsonl_reader(reader)
32}
33
34pub fn analyze_jsonl_file(path: impl AsRef<Path>) -> Result<HarnessTraceSummary> {
36 let path = path.as_ref();
37 let file = File::open(path).with_context(|| format!("read trace file {}", path.display()))?;
38 analyze_jsonl_reader(BufReader::new(file)).with_context(|| format!("analyze trace file {}", path.display()))
39}
40
41fn record_value(
42 value: &Value,
43 summary: &mut HarnessTraceSummary,
44 latencies: &mut LatencyAccumulator,
45 timing: &mut LifecycleTiming,
46 usage: &mut UsageAccounting,
47) -> bool {
48 let Some(object) = value.as_object() else {
49 return false;
50 };
51 let payload = object
52 .get("data")
53 .and_then(Value::as_object)
54 .or_else(|| object.get("event").and_then(Value::as_object))
55 .unwrap_or(object);
56 let event_type =
57 string_field(object, &["type", "event", "kind"]).or_else(|| string_field(payload, &["type", "event", "kind"]));
58 let is_thread_event = event_type.is_some_and(is_known_event_type);
59
60 let mut recognized = is_thread_event;
61 let mut error_recorded = false;
62 timing.record(event_type, object, payload, latencies);
63 if matches!(event_type, Some("turn.started" | "turn/start")) {
64 summary.turns = summary.turns.saturating_add(1);
65 }
66 if matches!(event_type, Some("thread.completed"))
67 && let Some(num_turns) = number_field_from(object, payload, &["num_turns"])
68 {
69 summary.turns = summary.turns.max(num_turns);
70 }
71 if matches!(event_type, Some("error" | "turn.failed")) {
72 add_error(summary, error_category(payload, object).unwrap_or(GENERIC_ERROR_CATEGORY));
73 recognized = true;
74 error_recorded = true;
75 }
76
77 if matches!(event_type, Some("tool/result")) {
78 if let Some(bytes) = tool_result_output_bytes(payload) {
79 summary.output_bytes = summary.output_bytes.saturating_add(bytes);
80 }
81 if let Some(category) = tool_result_error_category(payload) {
82 let category = canonical_error_category(category);
83 add_error(
84 summary,
85 if category == GENERIC_ERROR_CATEGORY {
86 "tool_error"
87 } else {
88 category
89 },
90 );
91 error_recorded = true;
92 }
93 recognized = true;
94 }
95
96 if let Some(item) = payload.get("item").and_then(Value::as_object) {
97 if matches!(event_type, Some("item.completed")) {
98 recognized |= record_item(item, summary);
99 if let Some(bytes) = output_bytes(item) {
100 summary.output_bytes = summary.output_bytes.saturating_add(bytes);
101 }
102 } else {
103 recognized = true;
104 }
105 }
106
107 let deepseek_tool = string_field_from(payload, object, &["tool", "tool_name", "name"])
108 .or_else(|| {
109 payload
110 .get("function")
111 .and_then(Value::as_object)
112 .and_then(|f| string_field(f, &["name"]))
113 })
114 .or_else(|| {
115 object
116 .get("function")
117 .and_then(Value::as_object)
118 .and_then(|f| string_field(f, &["name"]))
119 });
120 let has_step = payload.contains_key("step")
121 || payload.contains_key("step_id")
122 || object.contains_key("step")
123 || object.contains_key("step_id");
124 let is_step_start = matches!(event_type, Some("step/start"));
125 let is_synthetic_record = event_type.is_none();
126 let is_tool_call = is_synthetic_record || matches!(event_type, Some("tool" | "tool_call" | "tool/call"));
127 if is_step_start || (is_synthetic_record && (deepseek_tool.is_some() || has_step)) {
128 recognized = true;
129 if is_step_start || has_step {
130 summary.steps = summary.steps.saturating_add(1);
131 }
132 }
133 if is_tool_call && let Some(tool) = deepseek_tool {
134 recognized = true;
135 record_tool(summary, tool);
136 }
137
138 if let Some(latency) = number_field_from(object, payload, &["latency_ms", "duration_ms", "latency"]) {
139 latencies.record(latency);
140 recognized = true;
141 }
142 if is_tool_call && let Some(bytes) = output_bytes(payload).or_else(|| output_bytes(object)) {
143 summary.output_bytes = summary.output_bytes.saturating_add(bytes);
144 recognized = true;
145 }
146 if let Some(category) = error_category(payload, object) {
147 if !error_recorded {
148 add_error(summary, category);
149 }
150 recognized = true;
151 }
152 usage.record(event_type, payload);
153 recognized || has_usage(payload)
154}
155
156fn record_item(item: &serde_json::Map<String, Value>, summary: &mut HarnessTraceSummary) -> bool {
157 let Some(details_type) = item.get("type").and_then(Value::as_str) else {
158 return false;
159 };
160 match details_type {
161 "tool_invocation" | "mcp_tool_call" => {
162 summary.steps = summary.steps.saturating_add(1);
163 if let Some(tool) = string_field(item, &["tool_name", "name"]) {
164 record_tool(summary, tool);
165 }
166 if string_field(item, &["status", "outcome"])
167 .is_some_and(|status| status != "completed" && status != "success")
168 {
169 add_error(summary, string_field(item, &["outcome", "status"]).unwrap_or("tool_error"));
170 }
171 true
172 }
173 "command_execution" => {
174 summary.steps = summary.steps.saturating_add(1);
175 record_tool(summary, "command_execution");
176 if item
177 .get("status")
178 .and_then(Value::as_str)
179 .is_some_and(|status| status == "failed")
180 {
181 add_error(summary, "command_failed");
182 }
183 true
184 }
185 "tool_output" => true,
186 "harness" => {
187 if let Some(event) = item.get("event").and_then(Value::as_str)
188 && event.ends_with("failed")
189 {
190 add_error(summary, item.get("error_category").and_then(Value::as_str).unwrap_or(event));
191 }
192 true
193 }
194 "error" => {
195 add_error(summary, "error");
196 true
197 }
198 _ => false,
199 }
200}
201
202fn is_known_event_type(event: &str) -> bool {
203 event.starts_with("thread.")
204 || event.starts_with("turn.")
205 || event.starts_with("item.")
206 || event.starts_with("turn/")
207 || event.starts_with("step/")
208 || event.starts_with("agent/")
209 || event.starts_with("agent-preset/")
210 || event.starts_with("approval/")
211 || event.starts_with("assistant/")
212 || event.starts_with("command/")
213 || event.starts_with("goal/")
214 || event.starts_with("permission/")
215 || event.starts_with("request/")
216 || event.starts_with("sandbox/")
217 || event.starts_with("session/")
218 || event.starts_with("todo/")
219 || event.starts_with("user/")
220 || event.starts_with("web/")
221 || matches!(
222 event,
223 "error"
224 | "context.reset"
225 | "permission.requested"
226 | "permission.resolved"
227 | "tool/call"
228 | "tool/result"
229 | "assistant/message"
230 | "reasoning-chunks"
231 | "session"
232 | "text-chunks"
233 | "tool-call-chunks"
234 )
235}
236
237fn string_field_from<'a>(
238 primary: &'a serde_json::Map<String, Value>,
239 fallback: &'a serde_json::Map<String, Value>,
240 names: &[&str],
241) -> Option<&'a str> {
242 string_field(primary, names).or_else(|| string_field(fallback, names))
243}
244
245fn number_field_from(
246 primary: &serde_json::Map<String, Value>,
247 fallback: &serde_json::Map<String, Value>,
248 names: &[&str],
249) -> Option<u64> {
250 number_field(primary, names).or_else(|| number_field(fallback, names))
251}
252
253fn tool_result_output_bytes(object: &serde_json::Map<String, Value>) -> Option<u64> {
254 let content = object
255 .get("message")
256 .and_then(Value::as_object)
257 .and_then(|message| message.get("content"))
258 .and_then(Value::as_array)?;
259 let mut total = 0_u64;
260 let mut found = false;
261 for block in content {
262 let Some(block) = block.as_object() else {
263 continue;
264 };
265 let Some(fragments) = block.get("content") else {
266 continue;
267 };
268 match fragments {
269 Value::String(text) => {
270 total = total.saturating_add(text.len() as u64);
271 found = true;
272 }
273 Value::Array(fragments) => {
274 for fragment in fragments {
275 if let Some(text) = fragment.as_str().or_else(|| {
276 fragment
277 .as_object()
278 .and_then(|fragment| fragment.get("text"))
279 .and_then(Value::as_str)
280 }) {
281 total = total.saturating_add(text.len() as u64);
282 found = true;
283 }
284 }
285 }
286 _ => {}
287 }
288 }
289 found.then_some(total)
290}
291
292fn tool_result_error_category(object: &serde_json::Map<String, Value>) -> Option<&str> {
293 let content = object
294 .get("message")
295 .and_then(Value::as_object)
296 .and_then(|message| message.get("content"))
297 .and_then(Value::as_array)?;
298 for block in content {
299 let Some(block) = block.as_object() else {
300 continue;
301 };
302 if !block.get("isError").and_then(Value::as_bool).unwrap_or(false) {
303 continue;
304 }
305 if let Some(category) = block.get("content").and_then(Value::as_array).and_then(|fragments| {
306 fragments.iter().find_map(|fragment| {
307 fragment
308 .as_object()
309 .and_then(|fragment| fragment.get("text"))
310 .and_then(Value::as_str)
311 .or_else(|| fragment.as_str())
312 })
313 }) {
314 return Some(category);
315 }
316 return Some("tool_error");
317 }
318 None
319}
320
321fn record_tool(summary: &mut HarnessTraceSummary, tool: &str) {
322 let normalized_tool = safe_tool_name(tool);
323 let tool = if summary.tool_counts.contains_key(normalized_tool)
324 || summary.tool_counts.len() < MAX_DISTINCT_TOOL_LABELS.saturating_sub(1)
325 {
326 normalized_tool
327 } else {
328 OTHER_TOOL_LABEL
329 };
330 summary.tool_calls = summary.tool_calls.saturating_add(1);
331 let count = summary.tool_counts.entry(tool.to_owned()).or_default();
332 if *count > 0 {
333 summary.repeated_calls = summary.repeated_calls.saturating_add(1);
334 let repeated_count = summary.repeated_tool_counts.entry(tool.to_owned()).or_default();
335 *repeated_count = repeated_count.saturating_add(1);
336 }
337 *count = count.saturating_add(1);
338}
339
340fn add_error(summary: &mut HarnessTraceSummary, category: &str) {
341 let category = canonical_error_category(category);
342 let count = summary.error_categories.entry(category.to_owned()).or_default();
343 *count = count.saturating_add(1);
344}
345
346fn safe_tool_name(tool: &str) -> &'static str {
347 match tool.trim() {
348 "apply_patch" => "apply_patch",
349 "bash" => "bash",
350 "code_search" => "code_search",
351 "command_execution" => "command_execution",
352 "create_goal" => "create_goal",
353 "edit_file" => "edit_file",
354 "edit" => "edit",
355 "exec" => "exec",
356 "exec_command" => "exec_command",
357 "exec_pty_cmd" => "exec_pty_cmd",
358 "fetch" => "fetch",
359 "fetch_url" => "fetch_url",
360 "get_goal" => "get_goal",
361 "grep" => "grep",
362 "grep_file" => "grep_file",
363 "job_output" => "job_output",
364 "list" => "list",
365 "list_agents" => "list_agents",
366 "list_dir" => "list_dir",
367 "list_files" => "list_files",
368 "mcp" => "mcp",
369 "mcp_tool_call" => "mcp_tool_call",
370 "read" => "read",
371 "read_file" => "read_file",
372 "search" => "search",
373 "shell" => "shell",
374 "skill" => "skill",
375 "subagent" => "subagent",
376 "task_tracker" => "task_tracker",
377 "todo_write" => "todo_write",
378 "update_goal" => "update_goal",
379 "web_fetch" => "web_fetch",
380 "web_search" => "web_search",
381 "write" => "write",
382 "write_file" => "write_file",
383 "write_stdin" => "write_stdin",
384 _ => OTHER_TOOL_LABEL,
385 }
386}
387
388fn string_field<'a>(object: &'a serde_json::Map<String, Value>, names: &[&str]) -> Option<&'a str> {
389 names.iter().find_map(|name| object.get(*name).and_then(Value::as_str))
390}
391
392fn number_field(object: &serde_json::Map<String, Value>, names: &[&str]) -> Option<u64> {
393 names.iter().find_map(|name| object.get(*name).and_then(Value::as_u64))
394}
395
396fn nested_number_field(object: &serde_json::Map<String, Value>, containers: &[&str], names: &[&str]) -> Option<u64> {
397 containers.iter().find_map(|container| {
398 object
399 .get(*container)
400 .and_then(Value::as_object)
401 .and_then(|details| number_field(details, names))
402 })
403}
404
405fn output_bytes(object: &serde_json::Map<String, Value>) -> Option<u64> {
406 ["output", "aggregated_output", "tool_output"]
407 .iter()
408 .find_map(|name| object.get(*name).and_then(Value::as_str).map(|output| output.len() as u64))
409}
410
411fn error_category<'a>(
412 primary: &'a serde_json::Map<String, Value>,
413 fallback: &'a serde_json::Map<String, Value>,
414) -> Option<&'a str> {
415 if let Some(category) = string_field_from(primary, fallback, &["error_category", "error_code"]) {
416 return Some(category);
417 }
418 match primary.get("error").or_else(|| fallback.get("error")) {
419 Some(Value::String(error)) => Some(error),
420 Some(Value::Object(error)) => string_field(error, &["category", "code", "type"]),
421 _ => None,
422 }
423}
424
425fn canonical_error_category(error: &str) -> &'static str {
426 let error = error.trim().to_ascii_lowercase();
427 if error.contains("fs_not_observed") || error.contains("fs-not-observed") {
428 "fs_not_observed"
429 } else if error.contains("fs_stale_version") || error.contains("fs-stale-version") {
430 "fs_stale_version"
431 } else if error.contains("unknown_job") || error.contains("unknown-job") || error.contains("unknown job") {
432 "unknown_job"
433 } else if error.contains("invalid_goal") || error.contains("invalid-goal") || error.contains("invalid goal") {
434 "invalid_goal_update"
435 } else if error.contains("timeout") || error.contains("timed out") {
436 "timeout"
437 } else if error.contains("permission") || error.contains("denied") {
438 "permission_denied"
439 } else if error.contains("network") || error.contains("connection") {
440 "network"
441 } else if error.contains("parse") || error.contains("json") {
442 "parse"
443 } else if error.contains("rate_limit") || error.contains("rate-limit") || error.contains("ratelimit") {
444 "rate_limit"
445 } else if error.contains("command_failed") || error.contains("command-failed") {
446 "command_failed"
447 } else if error.contains("tool_error") || error.contains("tool-error") {
448 "tool_error"
449 } else {
450 GENERIC_ERROR_CATEGORY
451 }
452}
453
454fn usage_value(object: &serde_json::Map<String, Value>) -> Option<&serde_json::Map<String, Value>> {
455 ["usage", "tokens"]
456 .iter()
457 .find_map(|name| object.get(*name).and_then(Value::as_object))
458}
459
460fn has_usage(object: &serde_json::Map<String, Value>) -> bool {
461 usage_value(object).is_some()
462 || object.keys().any(|key| {
463 key.ends_with("_tokens")
464 || matches!(key.as_str(), "inputTokens" | "outputTokens" | "cacheReadTokens" | "reasoningTokens")
465 })
466}
467
468fn usage_sample(object: &serde_json::Map<String, Value>) -> Option<TokenUsage> {
469 if !has_usage(object) {
470 return None;
471 }
472
473 let source = usage_value(object).unwrap_or(object);
474 Some(TokenUsage {
475 input_tokens: number_field(source, &["input", "input_tokens", "prompt_tokens", "inputTokens"]).unwrap_or(0),
476 output_tokens: number_field(source, &["output", "output_tokens", "completion_tokens", "outputTokens"])
477 .unwrap_or(0),
478 cached_input_tokens: number_field(
479 source,
480 &[
481 "cached",
482 "cached_tokens",
483 "cached_input_tokens",
484 "cacheReadTokens",
485 "cache_read_tokens",
486 "prompt_cache_hit_tokens",
487 ],
488 )
489 .or_else(|| {
490 nested_number_field(
491 source,
492 &["input_tokens_details", "prompt_tokens_details"],
493 &["cached_tokens", "cache_read_tokens", "cacheReadTokens"],
494 )
495 })
496 .unwrap_or(0),
497 cache_creation_tokens: number_field(
498 source,
499 &[
500 "cache_creation",
501 "cache_creation_tokens",
502 "cacheCreationTokens",
503 "prompt_cache_creation_tokens",
504 "cache_write_tokens",
505 ],
506 )
507 .or_else(|| {
508 nested_number_field(
509 source,
510 &["input_tokens_details", "prompt_tokens_details"],
511 &["cache_creation_tokens", "cache_write_tokens", "cacheWriteTokens"],
512 )
513 })
514 .unwrap_or(0),
515 reasoning_tokens: number_field(source, &["reasoning", "reasoning_tokens", "reasoningTokens"])
516 .or_else(|| {
517 nested_number_field(
518 source,
519 &["output_tokens_details", "completion_tokens_details"],
520 &["reasoning_tokens", "reasoningTokens"],
521 )
522 })
523 .unwrap_or(0),
524 })
525}
526
527fn add_usage(total: &mut TokenUsage, sample: &TokenUsage) {
528 total.input_tokens = total.input_tokens.saturating_add(sample.input_tokens);
529 total.output_tokens = total.output_tokens.saturating_add(sample.output_tokens);
530 total.cached_input_tokens = total.cached_input_tokens.saturating_add(sample.cached_input_tokens);
531 total.cache_creation_tokens = total.cache_creation_tokens.saturating_add(sample.cache_creation_tokens);
532 total.reasoning_tokens = total.reasoning_tokens.saturating_add(sample.reasoning_tokens);
533}
534
535fn deterministic_reservoir_index(sample_number: u64) -> u64 {
536 let mut mixed = sample_number;
537 mixed ^= mixed >> 30;
538 mixed = mixed.wrapping_mul(0xbf58_476d_1ce4_e5b9);
539 mixed ^= mixed >> 27;
540 mixed = mixed.wrapping_mul(0x94d0_49bb_1331_11eb);
541 mixed ^= mixed >> 31;
542 mixed % sample_number.max(1)
543}
544
545#[cfg(test)]
546mod latency_tests {
547 use super::{analyze_jsonl, metrics::MAX_LATENCY_SAMPLES};
548
549 #[test]
550 fn bounds_latency_sample_storage_while_retaining_all_counts() {
551 let input = (0..(MAX_LATENCY_SAMPLES * 2))
552 .map(|latency| format!(r#"{{"latency_ms":{latency}}}"#))
553 .collect::<Vec<_>>()
554 .join("\n");
555
556 let summary = analyze_jsonl(&input).expect("trace should parse");
557
558 assert_eq!(summary.latency.count, (MAX_LATENCY_SAMPLES * 2) as u64);
559 assert_eq!(summary.latency.max_ms, Some((MAX_LATENCY_SAMPLES * 2 - 1) as u64));
560 }
561}
562
563#[cfg(test)]
564mod tests {
565 use std::io::{BufReader, Cursor};
566
567 use super::*;
568
569 const THREAD_EVENT_TRACE: &str = r#"{"type":"turn.started"}
570{"type":"item.completed","item":{"id":"1","type":"tool_invocation","tool_name":"exec_command","status":"completed"}}
571{"type":"item.completed","item":{"id":"2","type":"tool_invocation","tool_name":"exec_command","status":"failed","outcome":"timeout"}}
572{"type":"turn.completed","usage":{"input_tokens":8,"cached_input_tokens":3,"cache_creation_tokens":1,"output_tokens":5}}
573{"type":"item.completed","item":{"id":"3","type":"tool_output","output":"private output"}}
574"#;
575
576 #[test]
577 fn summarizes_deepseek_baseline_without_retaining_raw_text() {
578 let trace = r#"
579{"step":1,"tool":"exec_command","latency_ms":12,"output":"secret command output","tokens":{"input":100,"output":20,"cached":40}}
580{"step":2,"tool":"read_file","latency_ms":20,"error":"timeout"}
581"#;
582
583 let summary = analyze_jsonl(trace).expect("trace should parse");
584
585 assert_eq!(summary.steps, 2);
586 assert_eq!(summary.tool_calls, 2);
587 assert_eq!(summary.tool_counts["exec_command"], 1);
588 assert_eq!(summary.tool_counts["read_file"], 1);
589 assert_eq!(summary.error_categories["timeout"], 1);
590 assert_eq!(summary.output_bytes, 21);
591 assert_eq!(summary.token_usage.input_tokens, 100);
592 assert_eq!(summary.token_usage.output_tokens, 20);
593 assert_eq!(summary.token_usage.cached_input_tokens, 40);
594 assert!(
595 !serde_json::to_string(&summary)
596 .expect("summary should serialize")
597 .contains("secret command output")
598 );
599 }
600
601 #[test]
602 fn matches_known_deepseek_baseline_counts_with_compact_fixture() {
603 let mut trace = String::new();
604 for step in 1..=453 {
605 trace.push_str(&format!(
606 r#"{{"step":{step},"tool":"exec_command"}}
607"#
608 ));
609 }
610 trace.push_str(
611 r#"{"tool":"exec_command"}
612{"tool":"read_file"}
613{"tool":"read_file"}
614{"tool":"write_file"}
615{"tool":"write_file"}
616{"tool":"search"}
617{"tool":"search"}
618{"tool":"search"}
619{"tool":"search"}
620{"tool":"search"}
621{"tool":"search"}
622{"tool":"search"}
623{"tool":"search"}
624{"tool":"search"}
625{"tool":"search"}
626{"error":"timeout"}
627{"error":"timeout"}
628{"error":"timeout"}
629{"error":"timeout"}
630{"error":"timeout"}
631{"error":"timeout"}
632{"error":"timeout"}
633{"error":"timeout"}
634{"error":"timeout"}
635{"error":"timeout"}
636{"error":"timeout"}
637{"error":"timeout"}
638{"error":"timeout"}
639{"error":"timeout"}
640{"error":"timeout"}
641{"error":"timeout"}
642{"error":"timeout"}
643{"error":"timeout"}
644{"error":"timeout"}
645{"error":"timeout"}
646"#,
647 );
648
649 let summary = analyze_jsonl(&trace).expect("baseline fixture should parse");
650
651 assert_eq!(summary.steps, 453);
652 assert_eq!(summary.tool_calls, 468);
653 assert_eq!(summary.error_categories.values().sum::<u64>(), 20);
654 }
655
656 #[test]
657 fn parses_thread_events_and_skips_bad_or_unknown_lines() {
658 let input = format!("{}\nnot json\n{{\"future\":true}}\n", THREAD_EVENT_TRACE);
659 let summary = analyze_jsonl(&input).expect("trace should parse");
660
661 assert_eq!(summary.turns, 1);
662 assert_eq!(summary.steps, 2);
663 assert_eq!(summary.tool_calls, 2);
664 assert_eq!(summary.repeated_calls, 1);
665 assert_eq!(summary.error_categories["timeout"], 1);
666 assert_eq!(summary.output_bytes, 14);
667 assert_eq!(summary.token_usage.input_tokens, 8);
668 assert_eq!(summary.token_usage.cached_input_tokens, 3);
669 assert_eq!(summary.malformed_lines, 1);
670 assert_eq!(summary.unrecognized_lines, 1);
671 }
672
673 #[test]
674 fn reports_latency_statistics_and_file_errors_with_context() {
675 let summary = analyze_jsonl(
676 "{\"step\":1,\"latency_ms\":30}\n{\"step\":2,\"latency_ms\":10}\n{\"step\":3,\"latency_ms\":20}\n",
677 )
678 .expect("latency trace should parse");
679 assert_eq!(summary.latency.count, 3);
680 assert_eq!(summary.latency.total_ms, 60);
681 assert_eq!(summary.latency.p50_ms, Some(20));
682 assert_eq!(summary.latency.p95_ms, Some(30));
683
684 let missing = analyze_jsonl_file("/path/that/does/not/exist.jsonl").expect_err("missing file should fail");
685 assert!(missing.to_string().contains("read trace file"));
686 }
687
688 #[test]
689 fn redacts_untrusted_labels_and_counts_event_errors_once() {
690 let input = r#"{"type":"error","error":"secret command output timeout"}
691{"type":"turn.failed","error_category":"FS_STALE_VERSION","error":"private details"}
692{"step":1,"tool":"rm /sensitive/project","error":{"type":"secret_error"}}
693"#;
694
695 let summary = analyze_jsonl(input).expect("trace should parse");
696 let serialized = serde_json::to_string(&summary).expect("summary should serialize");
697
698 assert_eq!(summary.error_categories.values().sum::<u64>(), 3);
699 assert_eq!(summary.error_categories["timeout"], 1);
700 assert_eq!(summary.error_categories["fs_stale_version"], 1);
701 assert_eq!(summary.error_categories["error"], 1);
702 assert_eq!(summary.tool_counts["other_tool"], 1);
703 assert!(!serialized.contains("secret command output"));
704 assert!(!serialized.contains("/sensitive/project"));
705 assert!(!serialized.contains("secret_error"));
706 }
707
708 #[test]
709 fn classifies_unrelated_goal_text_by_its_actual_error() {
710 let summary = analyze_jsonl(
711 r#"{"error":"goal timeout"}
712{"error":"invalid goal update"}
713"#,
714 )
715 .expect("trace should parse");
716
717 assert_eq!(summary.error_categories["timeout"], 1);
718 assert_eq!(summary.error_categories["invalid_goal_update"], 1);
719 }
720
721 #[test]
722 fn saturates_latency_total_on_overflow() {
723 let input = format!("{{\"latency_ms\":{0}}}\n{{\"latency_ms\":{0}}}\n", u64::MAX);
724
725 let summary = analyze_jsonl(&input).expect("trace should parse");
726
727 assert_eq!(summary.latency.total_ms, u64::MAX);
728 assert_eq!(summary.latency.max_ms, Some(u64::MAX));
729 }
730
731 #[test]
732 fn prefers_per_turn_usage_over_thread_aggregate_and_reads_thread_turn_count() {
733 let input = r#"{"type":"turn.completed","usage":{"input_tokens":3,"cached_input_tokens":1,"output_tokens":2}}
734{"type":"thread.completed","num_turns":7,"usage":{"input_tokens":10,"cached_input_tokens":4,"output_tokens":8},"result":"private assistant result"}
735"#;
736
737 let summary = analyze_jsonl(input).expect("trace should parse");
738
739 assert_eq!(summary.turns, 7);
740 assert_eq!(summary.token_usage.input_tokens, 3);
741 assert_eq!(summary.token_usage.cached_input_tokens, 1);
742 assert_eq!(summary.token_usage.output_tokens, 2);
743 assert_eq!(summary.output_bytes, 0);
744 }
745
746 #[test]
747 fn falls_back_to_thread_aggregate_usage_when_turn_usage_is_missing() {
748 let input = r#"{"type":"thread.completed","num_turns":4,"usage":{"input_tokens":10,"cached_input_tokens":4,"cache_creation_tokens":2,"output_tokens":8}}
749"#;
750
751 let summary = analyze_jsonl(input).expect("trace should parse");
752
753 assert_eq!(summary.turns, 4);
754 assert_eq!(summary.token_usage.input_tokens, 10);
755 assert_eq!(summary.token_usage.cached_input_tokens, 4);
756 assert_eq!(summary.token_usage.cache_creation_tokens, 2);
757 assert_eq!(summary.token_usage.output_tokens, 8);
758 }
759
760 #[test]
761 fn buffered_reader_api_matches_text_analysis() {
762 let input = "{\"step\":1,\"tool\":\"read_file\",\"latency_ms\":12}\n";
763
764 let from_text = analyze_jsonl(input).expect("text trace should parse");
765 let from_reader =
766 analyze_jsonl_reader(BufReader::new(Cursor::new(input.as_bytes()))).expect("buffered trace should parse");
767
768 assert_eq!(from_reader, from_text);
769 }
770
771 #[test]
772 fn parses_nested_deepseek_envelopes_and_camel_case_usage() {
773 let input = r#"{"type":"turn/start","time":100,"data":{"turn":1}}
774{"type":"step/start","time":110,"data":{"step":1,"turn":1}}
775{"type":"assistant/message","data":{"step":1,"usage":{"inputTokens":100,"outputTokens":20,"cacheReadTokens":40,"reasoningTokens":7}}}
776{"type":"tool/call","time":120,"data":{"name":"read_file","arguments":"private arguments","step":1,"turn":1}}
777{"type":"tool/result","time":150,"data":{"step":1,"turn":1,"message":{"content":[{"type":"text","content":["private output"],"isError":true}]}}}
778{"type":"step/end","time":180,"data":{"step":1,"turn":1}}
779{"type":"turn/end","time":200,"data":{"turn":1,"reason":{}}}
780"#;
781
782 let summary = analyze_jsonl(input).expect("nested trace should parse");
783 let serialized = serde_json::to_string(&summary).expect("summary should serialize");
784
785 assert_eq!(summary.turns, 1);
786 assert_eq!(summary.steps, 1);
787 assert_eq!(summary.tool_calls, 1);
788 assert_eq!(summary.tool_counts["read_file"], 1);
789 assert_eq!(summary.error_categories["tool_error"], 1);
790 assert_eq!(summary.output_bytes, 14);
791 assert_eq!(summary.latency.total_ms, 70);
792 assert_eq!(summary.token_usage.input_tokens, 100);
793 assert_eq!(summary.token_usage.output_tokens, 20);
794 assert_eq!(summary.token_usage.cached_input_tokens, 40);
795 assert_eq!(summary.token_usage.reasoning_tokens, 7);
796 assert!(!serialized.contains("private arguments"));
797 assert!(!serialized.contains("private output"));
798 }
799
800 #[test]
801 fn parses_versioned_vtcode_event_envelopes() {
802 let input = r#"{"schema_version":"0.11.0","event":{"type":"turn.started"}}
803{"schema_version":"0.11.0","event":{"type":"item.completed","item":{"id":"1","type":"tool_invocation","tool_name":"read_file","status":"completed"}}}
804{"schema_version":"0.11.0","event":{"type":"turn.completed","usage":{"input_tokens":8,"output_tokens":2}}}
805"#;
806
807 let summary = analyze_jsonl(input).expect("versioned trace should parse");
808
809 assert_eq!(summary.turns, 1);
810 assert_eq!(summary.steps, 1);
811 assert_eq!(summary.tool_counts["read_file"], 1);
812 assert_eq!(summary.token_usage.input_tokens, 8);
813 assert_eq!(summary.token_usage.output_tokens, 2);
814 assert_eq!(summary.unrecognized_lines, 0);
815 }
816
817 #[test]
818 fn terminal_usage_takes_precedence_over_intermediate_and_thread_usage() {
819 let input = r#"{"type":"assistant/message","data":{"usage":{"inputTokens":100,"outputTokens":20}}}
820{"type":"turn/end","data":{"usage":{"inputTokens":3,"outputTokens":2}}}
821{"type":"thread.completed","num_turns":1,"usage":{"input_tokens":10,"output_tokens":8}}
822"#;
823
824 let summary = analyze_jsonl(input).expect("usage trace should parse");
825
826 assert_eq!(summary.token_usage.input_tokens, 3);
827 assert_eq!(summary.token_usage.output_tokens, 2);
828 }
829
830 #[test]
831 fn parses_nested_provider_token_details() {
832 let input = r#"{"usage":{"prompt_tokens":100,"completion_tokens":20,"prompt_cache_hit_tokens":40,"input_tokens_details":{"cached_tokens":5},"output_tokens_details":{"reasoning_tokens":7}}}
833"#;
834
835 let summary = analyze_jsonl(input).expect("provider usage should parse");
836
837 assert_eq!(summary.token_usage.input_tokens, 100);
838 assert_eq!(summary.token_usage.output_tokens, 20);
839 assert_eq!(summary.token_usage.cached_input_tokens, 40);
840 assert_eq!(summary.token_usage.reasoning_tokens, 7);
841 }
842
843 #[test]
844 fn rejects_oversized_records_without_blocking_following_lines() {
845 let input = format!(
846 "{{\"payload\":\"{}\"}}\n{{\"step\":1,\"tool\":\"read_file\"}}\n",
847 "x".repeat(MAX_TRACE_LINE_BYTES)
848 );
849
850 let summary = analyze_jsonl(&input).expect("oversized trace should parse");
851
852 assert_eq!(summary.malformed_lines, 1);
853 assert_eq!(summary.steps, 1);
854 assert_eq!(summary.tool_calls, 1);
855 }
856}