Skip to main content

vtcode_exec_events/
atif.rs

1//! Agent Trajectory Interchange Format (ATIF) types and builder.
2//!
3//! Implements the [ATIF specification](https://github.com/laude-institute/harbor/blob/main/docs/rfcs/0001-trajectory-format.md)
4//! v1.4 for logging complete agent interaction histories in a standardized,
5//! JSON-based format usable across debugging, visualization, SFT, and RL
6//! pipelines.
7//!
8//! # Overview
9//!
10//! ATIF provides a complete session trajectory: user messages, agent responses,
11//! tool executions, observations, and per-step/aggregate LLM metrics. The
12//! [`AtifTrajectoryBuilder`] converts live [`ThreadEvent`]
13//! streams into a finished [`Trajectory`].
14//!
15//! # Example
16//!
17//! ```rust
18//! use vtcode_exec_events::atif::*;
19//!
20//! let agent = AtifAgent::new("vtcode", env!("CARGO_PKG_VERSION"));
21//! let mut builder = AtifTrajectoryBuilder::new(agent);
22//!
23//! // Feed ThreadEvents as they arrive …
24//! // builder.process_event(&event);
25//!
26//! let trajectory = builder.finish(None);
27//! let json = serde_json::to_string_pretty(&trajectory).unwrap();
28//! ```
29
30use chrono::{DateTime, Utc};
31use serde::{Deserialize, Serialize};
32use serde_json::Value;
33
34use crate::{ThreadEvent, ThreadItemDetails, ToolCallStatus};
35
36/// Current ATIF schema version supported by this implementation.
37const ATIF_SCHEMA_VERSION: &str = "ATIF-v1.4";
38
39// ============================================================================
40// Core ATIF Types
41// ============================================================================
42
43/// Root-level ATIF trajectory object.
44#[derive(Debug, Clone, Serialize, Deserialize)]
45pub struct Trajectory {
46    /// ATIF schema version (e.g., "ATIF-v1.4").
47    schema_version: String,
48    /// Unique identifier for the entire agent run.
49    session_id: String,
50    /// Agent configuration for this trajectory.
51    agent: AtifAgent,
52    /// Ordered interaction steps.
53    steps: Vec<Step>,
54    /// Optional developer notes.
55    #[serde(skip_serializing_if = "Option::is_none")]
56    notes: Option<String>,
57    /// Aggregate metrics for the full trajectory.
58    #[serde(skip_serializing_if = "Option::is_none")]
59    pub final_metrics: Option<FinalMetrics>,
60    /// Optional custom root-level metadata.
61    #[serde(skip_serializing_if = "Option::is_none")]
62    extra: Option<Value>,
63}
64
65/// Agent configuration metadata.
66#[derive(Debug, Clone, Serialize, Deserialize)]
67pub struct AtifAgent {
68    /// Agent system name (e.g., "vtcode").
69    name: String,
70    /// Agent system version.
71    version: String,
72    /// Default LLM model used. Step-level `model_name` overrides this.
73    #[serde(skip_serializing_if = "Option::is_none")]
74    model_name: Option<String>,
75    /// Optional custom agent metadata.
76    #[serde(skip_serializing_if = "Option::is_none")]
77    extra: Option<Value>,
78}
79
80impl AtifAgent {
81    /// Create a new agent descriptor.
82    pub fn new(name: impl Into<String>, version: impl Into<String>) -> Self {
83        Self {
84            name: name.into(),
85            version: version.into(),
86            model_name: None,
87            extra: None,
88        }
89    }
90
91    /// Create a vtcode agent descriptor using the crate version.
92    pub fn vtcode() -> Self {
93        Self::new("vtcode", env!("CARGO_PKG_VERSION"))
94    }
95
96    /// Set the default model name.
97    pub fn with_model(mut self, model: impl Into<String>) -> Self {
98        self.model_name = Some(model.into());
99        self
100    }
101}
102
103/// The originator of a step.
104#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
105#[serde(rename_all = "lowercase")]
106pub enum StepSource {
107    /// System prompt or system-initiated operation.
108    System,
109    /// User message.
110    User,
111    /// Agent response.
112    Agent,
113}
114
115/// Individual interaction step.
116#[derive(Debug, Clone, Serialize, Deserialize)]
117pub struct Step {
118    /// Ordinal index (starting from 1).
119    step_id: u64,
120    /// ISO 8601 timestamp.
121    #[serde(skip_serializing_if = "Option::is_none")]
122    timestamp: Option<String>,
123    /// Originator of this step.
124    source: StepSource,
125    /// LLM model used for this step (agent steps only).
126    #[serde(skip_serializing_if = "Option::is_none")]
127    model_name: Option<String>,
128    /// Step content — text message or array.
129    #[serde(skip_serializing_if = "Option::is_none")]
130    message: Option<String>,
131    /// Agent internal reasoning content.
132    #[serde(skip_serializing_if = "Option::is_none")]
133    reasoning_content: Option<String>,
134    /// Tool/function invocations (agent steps only).
135    #[serde(skip_serializing_if = "Option::is_none")]
136    tool_calls: Option<Vec<AtifToolCall>>,
137    /// Environment feedback after actions.
138    #[serde(skip_serializing_if = "Option::is_none")]
139    observation: Option<Observation>,
140    /// LLM operational metrics (agent steps only).
141    #[serde(skip_serializing_if = "Option::is_none")]
142    metrics: Option<StepMetrics>,
143    /// Custom step-level metadata.
144    #[serde(skip_serializing_if = "Option::is_none")]
145    extra: Option<Value>,
146}
147
148impl Step {
149    /// Create a user step.
150    fn user(step_id: u64, message: impl Into<String>) -> Self {
151        Self {
152            step_id,
153            timestamp: Some(Utc::now().to_rfc3339()),
154            source: StepSource::User,
155            model_name: None,
156            message: Some(message.into()),
157            reasoning_content: None,
158            tool_calls: None,
159            observation: None,
160            metrics: None,
161            extra: None,
162        }
163    }
164
165    /// Create an agent step.
166    fn agent(step_id: u64, message: impl Into<String>) -> Self {
167        Self {
168            step_id,
169            timestamp: Some(Utc::now().to_rfc3339()),
170            source: StepSource::Agent,
171            model_name: None,
172            message: Some(message.into()),
173            reasoning_content: None,
174            tool_calls: None,
175            observation: None,
176            metrics: None,
177            extra: None,
178        }
179    }
180
181    /// Create a system step.
182    fn system(step_id: u64, message: impl Into<String>) -> Self {
183        Self {
184            step_id,
185            timestamp: Some(Utc::now().to_rfc3339()),
186            source: StepSource::System,
187            model_name: None,
188            message: Some(message.into()),
189            reasoning_content: None,
190            tool_calls: None,
191            observation: None,
192            metrics: None,
193            extra: None,
194        }
195    }
196}
197
198/// Structured tool/function invocation.
199#[derive(Debug, Clone, Serialize, Deserialize)]
200pub struct AtifToolCall {
201    /// Unique identifier for the tool call.
202    tool_call_id: String,
203    /// Function/tool name.
204    function_name: String,
205    /// Arguments passed to the tool.
206    #[serde(skip_serializing_if = "Option::is_none")]
207    arguments: Option<Value>,
208}
209
210/// Environment feedback container.
211#[derive(Debug, Clone, Serialize, Deserialize)]
212pub struct Observation {
213    /// Results from tool calls or system operations.
214    results: Vec<ObservationResult>,
215}
216
217/// Individual observation result tied to a tool call.
218#[derive(Debug, Clone, Serialize, Deserialize)]
219pub struct ObservationResult {
220    /// Identifier of the originating tool call.
221    source_call_id: String,
222    /// Content/output of the observation.
223    content: String,
224}
225
226/// Per-step LLM operational metrics.
227#[derive(Debug, Clone, Serialize, Deserialize)]
228pub struct StepMetrics {
229    /// Total input tokens for this step (cached + non-cached).
230    #[serde(skip_serializing_if = "Option::is_none")]
231    prompt_tokens: Option<u64>,
232    /// Completion tokens generated.
233    #[serde(skip_serializing_if = "Option::is_none")]
234    completion_tokens: Option<u64>,
235    /// Subset of prompt_tokens that were cache hits.
236    #[serde(skip_serializing_if = "Option::is_none")]
237    cached_tokens: Option<u64>,
238    /// Estimated cost in USD for this step.
239    #[serde(skip_serializing_if = "Option::is_none")]
240    cost_usd: Option<f64>,
241    /// Log probabilities for completion tokens.
242    #[serde(skip_serializing_if = "Option::is_none")]
243    logprobs: Option<Vec<f64>>,
244    /// Completion token IDs for RL training.
245    #[serde(skip_serializing_if = "Option::is_none")]
246    completion_token_ids: Option<Vec<u64>>,
247    /// Prompt token IDs.
248    #[serde(skip_serializing_if = "Option::is_none")]
249    prompt_token_ids: Option<Vec<u64>>,
250    /// Custom metrics.
251    #[serde(skip_serializing_if = "Option::is_none")]
252    extra: Option<Value>,
253}
254
255impl StepMetrics {
256    /// Create metrics from vtcode Usage.
257    fn from_usage(usage: &crate::Usage) -> Self {
258        Self {
259            prompt_tokens: Some(usage.input_tokens),
260            completion_tokens: Some(usage.output_tokens),
261            cached_tokens: if usage.cached_input_tokens > 0 {
262                Some(usage.cached_input_tokens)
263            } else {
264                None
265            },
266            cost_usd: None,
267            logprobs: None,
268            completion_token_ids: None,
269            prompt_token_ids: None,
270            extra: if usage.cache_creation_tokens > 0 {
271                Some(serde_json::json!({
272                    "cache_creation_tokens": usage.cache_creation_tokens
273                }))
274            } else {
275                None
276            },
277        }
278    }
279}
280
281/// Trajectory-level aggregate metrics.
282#[derive(Debug, Clone, Default, Serialize, Deserialize)]
283pub struct FinalMetrics {
284    /// Sum of all prompt tokens across steps.
285    #[serde(skip_serializing_if = "Option::is_none")]
286    pub total_prompt_tokens: Option<u64>,
287    /// Sum of all completion tokens across steps.
288    #[serde(skip_serializing_if = "Option::is_none")]
289    pub total_completion_tokens: Option<u64>,
290    /// Sum of all cached tokens across steps.
291    #[serde(skip_serializing_if = "Option::is_none")]
292    pub total_cached_tokens: Option<u64>,
293    /// Total estimated cost in USD.
294    #[serde(skip_serializing_if = "Option::is_none")]
295    total_cost_usd: Option<f64>,
296    /// Total number of steps.
297    #[serde(skip_serializing_if = "Option::is_none")]
298    total_steps: Option<u64>,
299    /// Custom aggregate metrics.
300    #[serde(skip_serializing_if = "Option::is_none")]
301    extra: Option<Value>,
302}
303
304// ============================================================================
305// Builder — converts live ThreadEvent streams into ATIF Trajectory
306// ============================================================================
307
308/// Stateful collector that converts a live [`ThreadEvent`] stream into an
309/// ATIF-compliant [`Trajectory`].
310///
311/// Feed events via [`process_event`](Self::process_event) (timestamps at
312/// observation time) or [`process_event_at`](Self::process_event_at)
313/// (deterministic timestamps for tests). Call [`finish`](Self::finish) to
314/// produce the final trajectory.
315pub struct AtifTrajectoryBuilder {
316    completed_at: Option<String>,
317    agent: AtifAgent,
318    session_id: Option<String>,
319    steps: Vec<Step>,
320    next_step_id: u64,
321    // Running token accumulators for final metrics
322    total_input_tokens: u64,
323    total_output_tokens: u64,
324    total_cached_tokens: u64,
325    num_turns: usize,
326    /// Whether any per-turn `TurnCompleted`/`TurnFailed` usage was observed.
327    /// Guards the `ThreadCompleted` aggregate from double-counting usage.
328    saw_per_turn_usage: bool,
329    /// Pending tool invocations awaiting matching ToolOutput.
330    pending_tool_calls: Vec<PendingToolCall>,
331}
332
333struct PendingToolCall {
334    call_id: String,
335    tool_call_id: Option<String>,
336    tool_name: String,
337    arguments: Option<Value>,
338    timestamp: String,
339}
340
341impl AtifTrajectoryBuilder {
342    /// Create a new builder for the given agent.
343    pub fn new(agent: AtifAgent) -> Self {
344        Self {
345            agent,
346            completed_at: None,
347            session_id: None,
348            steps: Vec::new(),
349            next_step_id: 1,
350            total_input_tokens: 0,
351            total_output_tokens: 0,
352            total_cached_tokens: 0,
353            num_turns: 0,
354            saw_per_turn_usage: false,
355            pending_tool_calls: Vec::new(),
356        }
357    }
358
359    /// Set the session ID explicitly. If not set, it will be derived from
360    /// `ThreadStarted` or `ThreadCompleted` events.
361    pub fn set_session_id(&mut self, id: impl Into<String>) {
362        self.session_id = Some(id.into());
363    }
364
365    /// Process a thread event using the current wall-clock time.
366    pub fn process_event(&mut self, event: &ThreadEvent) {
367        self.process_event_at(event, Utc::now());
368    }
369
370    /// Process a thread event with an explicit timestamp (for deterministic tests).
371    pub fn process_event_at(&mut self, event: &ThreadEvent, ts: DateTime<Utc>) {
372        let first_step = self.steps.len();
373        let ts_str = ts.to_rfc3339();
374        match event {
375            ThreadEvent::ThreadStarted(e) => {
376                if self.session_id.is_none() {
377                    self.session_id = Some(e.thread_id.clone());
378                }
379            }
380            ThreadEvent::ThreadCompleted(e) => {
381                self.completed_at = e.completed_at.clone();
382                if self.session_id.is_none() {
383                    self.session_id = Some(e.session_id.clone());
384                }
385                self.num_turns = e.num_turns;
386                // `ThreadCompleted` carries the *aggregate* usage, which equals
387                // the sum of per-turn usage already accumulated by the
388                // `TurnCompleted`/`TurnFailed` arms. Only fall back to it when
389                // no per-turn usage was observed (e.g. a harness that emits
390                // only the thread aggregate), otherwise totals double-count.
391                if !self.saw_per_turn_usage {
392                    self.total_input_tokens = self.total_input_tokens.saturating_add(e.usage.input_tokens);
393                    self.total_output_tokens = self.total_output_tokens.saturating_add(e.usage.output_tokens);
394                    self.total_cached_tokens = self.total_cached_tokens.saturating_add(e.usage.cached_input_tokens);
395                }
396            }
397            ThreadEvent::TurnCompleted(e) => {
398                self.saw_per_turn_usage = true;
399                self.total_input_tokens = self.total_input_tokens.saturating_add(e.usage.input_tokens);
400                self.total_output_tokens = self.total_output_tokens.saturating_add(e.usage.output_tokens);
401                self.total_cached_tokens = self.total_cached_tokens.saturating_add(e.usage.cached_input_tokens);
402                self.num_turns += 1;
403
404                let mut step = Step::system(self.next_step_id, "turn_completed");
405                step.timestamp = Some(ts_str);
406                step.metrics = Some(StepMetrics::from_usage(&e.usage));
407                if !e.in_progress_exec_sessions.is_empty() {
408                    step.extra = Some(serde_json::json!({
409                        "in_progress_exec_sessions": e.in_progress_exec_sessions,
410                    }));
411                }
412                self.push_step(step);
413            }
414            ThreadEvent::TurnFailed(e) => {
415                if let Some(usage) = &e.usage {
416                    self.saw_per_turn_usage = true;
417                    self.total_input_tokens = self.total_input_tokens.saturating_add(usage.input_tokens);
418                    self.total_output_tokens = self.total_output_tokens.saturating_add(usage.output_tokens);
419                    self.total_cached_tokens = self.total_cached_tokens.saturating_add(usage.cached_input_tokens);
420                }
421                let mut step = Step::system(self.next_step_id, &e.message);
422                step.timestamp = Some(ts_str);
423                step.metrics = e.usage.as_ref().map(StepMetrics::from_usage);
424                self.push_step(step);
425            }
426            ThreadEvent::TurnBlocked(e) => {
427                if let Some(usage) = &e.usage {
428                    self.saw_per_turn_usage = true;
429                    self.total_input_tokens = self.total_input_tokens.saturating_add(usage.input_tokens);
430                    self.total_output_tokens = self.total_output_tokens.saturating_add(usage.output_tokens);
431                    self.total_cached_tokens = self.total_cached_tokens.saturating_add(usage.cached_input_tokens);
432                }
433                let mut step = Step::system(self.next_step_id, &e.message);
434                step.timestamp = Some(ts_str);
435                step.metrics = e.usage.as_ref().map(StepMetrics::from_usage);
436                step.extra = Some(serde_json::json!({
437                    "last_tool": e.last_tool,
438                    "blocked_streak": e.blocked_streak,
439                    "blocked_total": e.blocked_total,
440                    "consecutive_cap": e.consecutive_cap,
441                    "total_cap": e.total_cap,
442                    "recovery_active": e.recovery_active,
443                }));
444                self.push_step(step);
445            }
446            ThreadEvent::ItemCompleted(e) => {
447                self.process_item_completed(&e.item.id, &e.item.details, &ts_str);
448            }
449            ThreadEvent::ThreadCompactBoundary(e) => {
450                let msg = format!(
451                    "context_compaction: {} messages -> {} messages ({})",
452                    e.original_message_count,
453                    e.compacted_message_count,
454                    e.trigger.as_str()
455                );
456                let mut step = Step::system(self.next_step_id, msg);
457                step.timestamp = Some(ts_str);
458                self.push_step(step);
459            }
460            ThreadEvent::ContextReset(e) => {
461                let msg = format!(
462                    "context_reset: {}% context used; plan preserved: {}; tool budget reset: {}",
463                    e.previous_context_usage_percent, e.plan_preserved, e.tool_budget_reset
464                );
465                let mut step = Step::system(self.next_step_id, msg);
466                step.timestamp = Some(ts_str);
467                step.extra = Some(serde_json::json!({
468                    "thread_id": e.thread_id,
469                    "turn_id": e.turn_id,
470                    "trigger": e.trigger,
471                    "plan_preserved": e.plan_preserved,
472                    "previous_context_usage_percent": e.previous_context_usage_percent,
473                    "tool_budget_reset": e.tool_budget_reset,
474                }));
475                self.push_step(step);
476            }
477            ThreadEvent::Error(e) => {
478                let mut step = Step::system(self.next_step_id, &e.message);
479                step.timestamp = Some(ts_str);
480                self.push_step(step);
481            }
482            // Skip streaming/lifecycle events that don't map to ATIF steps
483            ThreadEvent::TurnStarted(e) => {
484                if let Some(context) = &e.context {
485                    let mut step = Step::agent(self.next_step_id, context.goal.as_deref().unwrap_or("Turn started"));
486                    step.source = StepSource::User;
487                    step.timestamp = Some(context.timestamp.clone());
488                    step.extra = Some(serde_json::json!({"execution_context": context}));
489                    self.push_step(step);
490                }
491            }
492            ThreadEvent::MatrixUpdated(_)
493            | ThreadEvent::ItemStarted(_)
494            | ThreadEvent::ItemUpdated(_)
495            | ThreadEvent::PlanDelta(_)
496            | ThreadEvent::PlanApprovalRequested(_)
497            | ThreadEvent::PlanApprovalResolved(_)
498            | ThreadEvent::PermissionRequested(_)
499            | ThreadEvent::PermissionResolved(_)
500            | ThreadEvent::Interjected(_)
501            | ThreadEvent::Unknown => {}
502        }
503        for step in self.steps.iter_mut().skip(first_step) {
504            let context = match event {
505                ThreadEvent::ItemCompleted(e) => e.item.context.as_deref(),
506                _ => None,
507            };
508            if let Some(context) = context {
509                step.timestamp = Some(context.timestamp.clone());
510                let extra = step.extra.get_or_insert_with(|| serde_json::json!({}));
511                if let Some(extra) = extra.as_object_mut() {
512                    let _ = extra.insert("item_context".into(), serde_json::json!(context));
513                }
514            }
515            match event {
516                ThreadEvent::TurnCompleted(e) => {
517                    if let Some(ts) = &e.completed_at {
518                        step.timestamp = Some(ts.to_string());
519                    }
520                }
521                ThreadEvent::TurnFailed(e) => {
522                    if let Some(ts) = &e.completed_at {
523                        step.timestamp = Some(ts.to_string());
524                    }
525                }
526                ThreadEvent::TurnBlocked(e) => {
527                    if let Some(ts) = &e.completed_at {
528                        step.timestamp = Some(ts.clone());
529                    }
530                }
531                _ => {}
532            }
533        }
534    }
535
536    fn process_item_completed(&mut self, item_id: &str, details: &ThreadItemDetails, ts: &str) {
537        match details {
538            ThreadItemDetails::Decision(d) => {
539                let mut step = Step::agent(self.next_step_id, &d.summary);
540                step.timestamp = Some(ts.to_owned());
541                step.extra = Some(
542                    serde_json::json!({"vtcode_item_type": "decision", "public_rationale": d.rationale, "alternatives": d.alternatives, "evidence_ids": d.evidence_ids}),
543                );
544                self.push_step(step);
545            }
546            ThreadItemDetails::AgentMessage(msg) => {
547                let mut step = Step::agent(self.next_step_id, &msg.text);
548                step.timestamp = Some(ts.to_string());
549                self.push_step(step);
550            }
551            ThreadItemDetails::Plan(plan) => {
552                let mut step = Step::agent(self.next_step_id, &plan.text);
553                step.timestamp = Some(ts.to_string());
554                step.extra = Some(serde_json::json!({ "vtcode_item_type": "plan" }));
555                self.push_step(step);
556            }
557            ThreadItemDetails::Reasoning(r) => {
558                let mut step = Step::agent(self.next_step_id, "");
559                step.timestamp = Some(ts.to_string());
560                step.reasoning_content = Some(r.text.clone());
561                step.message = None;
562                self.push_step(step);
563            }
564            ThreadItemDetails::ToolInvocation(inv) => {
565                // Buffer the invocation; we'll pair it with the ToolOutput
566                self.pending_tool_calls.push(PendingToolCall {
567                    call_id: item_id.to_string(),
568                    tool_call_id: inv.tool_call_id.clone(),
569                    tool_name: inv.tool_name.clone(),
570                    arguments: inv.arguments.clone(),
571                    timestamp: ts.to_string(),
572                });
573            }
574            ThreadItemDetails::ToolOutput(output) => {
575                // Find the matching pending invocation
576                let pending_idx = self.pending_tool_calls.iter().position(|p| p.call_id == output.call_id);
577
578                let (tool_name, arguments, tool_call_id, inv_ts) = if let Some(idx) = pending_idx {
579                    let p = self.pending_tool_calls.remove(idx);
580                    (p.tool_name, p.arguments, p.tool_call_id, p.timestamp)
581                } else {
582                    ("unknown".to_string(), None, output.tool_call_id.clone(), ts.to_string())
583                };
584
585                let call_id = tool_call_id.clone().unwrap_or_else(|| output.call_id.clone());
586
587                let mut step = Step::agent(self.next_step_id, "");
588                step.timestamp = Some(inv_ts);
589                step.message = None;
590                step.tool_calls = Some(vec![AtifToolCall {
591                    tool_call_id: call_id.clone(),
592                    function_name: tool_name,
593                    arguments,
594                }]);
595
596                let status_suffix = match output.status {
597                    ToolCallStatus::Failed => " [FAILED]",
598                    ToolCallStatus::InProgress => " [IN_PROGRESS]",
599                    ToolCallStatus::Completed => "",
600                };
601                let content = format!("{}{}", output.output, status_suffix);
602                step.observation = Some(Observation {
603                    results: vec![ObservationResult { source_call_id: call_id, content }],
604                });
605                self.push_step(step);
606            }
607            ThreadItemDetails::CommandExecution(cmd) => {
608                let call_id = item_id.to_string();
609                let mut step = Step::agent(self.next_step_id, "");
610                step.timestamp = Some(ts.to_string());
611                step.message = None;
612                step.tool_calls = Some(vec![AtifToolCall {
613                    tool_call_id: call_id.clone(),
614                    function_name: "command_execution".to_string(),
615                    arguments: Some(serde_json::json!({
616                        "command": cmd.command,
617                        "arguments": cmd.arguments,
618                    })),
619                }]);
620                step.observation = Some(Observation {
621                    results: vec![ObservationResult {
622                        source_call_id: call_id,
623                        content: cmd.aggregated_output.clone(),
624                    }],
625                });
626                if let Some(exit_code) = cmd.exit_code {
627                    step.extra = Some(serde_json::json!({ "exit_code": exit_code }));
628                }
629                self.push_step(step);
630            }
631            ThreadItemDetails::McpToolCall(mcp) => {
632                let call_id = item_id.to_string();
633                let mut step = Step::agent(self.next_step_id, "");
634                step.timestamp = Some(ts.to_string());
635                step.message = None;
636                step.tool_calls = Some(vec![AtifToolCall {
637                    tool_call_id: call_id.clone(),
638                    function_name: mcp.tool_name.clone(),
639                    arguments: mcp.arguments.clone(),
640                }]);
641                if let Some(result) = &mcp.result {
642                    step.observation = Some(Observation {
643                        results: vec![ObservationResult { source_call_id: call_id, content: result.clone() }],
644                    });
645                }
646                self.push_step(step);
647            }
648            ThreadItemDetails::FileChange(fc) => {
649                let changes: Vec<String> = fc.changes.iter().map(|c| format!("{}: {:?}", c.path, c.kind)).collect();
650                let msg = format!("file_changes: {}", changes.join(", "));
651                let mut step = Step::system(self.next_step_id, msg);
652                step.timestamp = Some(ts.to_string());
653                self.push_step(step);
654            }
655            ThreadItemDetails::WebSearch(ws) => {
656                let mut step = Step::system(self.next_step_id, format!("web_search: {}", ws.query));
657                step.timestamp = Some(ts.to_string());
658                if let Some(results) = &ws.results {
659                    step.observation = Some(Observation {
660                        results: results
661                            .iter()
662                            .enumerate()
663                            .map(|(i, r)| ObservationResult {
664                                source_call_id: format!("search_{i}"),
665                                content: r.clone(),
666                            })
667                            .collect(),
668                    });
669                }
670                self.push_step(step);
671            }
672            ThreadItemDetails::Harness(h) => {
673                let msg = format!("harness: {:?}", h.event);
674                let mut step = Step::system(self.next_step_id, msg);
675                step.timestamp = Some(ts.to_string());
676                let mut extra = serde_json::Map::new();
677                if let Some(m) = &h.message {
678                    let _ = extra.insert("harness_message".to_string(), Value::String(m.clone()));
679                }
680                if h.event == crate::HarnessEventKind::BackgroundSubprocessCompleted {
681                    for (key, value) in [
682                        ("task_id", h.task_id.as_ref()),
683                        ("session_id", h.session_id.as_ref()),
684                        ("exec_session_id", h.exec_session_id.as_ref()),
685                        ("status", h.status.as_ref()),
686                        ("transcript_path", h.transcript_path.as_ref()),
687                        ("archive_path", h.archive_path.as_ref()),
688                        ("error_category", h.error_category.as_ref()),
689                    ] {
690                        if let Some(value) = value {
691                            let _ = extra.insert(key.to_string(), Value::String(value.clone()));
692                        }
693                    }
694                    if let Some(exit_code) = h.exit_code {
695                        let _ = extra.insert("exit_code".to_string(), Value::from(exit_code));
696                    }
697                }
698                if !extra.is_empty() {
699                    step.extra = Some(Value::Object(extra));
700                }
701                self.push_step(step);
702            }
703            ThreadItemDetails::Error(e) => {
704                let mut step = Step::system(self.next_step_id, &e.message);
705                step.timestamp = Some(ts.to_string());
706                self.push_step(step);
707            }
708        }
709    }
710
711    fn push_step(&mut self, step: Step) {
712        self.next_step_id = step.step_id + 1;
713        self.steps.push(step);
714    }
715
716    /// Consume the builder and produce the final ATIF trajectory.
717    ///
718    /// Pass optional `FinalMetrics` to override the accumulated values.
719    /// If `None`, final metrics are derived from observed events.
720    pub fn finish(self, override_metrics: Option<FinalMetrics>) -> Trajectory {
721        let final_metrics = override_metrics.unwrap_or_else(|| FinalMetrics {
722            total_prompt_tokens: Some(self.total_input_tokens),
723            total_completion_tokens: Some(self.total_output_tokens),
724            total_cached_tokens: if self.total_cached_tokens > 0 {
725                Some(self.total_cached_tokens)
726            } else {
727                None
728            },
729            total_cost_usd: None,
730            total_steps: Some(self.steps.len() as u64),
731            extra: Some(serde_json::json!({ "num_turns": self.num_turns })),
732        });
733
734        Trajectory {
735            schema_version: ATIF_SCHEMA_VERSION.to_string(),
736            session_id: self.session_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()),
737            agent: self.agent,
738            steps: self.steps,
739            notes: None,
740            final_metrics: Some(final_metrics),
741            extra: self.completed_at.map(|timestamp| serde_json::json!({"completed_at":timestamp})),
742        }
743    }
744
745    /// Returns the number of steps collected so far.
746    pub fn step_count(&self) -> usize {
747        self.steps.len()
748    }
749}
750
751impl crate::EventEmitter for AtifTrajectoryBuilder {
752    fn emit(&mut self, event: &ThreadEvent) {
753        self.process_event(event);
754    }
755}
756
757#[cfg(test)]
758mod tests {
759    use super::*;
760    use crate::{
761        AgentMessageItem, CompactionMode, CompactionTrigger, HarnessEventItem, HarnessEventKind, ItemCompletedEvent,
762        ThreadCompactBoundaryEvent, ThreadItem, ThreadStartedEvent, ToolInvocationItem, ToolOutputItem,
763        TurnCompletedEvent, TurnStartedEvent, Usage,
764    };
765
766    fn fixed_ts() -> DateTime<Utc> {
767        "2025-01-15T10:30:00Z".parse().unwrap()
768    }
769
770    #[test]
771    fn terminal_timestamps_survive_export_without_synthetic_steps() {
772        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
773        let timestamp = "2026-10-03T01:02:03Z";
774        let blocked = serde_json::from_value(
775            serde_json::json!({"type":"turn.blocked", "message":"Blocked", "completed_at":timestamp}),
776        )
777        .unwrap();
778        builder.process_event_at(&blocked, fixed_ts());
779        let completed = serde_json::from_value(serde_json::json!({"type":"thread.completed", "completed_at":timestamp, "thread_id":"thread", "session_id":"session", "subtype":"success", "outcome_code":"completed", "usage":Usage::default(), "num_turns":1})).unwrap();
780        builder.process_event_at(&completed, fixed_ts());
781        let trajectory = builder.finish(None);
782        assert_eq!(trajectory.steps.len(), 1);
783        assert_eq!(trajectory.steps[0].timestamp.as_deref(), Some(timestamp));
784        assert_eq!(trajectory.extra.as_ref().unwrap()["completed_at"], timestamp);
785    }
786
787    #[test]
788    fn explanation_metadata_and_public_rationale_preserve_recorded_time() {
789        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
790        let context = crate::ItemContext {
791            task_id: "task-1".into(),
792            turn_id: "turn-1".into(),
793            actor_id: "root".into(),
794            parent_actor_id: None,
795            timestamp: "2026-10-03T01:02:03Z".into(),
796            activity: None,
797        };
798        builder.process_event_at(
799            &ThreadEvent::ItemCompleted(ItemCompletedEvent {
800                item: ThreadItem {
801                    id: "decision-1".into(),
802                    context: Some(Box::new(context.clone())),
803                    details: ThreadItemDetails::Decision(Box::new(crate::DecisionItem {
804                        summary: "Reuse parser".into(),
805                        rationale: "Preserve checks".into(),
806                        alternatives: vec!["Replace parser".into()],
807                        evidence_ids: vec!["read-1".into()],
808                    })),
809                },
810            }),
811            fixed_ts(),
812        );
813        let trajectory = builder.finish(None);
814        assert_eq!(trajectory.steps.len(), 1);
815        let step = &trajectory.steps[0];
816        assert_eq!(step.timestamp.as_deref(), Some(context.timestamp.as_str()));
817        let extra = step.extra.as_ref().unwrap();
818        assert_eq!(extra["public_rationale"], "Preserve checks");
819        assert_eq!(extra["item_context"]["task_id"], "task-1");
820        assert_eq!(extra["evidence_ids"][0], "read-1");
821    }
822
823    #[test]
824    fn trajectory_round_trip() {
825        let trajectory = Trajectory {
826            schema_version: ATIF_SCHEMA_VERSION.to_string(),
827            session_id: "test-session".to_string(),
828            agent: AtifAgent::vtcode(),
829            steps: vec![Step::user(1, "hello")],
830            notes: None,
831            final_metrics: None,
832            extra: None,
833        };
834
835        let json = serde_json::to_string_pretty(&trajectory).unwrap();
836        let restored: Trajectory = serde_json::from_str(&json).unwrap();
837        assert_eq!(restored.schema_version, ATIF_SCHEMA_VERSION);
838        assert_eq!(restored.session_id, "test-session");
839        assert_eq!(restored.steps.len(), 1);
840    }
841
842    #[test]
843    fn builder_thread_started_sets_session_id() {
844        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
845        let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread-abc".to_string() });
846        builder.process_event_at(&event, fixed_ts());
847        let trajectory = builder.finish(None);
848        assert_eq!(trajectory.session_id, "thread-abc");
849    }
850
851    #[test]
852    fn builder_agent_message_step() {
853        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
854        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
855            item: ThreadItem {
856                context: None,
857                id: "msg-1".to_string(),
858                details: ThreadItemDetails::AgentMessage(AgentMessageItem { text: "Hello, world!".to_string() }),
859            },
860        });
861        builder.process_event_at(&event, fixed_ts());
862        let trajectory = builder.finish(None);
863
864        assert_eq!(trajectory.steps.len(), 1);
865        let step = &trajectory.steps[0];
866        assert_eq!(step.step_id, 1);
867        assert_eq!(step.source, StepSource::Agent);
868        assert_eq!(step.message.as_deref(), Some("Hello, world!"));
869    }
870
871    #[test]
872    fn background_completion_preserves_identity_in_atif_extra() {
873        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
874        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
875            item: ThreadItem {
876                context: None,
877                id: "background-completion:task-1:exec-1".to_string(),
878                details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
879                    event: HarnessEventKind::BackgroundSubprocessCompleted,
880                    message: Some("completed successfully".to_string()),
881                    command: None,
882                    path: None,
883                    exit_code: Some(0),
884                    attempt: None,
885                    error_category: Some("background_subprocess".to_string()),
886                    duration_ms: None,
887                    task_id: Some("task-1".to_string()),
888                    session_id: Some("session-1".to_string()),
889                    exec_session_id: Some("exec-1".to_string()),
890                    status: Some("stopped".to_string()),
891                    transcript_path: Some("/tmp/transcript.jsonl".to_string()),
892                    archive_path: Some("/tmp/archive.json".to_string()),
893                })),
894            },
895        });
896
897        builder.process_event_at(&event, fixed_ts());
898        let trajectory = builder.finish(None);
899        let step = trajectory.steps.first().expect("background completion step");
900        let extra = step.extra.as_ref().expect("background completion metadata");
901        assert_eq!(step.message.as_deref(), Some("harness: BackgroundSubprocessCompleted"));
902        assert_eq!(extra["harness_message"], "completed successfully");
903        assert_eq!(extra["task_id"], "task-1");
904        assert_eq!(extra["session_id"], "session-1");
905        assert_eq!(extra["exec_session_id"], "exec-1");
906        assert_eq!(extra["status"], "stopped");
907        assert_eq!(extra["exit_code"], 0);
908        assert_eq!(extra["transcript_path"], "/tmp/transcript.jsonl");
909        assert_eq!(extra["archive_path"], "/tmp/archive.json");
910        assert_eq!(extra["error_category"], "background_subprocess");
911    }
912
913    #[test]
914    fn builder_tool_invocation_with_output() {
915        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
916        let ts = fixed_ts();
917
918        // Tool invocation
919        let inv_event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
920            item: ThreadItem {
921                context: None,
922                id: "tool_1".to_string(),
923                details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
924                    tool_name: "read_file".to_string(),
925                    arguments: Some(serde_json::json!({"path": "README.md"})),
926                    tool_call_id: Some("tc_0".to_string()),
927                    status: ToolCallStatus::Completed,
928                    outcome: None,
929                })),
930            },
931        });
932        builder.process_event_at(&inv_event, ts);
933
934        // Tool output
935        let out_event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
936            item: ThreadItem {
937                context: None,
938                id: "tool_1:output".to_string(),
939                details: ThreadItemDetails::ToolOutput(Box::new(ToolOutputItem {
940                    call_id: "tool_1".to_string(),
941                    tool_call_id: Some("tc_0".to_string()),
942                    spool_path: None,
943                    output: "file contents here".to_string(),
944                    exit_code: Some(0),
945                    status: ToolCallStatus::Completed,
946                })),
947            },
948        });
949        builder.process_event_at(&out_event, ts);
950
951        let trajectory = builder.finish(None);
952        // Only one step: the invocation is buffered until output arrives
953        assert_eq!(trajectory.steps.len(), 1);
954        let step = &trajectory.steps[0];
955        assert_eq!(step.source, StepSource::Agent);
956
957        let calls = step.tool_calls.as_ref().unwrap();
958        assert_eq!(calls.len(), 1);
959        assert_eq!(calls[0].function_name, "read_file");
960        assert_eq!(calls[0].tool_call_id, "tc_0");
961
962        let obs = step.observation.as_ref().unwrap();
963        assert_eq!(obs.results.len(), 1);
964        assert_eq!(obs.results[0].content, "file contents here");
965    }
966
967    #[test]
968    fn builder_turn_completed_accumulates_metrics() {
969        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
970        let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
971            completed_at: None,
972            usage: Usage {
973                input_tokens: 500,
974                cached_input_tokens: 100,
975                cache_creation_tokens: 0,
976                output_tokens: 200,
977            },
978            in_progress_exec_sessions: Vec::new(),
979        });
980        builder.process_event_at(&event, fixed_ts());
981
982        let trajectory = builder.finish(None);
983        let fm = trajectory.final_metrics.as_ref().unwrap();
984        assert_eq!(fm.total_prompt_tokens, Some(500));
985        assert_eq!(fm.total_completion_tokens, Some(200));
986        assert_eq!(fm.total_cached_tokens, Some(100));
987    }
988
989    #[test]
990    fn builder_turn_completed_preserves_in_progress_sessions_for_resume() {
991        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
992        let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
993            completed_at: None,
994            usage: Usage::default(),
995            in_progress_exec_sessions: vec!["run-7".to_string()],
996        });
997        builder.process_event_at(&event, fixed_ts());
998
999        let trajectory = builder.finish(None);
1000        let step = trajectory.steps.last().expect("turn_completed step");
1001        let extra = step.extra.clone().expect("extra carries resume ids");
1002        assert_eq!(extra["in_progress_exec_sessions"], serde_json::json!(["run-7"]));
1003
1004        // Empty ids stay omitted so steady-state export is unchanged.
1005        let mut empty_builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1006        let empty = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1007            completed_at: None,
1008            usage: Usage::default(),
1009            in_progress_exec_sessions: Vec::new(),
1010        });
1011        empty_builder.process_event_at(&empty, fixed_ts());
1012        let empty_trajectory = empty_builder.finish(None);
1013        assert!(empty_trajectory.steps.last().expect("step").extra.is_none());
1014    }
1015
1016    #[test]
1017    fn step_metrics_from_usage() {
1018        let usage = Usage {
1019            input_tokens: 1000,
1020            cached_input_tokens: 200,
1021            cache_creation_tokens: 50,
1022            output_tokens: 300,
1023        };
1024        let metrics = StepMetrics::from_usage(&usage);
1025        assert_eq!(metrics.prompt_tokens, Some(1000));
1026        assert_eq!(metrics.completion_tokens, Some(300));
1027        assert_eq!(metrics.cached_tokens, Some(200));
1028        assert!(metrics.extra.is_some());
1029    }
1030
1031    #[test]
1032    fn builder_implements_event_emitter() {
1033        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1034        let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "t-1".to_string() });
1035        // Use EventEmitter trait
1036        crate::EventEmitter::emit(&mut builder, &event);
1037        assert_eq!(builder.step_count(), 0); // ThreadStarted doesn't create a step
1038    }
1039
1040    #[test]
1041    fn skips_lifecycle_events() {
1042        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1043        builder.process_event(&ThreadEvent::TurnStarted(TurnStartedEvent::default()));
1044        assert_eq!(builder.step_count(), 0);
1045    }
1046
1047    #[test]
1048    fn compact_boundary_with_segment_metadata_preserves_atif_export() {
1049        let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1050        let event = ThreadEvent::ThreadCompactBoundary(Box::new(ThreadCompactBoundaryEvent {
1051            thread_id: "thread-1".to_string(),
1052            trigger: CompactionTrigger::Auto,
1053            mode: CompactionMode::Local,
1054            original_message_count: 12,
1055            compacted_message_count: 5,
1056            history_artifact_path: None,
1057            previous_segment_id: Some("segment-0001".to_string()),
1058            new_segment_id: Some("segment-0002".to_string()),
1059            previous_prefix_hash: Some("prefix-before".to_string()),
1060            new_prefix_hash: Some("prefix-after".to_string()),
1061            previous_catalog_hash: Some("catalog-before".to_string()),
1062            new_catalog_hash: Some("catalog-after".to_string()),
1063        }));
1064
1065        builder.process_event_at(&event, fixed_ts());
1066        let trajectory = builder.finish(None);
1067
1068        assert_eq!(trajectory.steps.len(), 1);
1069        assert_eq!(trajectory.steps[0].source, StepSource::System);
1070        assert_eq!(
1071            trajectory.steps[0].message.as_deref(),
1072            Some("context_compaction: 12 messages -> 5 messages (auto)")
1073        );
1074    }
1075}