1use crate::capability::CapabilityLedgerEntry;
2use crate::event::{AgentEvent, ApprovalDecision};
3use crate::security::redact_secrets;
4use crate::tool::{ToolInvocation, ToolResult};
5use serde::{Deserialize, Serialize};
6use std::collections::HashMap;
7use std::path::{Path, PathBuf};
8use std::time::{SystemTime, UNIX_EPOCH};
9
10#[derive(Debug, Clone, Serialize, Deserialize)]
15pub struct TurnTrace {
16 pub version: u32,
18
19 pub turn_id: String,
21
22 pub session_id: String,
24
25 pub started_at: u64,
27
28 pub ended_at: u64,
30
31 pub model_provider: String,
33
34 pub model_name: String,
36
37 pub task: String,
39
40 pub visible_tool_count: usize,
42
43 #[serde(default)]
45 pub visible_tools: Vec<String>,
46
47 #[serde(default)]
49 pub deferred_tools_discovered: Vec<String>,
50
51 #[serde(default)]
53 pub tool_calls: Vec<ToolCallTrace>,
54
55 #[serde(default)]
57 pub approvals: Vec<ApprovalTrace>,
58
59 #[serde(default)]
61 pub capabilities: Vec<CapabilityLedgerEntry>,
62
63 #[serde(default)]
65 pub verifier_results: Vec<VerifierTrace>,
66
67 pub metrics: TurnMetrics,
69
70 pub outcome: TurnOutcome,
72
73 #[serde(default)]
75 pub final_message: String,
76}
77
78#[derive(Debug, Clone, Serialize, Deserialize)]
80pub struct ToolCallTrace {
81 pub invocation: ToolInvocation,
83 pub result: ToolResult,
85 pub duration_ms: u64,
87 #[serde(default)]
89 pub was_retry: bool,
90 #[serde(default)]
92 pub was_rolled_back: bool,
93 #[serde(default)]
95 pub error_code: Option<String>,
96}
97
98#[derive(Debug, Clone, Serialize, Deserialize)]
100pub struct ApprovalTrace {
101 pub tool_name: String,
103 pub risk: String,
105 pub decision: String,
107 pub duration_ms: u64,
109}
110
111#[derive(Debug, Clone, Serialize, Deserialize)]
113pub struct VerifierTrace {
114 pub verifier: String,
116 pub command: String,
118 pub passed: bool,
120 pub duration_ms: u64,
122 pub exit_code: Option<i32>,
124}
125
126#[derive(Debug, Clone, Serialize, Deserialize)]
128pub struct TurnMetrics {
129 pub input_tokens: u64,
131 pub output_tokens: u64,
133 pub cache_creation_tokens: u64,
135 pub cache_read_tokens: u64,
137 pub tool_call_count: usize,
139 pub failed_tool_calls: usize,
141 pub approval_count: usize,
143 pub verifier_count: usize,
145 pub retry_count: usize,
147 pub rollback_count: usize,
149 pub wall_time_ms: u64,
151}
152
153impl Default for TurnMetrics {
154 fn default() -> Self {
155 Self {
156 input_tokens: 0,
157 output_tokens: 0,
158 cache_creation_tokens: 0,
159 cache_read_tokens: 0,
160 tool_call_count: 0,
161 failed_tool_calls: 0,
162 approval_count: 0,
163 verifier_count: 0,
164 retry_count: 0,
165 rollback_count: 0,
166 wall_time_ms: 0,
167 }
168 }
169}
170
171#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
173pub enum TurnOutcome {
174 Success,
176 PartialSuccess,
178 Stopped(String),
180 Failed(String),
182}
183
184impl TurnTrace {
185 pub const CURRENT_VERSION: u32 = 1;
187
188 pub fn new(
190 turn_id: impl Into<String>,
191 session_id: impl Into<String>,
192 model_provider: impl Into<String>,
193 model_name: impl Into<String>,
194 task: impl Into<String>,
195 ) -> Self {
196 Self {
197 version: Self::CURRENT_VERSION,
198 turn_id: turn_id.into(),
199 session_id: session_id.into(),
200 started_at: current_unix_millis(),
201 ended_at: 0,
202 model_provider: model_provider.into(),
203 model_name: model_name.into(),
204 task: task.into(),
205 visible_tool_count: 0,
206 visible_tools: Vec::new(),
207 deferred_tools_discovered: Vec::new(),
208 tool_calls: Vec::new(),
209 approvals: Vec::new(),
210 capabilities: Vec::new(),
211 verifier_results: Vec::new(),
212 metrics: TurnMetrics::default(),
213 outcome: TurnOutcome::Success,
214 final_message: String::new(),
215 }
216 }
217
218 pub fn finalize(&mut self) {
220 self.ended_at = current_unix_millis();
221 self.metrics.wall_time_ms = self.ended_at.saturating_sub(self.started_at);
222 }
223
224 pub fn record_tool_call(
226 &mut self,
227 invocation: &ToolInvocation,
228 result: &ToolResult,
229 duration_ms: u64,
230 ) {
231 self.metrics.tool_call_count += 1;
232 if !result.ok {
233 self.metrics.failed_tool_calls += 1;
234 }
235 self.tool_calls.push(ToolCallTrace {
236 invocation: invocation.clone(),
237 result: result.clone(),
238 duration_ms,
239 was_retry: false,
240 was_rolled_back: false,
241 error_code: result
242 .output
243 .get("error_code")
244 .and_then(|v| v.as_str())
245 .map(|s| s.to_string()),
246 });
247 }
248
249 pub fn record_approval(
251 &mut self,
252 tool_name: &str,
253 risk: &str,
254 approved: bool,
255 duration_ms: u64,
256 ) {
257 self.metrics.approval_count += 1;
258 self.approvals.push(ApprovalTrace {
259 tool_name: tool_name.to_string(),
260 risk: risk.to_string(),
261 decision: if approved {
262 "approved".to_string()
263 } else {
264 "denied".to_string()
265 },
266 duration_ms,
267 });
268 }
269
270 pub fn record_capability(&mut self, entry: CapabilityLedgerEntry) {
272 self.capabilities.push(entry);
273 }
274
275 pub fn record_verifier(
277 &mut self,
278 verifier: &str,
279 command: &str,
280 passed: bool,
281 duration_ms: u64,
282 exit_code: Option<i32>,
283 ) {
284 self.metrics.verifier_count += 1;
285 self.verifier_results.push(VerifierTrace {
286 verifier: verifier.to_string(),
287 command: redact_secrets(command),
288 passed,
289 duration_ms,
290 exit_code,
291 });
292 }
293
294 pub fn redacted(&self) -> Self {
296 let mut trace = self.clone();
297 trace.task = redact_secrets(&trace.task);
298 trace.final_message = redact_secrets(&trace.final_message);
299 for call in &mut trace.tool_calls {
300 call.invocation.input = redact_json_value(&call.invocation.input);
301 call.result.output = redact_json_value(&call.result.output);
302 }
303 for verifier in &mut trace.verifier_results {
304 verifier.command = redact_secrets(&verifier.command);
305 }
306 for capability in &mut trace.capabilities {
307 capability.justification = redact_secrets(&capability.justification);
308 }
309 trace
310 }
311}
312
313#[derive(Debug, Clone)]
315pub struct TraceStore {
316 root: PathBuf,
317}
318
319impl TraceStore {
320 pub fn new(data_dir: &Path) -> Self {
322 Self {
323 root: data_dir.join("traces"),
324 }
325 }
326
327 pub fn root(&self) -> &Path {
329 &self.root
330 }
331
332 pub fn save_trace(&self, trace: &TurnTrace) -> std::io::Result<PathBuf> {
336 std::fs::create_dir_all(&self.root)?;
337 let path = self.root.join(format!("{}.jsonl", trace.session_id));
338 let mut file = std::fs::OpenOptions::new()
339 .create(true)
340 .append(true)
341 .open(&path)?;
342 let redacted = trace.redacted();
343 let line = serde_json::to_string(&redacted)?;
344 use std::io::Write;
345 writeln!(file, "{line}")?;
346 Ok(path)
347 }
348
349 pub fn save_session_traces(
351 &self,
352 session_id: &str,
353 traces: &[TurnTrace],
354 ) -> std::io::Result<PathBuf> {
355 std::fs::create_dir_all(&self.root)?;
356 let path = self.root.join(format!("{session_id}.jsonl"));
357 let mut file = std::fs::File::create(&path)?;
358 use std::io::Write;
359 for trace in traces {
360 let redacted = trace.redacted();
361 let line = serde_json::to_string(&redacted)?;
362 writeln!(file, "{line}")?;
363 }
364 Ok(path)
365 }
366
367 pub fn load_session_traces(&self, session_id: &str) -> Vec<TurnTrace> {
369 let path = self.root.join(format!("{session_id}.jsonl"));
370 let Ok(content) = std::fs::read_to_string(&path) else {
371 return Vec::new();
372 };
373 content
374 .lines()
375 .filter_map(|line| serde_json::from_str(line).ok())
376 .collect()
377 }
378
379 pub fn list_sessions(&self) -> Vec<String> {
381 let Ok(entries) = std::fs::read_dir(&self.root) else {
382 return Vec::new();
383 };
384 entries
385 .filter_map(|e| e.ok())
386 .filter(|e| e.path().extension().is_some_and(|ext| ext == "jsonl"))
387 .filter_map(|e| {
388 e.path()
389 .file_stem()
390 .and_then(|s| s.to_str())
391 .map(|s| s.to_string())
392 })
393 .collect()
394 }
395}
396
397pub fn turn_traces_from_events(
398 session_id: &str,
399 model_provider: &str,
400 model_name: &str,
401 events: &[AgentEvent],
402) -> Vec<TurnTrace> {
403 let mut traces = Vec::new();
404 let mut current: Option<TurnTrace> = None;
405 let mut invocations: HashMap<String, ToolInvocation> = HashMap::new();
406
407 for event in events {
408 match event {
409 AgentEvent::UserTaskSubmitted { text, .. } => {
410 if let Some(mut trace) = current.take() {
411 trace.finalize();
412 traces.push(trace);
413 }
414 invocations.clear();
415 current = Some(TurnTrace::new(
416 format!("{}-trace-{}", session_id, traces.len() + 1),
417 session_id.to_string(),
418 model_provider.to_string(),
419 model_name.to_string(),
420 text.clone(),
421 ));
422 }
423 AgentEvent::ToolRequested(invocation) => {
424 invocations.insert(invocation.id.clone(), invocation.clone());
425 }
426 AgentEvent::ToolCompleted(result) => {
427 if let Some(trace) = current.as_mut() {
428 let invocation =
429 invocations
430 .get(&result.invocation_id)
431 .cloned()
432 .unwrap_or(ToolInvocation {
433 id: result.invocation_id.clone(),
434 tool_name: result
435 .output
436 .get("tool")
437 .and_then(serde_json::Value::as_str)
438 .unwrap_or("unknown")
439 .to_string(),
440 input: serde_json::json!({}),
441 });
442 trace.record_tool_call(&invocation, result, 0);
443 if invocation.tool_name == "verifier" {
444 if let Some(verifier_trace) = verifier_trace_from_tool(&invocation, result)
445 {
446 trace.metrics.verifier_count += 1;
447 trace.verifier_results.push(verifier_trace);
448 }
449 }
450 }
451 }
452 AgentEvent::ApprovalRequested(request) => {
453 if let Some(trace) = current.as_mut() {
454 trace.metrics.approval_count += 1;
455 trace.approvals.push(ApprovalTrace {
456 tool_name: request.id.clone(),
457 risk: format!("{:?}", request.risk).to_lowercase(),
458 decision: "requested".to_string(),
459 duration_ms: 0,
460 });
461 }
462 }
463 AgentEvent::ApprovalResolved(decision) => {
464 if let Some(trace) = current.as_mut() {
465 let (id, label) = match decision {
466 ApprovalDecision::Approved { id } => (id, "approved"),
467 ApprovalDecision::Denied { id } => (id, "denied"),
468 };
469 trace.approvals.push(ApprovalTrace {
470 tool_name: id.clone(),
471 risk: String::new(),
472 decision: label.to_string(),
473 duration_ms: 0,
474 });
475 }
476 }
477 AgentEvent::CapabilityRecorded(entry) => {
478 if let Some(trace) = current.as_mut() {
479 trace.record_capability(entry.clone());
480 }
481 }
482 AgentEvent::UsageReported {
483 input_tokens,
484 output_tokens,
485 cache_creation_tokens,
486 cache_read_tokens,
487 } => {
488 if let Some(trace) = current.as_mut() {
489 trace.metrics.input_tokens =
490 trace.metrics.input_tokens.saturating_add(*input_tokens);
491 trace.metrics.output_tokens =
492 trace.metrics.output_tokens.saturating_add(*output_tokens);
493 trace.metrics.cache_creation_tokens = trace
494 .metrics
495 .cache_creation_tokens
496 .saturating_add(*cache_creation_tokens);
497 trace.metrics.cache_read_tokens = trace
498 .metrics
499 .cache_read_tokens
500 .saturating_add(*cache_read_tokens);
501 }
502 }
503 AgentEvent::ModelOutput { text, .. } => {
504 if let Some(trace) = current.as_mut() {
505 trace.final_message = text.clone();
506 }
507 }
508 AgentEvent::Error { message } => {
509 if let Some(trace) = current.as_mut() {
510 trace.outcome = TurnOutcome::Failed(message.clone());
511 }
512 }
513 AgentEvent::HarnessStopped { reason, .. } => {
514 if let Some(trace) = current.as_mut() {
515 trace.outcome = TurnOutcome::Stopped(reason.clone());
516 }
517 }
518 _ => {}
519 }
520 }
521
522 if let Some(mut trace) = current {
523 trace.finalize();
524 if trace.metrics.failed_tool_calls > 0 && trace.outcome == TurnOutcome::Success {
525 trace.outcome = TurnOutcome::PartialSuccess;
526 }
527 traces.push(trace);
528 }
529
530 traces
531}
532
533fn verifier_trace_from_tool(
534 invocation: &ToolInvocation,
535 result: &ToolResult,
536) -> Option<VerifierTrace> {
537 let command = result
538 .output
539 .get("command")
540 .and_then(serde_json::Value::as_str)
541 .or_else(|| {
542 invocation
543 .input
544 .get("command")
545 .and_then(serde_json::Value::as_str)
546 })?;
547 let status = result
548 .output
549 .get("status")
550 .and_then(serde_json::Value::as_str)
551 .unwrap_or(if result.ok { "pass" } else { "fail" });
552 Some(VerifierTrace {
553 verifier: invocation
554 .input
555 .get("verifier")
556 .and_then(serde_json::Value::as_str)
557 .unwrap_or("command")
558 .to_string(),
559 command: redact_secrets(command),
560 passed: matches!(status, "pass" | "skipped"),
561 duration_ms: result
562 .output
563 .get("duration_ms")
564 .and_then(serde_json::Value::as_u64)
565 .unwrap_or_default(),
566 exit_code: result
567 .output
568 .get("exit_code")
569 .and_then(serde_json::Value::as_i64)
570 .map(|value| value as i32),
571 })
572}
573
574fn current_unix_millis() -> u64 {
575 SystemTime::now()
576 .duration_since(UNIX_EPOCH)
577 .map(|d| d.as_millis() as u64)
578 .unwrap_or(0)
579}
580
581fn redact_json_value(value: &serde_json::Value) -> serde_json::Value {
582 match value {
583 serde_json::Value::String(text) => serde_json::Value::String(redact_secrets(text)),
584 serde_json::Value::Array(values) => {
585 serde_json::Value::Array(values.iter().map(redact_json_value).collect())
586 }
587 serde_json::Value::Object(map) => serde_json::Value::Object(
588 map.iter()
589 .map(|(key, value)| (key.clone(), redact_json_value(value)))
590 .collect(),
591 ),
592 other => other.clone(),
593 }
594}
595
596#[cfg(test)]
597mod tests {
598 use super::*;
599 use serde_json::json;
600
601 #[test]
602 fn turn_trace_creation_and_finalize() {
603 let mut trace = TurnTrace::new("turn-1", "session-1", "openai", "gpt-4", "test task");
604 assert_eq!(trace.version, TurnTrace::CURRENT_VERSION);
605 assert_eq!(trace.turn_id, "turn-1");
606 assert_eq!(trace.ended_at, 0);
607
608 trace.finalize();
609 assert!(trace.ended_at >= trace.started_at);
610 assert!(trace.metrics.wall_time_ms > 0 || trace.ended_at == trace.started_at);
611 }
612
613 #[test]
614 fn trace_records_tool_call() {
615 let mut trace = TurnTrace::new("t1", "s1", "p", "m", "task");
616 let inv = ToolInvocation {
617 id: "c1".to_string(),
618 tool_name: "read".to_string(),
619 input: json!({"path": "file.txt"}),
620 };
621 let result = ToolResult {
622 invocation_id: "c1".to_string(),
623 ok: true,
624 output: json!("content"),
625 };
626 trace.record_tool_call(&inv, &result, 100);
627 assert_eq!(trace.metrics.tool_call_count, 1);
628 assert_eq!(trace.metrics.failed_tool_calls, 0);
629 }
630
631 #[test]
632 fn trace_records_failed_tool() {
633 let mut trace = TurnTrace::new("t1", "s1", "p", "m", "task");
634 let inv = ToolInvocation {
635 id: "c1".to_string(),
636 tool_name: "bash".to_string(),
637 input: json!({"command": "false"}),
638 };
639 let result = ToolResult {
640 invocation_id: "c1".to_string(),
641 ok: false,
642 output: json!({"error_code": "command_failed", "message": "exit 1"}),
643 };
644 trace.record_tool_call(&inv, &result, 50);
645 assert_eq!(trace.metrics.failed_tool_calls, 1);
646 assert_eq!(
647 trace.tool_calls[0].error_code.as_deref(),
648 Some("command_failed")
649 );
650 }
651
652 #[test]
653 fn trace_records_approval_and_verifier() {
654 let mut trace = TurnTrace::new("t1", "s1", "p", "m", "task");
655 trace.record_approval("write", "write", true, 500);
656 trace.record_verifier("verify.test", "cargo test", true, 3000, Some(0));
657 assert_eq!(trace.metrics.approval_count, 1);
658 assert_eq!(trace.metrics.verifier_count, 1);
659 }
660
661 #[test]
662 fn trace_records_and_redacts_capability_entries() {
663 let mut trace = TurnTrace::new("t1", "s1", "p", "m", "task");
664 trace.record_capability(CapabilityLedgerEntry {
665 capability: crate::capability::Capability::RepoRead,
666 scope: crate::capability::CapabilityScope::Turn("t1".to_string()),
667 decision: crate::capability::CapabilityDecision::Granted,
668 at_ms: 1,
669 justification: "read token=sk-proj-1234567890abcdef".to_string(),
670 });
671
672 let redacted = trace.redacted();
673 assert_eq!(redacted.capabilities.len(), 1);
674 assert!(!redacted.capabilities[0].justification.contains("sk-proj"));
675 assert!(
676 redacted.capabilities[0]
677 .justification
678 .contains("<redacted>")
679 );
680 }
681
682 #[test]
683 fn trace_store_save_and_load() {
684 let dir = tempfile::tempdir().unwrap();
685 let store = TraceStore::new(dir.path());
686 let mut trace = TurnTrace::new("turn-1", "session-test", "openai", "gpt-4", "task");
687 trace.finalize();
688 store.save_trace(&trace).unwrap();
689 let loaded = store.load_session_traces("session-test");
690 assert_eq!(loaded.len(), 1);
691 assert_eq!(loaded[0].turn_id, "turn-1");
692 }
693
694 #[test]
695 fn trace_store_redacts_secrets_before_persisting() {
696 let dir = tempfile::tempdir().unwrap();
697 let store = TraceStore::new(dir.path());
698 let mut trace = TurnTrace::new(
699 "turn-secret",
700 "session-secret",
701 "openai",
702 "gpt-4",
703 "use OPENAI_API_KEY=sk-proj-1234567890abcdef",
704 );
705 trace.record_tool_call(
706 &ToolInvocation {
707 id: "c1".to_string(),
708 tool_name: "bash".to_string(),
709 input: json!({"command": "echo sk-proj-1234567890abcdef"}),
710 },
711 &ToolResult {
712 invocation_id: "c1".to_string(),
713 ok: true,
714 output: json!({"stdout": "sk-proj-1234567890abcdef"}),
715 },
716 1,
717 );
718 trace.finalize();
719
720 let path = store.save_trace(&trace).unwrap();
721 let persisted = std::fs::read_to_string(path).unwrap();
722 assert!(!persisted.contains("sk-proj-1234567890abcdef"));
723 assert!(persisted.contains("<redacted>"));
724 }
725
726 #[test]
727 fn trace_store_save_appends() {
728 let dir = tempfile::tempdir().unwrap();
729 let store = TraceStore::new(dir.path());
730 let mut t1 = TurnTrace::new("turn-1", "session-multi", "p", "m", "task1");
731 t1.finalize();
732 let mut t2 = TurnTrace::new("turn-2", "session-multi", "p", "m", "task2");
733 t2.finalize();
734 store.save_trace(&t1).unwrap();
735 store.save_trace(&t2).unwrap();
736 let loaded = store.load_session_traces("session-multi");
737 assert_eq!(loaded.len(), 2);
738 }
739
740 #[test]
741 fn trace_store_save_session_traces_replaces_existing_file() {
742 let dir = tempfile::tempdir().unwrap();
743 let store = TraceStore::new(dir.path());
744 let mut first = TurnTrace::new("t1", "session-replace", "p", "m", "task1");
745 first.finalize();
746 store.save_trace(&first).unwrap();
747
748 let mut second = TurnTrace::new("t2", "session-replace", "p", "m", "task2");
749 second.finalize();
750 store
751 .save_session_traces("session-replace", &[second])
752 .unwrap();
753
754 let loaded = store.load_session_traces("session-replace");
755 assert_eq!(loaded.len(), 1);
756 assert_eq!(loaded[0].turn_id, "t2");
757 }
758
759 #[test]
760 fn trace_store_list_sessions() {
761 let dir = tempfile::tempdir().unwrap();
762 let store = TraceStore::new(dir.path());
763 let mut t1 = TurnTrace::new("t1", "session-a", "p", "m", "task");
764 t1.finalize();
765 let mut t2 = TurnTrace::new("t2", "session-b", "p", "m", "task");
766 t2.finalize();
767 store.save_trace(&t1).unwrap();
768 store.save_trace(&t2).unwrap();
769 let sessions = store.list_sessions();
770 assert!(sessions.contains(&"session-a".to_string()));
771 assert!(sessions.contains(&"session-b".to_string()));
772 }
773
774 #[test]
775 fn trace_store_empty_for_missing_session() {
776 let dir = tempfile::tempdir().unwrap();
777 let store = TraceStore::new(dir.path());
778 let loaded = store.load_session_traces("nonexistent");
779 assert!(loaded.is_empty());
780 }
781
782 #[test]
783 fn session_events_generate_turn_trace_with_capabilities() {
784 let events = vec![
785 AgentEvent::UserTaskSubmitted {
786 text: "verify".to_string(),
787 content_parts: Vec::new(),
788 submitted_at: None,
789 },
790 AgentEvent::CapabilityRecorded(CapabilityLedgerEntry {
791 capability: crate::capability::Capability::RepoRead,
792 scope: crate::capability::CapabilityScope::SingleCall("read".to_string()),
793 decision: crate::capability::CapabilityDecision::Requested,
794 at_ms: 1,
795 justification: "read".to_string(),
796 }),
797 AgentEvent::ToolRequested(ToolInvocation {
798 id: "verify".to_string(),
799 tool_name: "verifier".to_string(),
800 input: json!({"verifier": "test", "command": "true"}),
801 }),
802 AgentEvent::ToolCompleted(ToolResult {
803 invocation_id: "verify".to_string(),
804 ok: true,
805 output: json!({"status": "pass", "command": "true", "duration_ms": 3, "exit_code": 0}),
806 }),
807 AgentEvent::UsageReported {
808 input_tokens: 100,
809 output_tokens: 20,
810 cache_creation_tokens: 0,
811 cache_read_tokens: 0,
812 },
813 AgentEvent::ModelOutput {
814 text: "done".to_string(),
815 thinking: None,
816 },
817 ];
818
819 let traces = turn_traces_from_events("s", "p", "m", &events);
820
821 assert_eq!(traces.len(), 1);
822 assert_eq!(traces[0].task, "verify");
823 assert_eq!(traces[0].tool_calls.len(), 1);
824 assert_eq!(traces[0].verifier_results.len(), 1);
825 assert_eq!(traces[0].capabilities.len(), 1);
826 assert_eq!(traces[0].metrics.input_tokens, 100);
827 }
828}