1use chrono::{DateTime, Utc};
31use serde::{Deserialize, Serialize};
32use serde_json::Value;
33
34use crate::{ThreadEvent, ThreadItemDetails, ToolCallStatus};
35
36const ATIF_SCHEMA_VERSION: &str = "ATIF-v1.4";
38
39#[derive(Debug, Clone, Serialize, Deserialize)]
45pub struct Trajectory {
46 schema_version: String,
48 session_id: String,
50 agent: AtifAgent,
52 steps: Vec<Step>,
54 #[serde(skip_serializing_if = "Option::is_none")]
56 notes: Option<String>,
57 #[serde(skip_serializing_if = "Option::is_none")]
59 pub final_metrics: Option<FinalMetrics>,
60 #[serde(skip_serializing_if = "Option::is_none")]
62 extra: Option<Value>,
63}
64
65#[derive(Debug, Clone, Serialize, Deserialize)]
67pub struct AtifAgent {
68 name: String,
70 version: String,
72 #[serde(skip_serializing_if = "Option::is_none")]
74 model_name: Option<String>,
75 #[serde(skip_serializing_if = "Option::is_none")]
77 extra: Option<Value>,
78}
79
80impl AtifAgent {
81 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 pub fn vtcode() -> Self {
93 Self::new("vtcode", env!("CARGO_PKG_VERSION"))
94 }
95
96 pub fn with_model(mut self, model: impl Into<String>) -> Self {
98 self.model_name = Some(model.into());
99 self
100 }
101}
102
103#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
105#[serde(rename_all = "lowercase")]
106pub enum StepSource {
107 System,
109 User,
111 Agent,
113}
114
115#[derive(Debug, Clone, Serialize, Deserialize)]
117pub struct Step {
118 step_id: u64,
120 #[serde(skip_serializing_if = "Option::is_none")]
122 timestamp: Option<String>,
123 source: StepSource,
125 #[serde(skip_serializing_if = "Option::is_none")]
127 model_name: Option<String>,
128 #[serde(skip_serializing_if = "Option::is_none")]
130 message: Option<String>,
131 #[serde(skip_serializing_if = "Option::is_none")]
133 reasoning_content: Option<String>,
134 #[serde(skip_serializing_if = "Option::is_none")]
136 tool_calls: Option<Vec<AtifToolCall>>,
137 #[serde(skip_serializing_if = "Option::is_none")]
139 observation: Option<Observation>,
140 #[serde(skip_serializing_if = "Option::is_none")]
142 metrics: Option<StepMetrics>,
143 #[serde(skip_serializing_if = "Option::is_none")]
145 extra: Option<Value>,
146}
147
148impl Step {
149 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 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 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#[derive(Debug, Clone, Serialize, Deserialize)]
200pub struct AtifToolCall {
201 tool_call_id: String,
203 function_name: String,
205 #[serde(skip_serializing_if = "Option::is_none")]
207 arguments: Option<Value>,
208}
209
210#[derive(Debug, Clone, Serialize, Deserialize)]
212pub struct Observation {
213 results: Vec<ObservationResult>,
215}
216
217#[derive(Debug, Clone, Serialize, Deserialize)]
219pub struct ObservationResult {
220 source_call_id: String,
222 content: String,
224}
225
226#[derive(Debug, Clone, Serialize, Deserialize)]
228pub struct StepMetrics {
229 #[serde(skip_serializing_if = "Option::is_none")]
231 prompt_tokens: Option<u64>,
232 #[serde(skip_serializing_if = "Option::is_none")]
234 completion_tokens: Option<u64>,
235 #[serde(skip_serializing_if = "Option::is_none")]
237 cached_tokens: Option<u64>,
238 #[serde(skip_serializing_if = "Option::is_none")]
240 cost_usd: Option<f64>,
241 #[serde(skip_serializing_if = "Option::is_none")]
243 logprobs: Option<Vec<f64>>,
244 #[serde(skip_serializing_if = "Option::is_none")]
246 completion_token_ids: Option<Vec<u64>>,
247 #[serde(skip_serializing_if = "Option::is_none")]
249 prompt_token_ids: Option<Vec<u64>>,
250 #[serde(skip_serializing_if = "Option::is_none")]
252 extra: Option<Value>,
253}
254
255impl StepMetrics {
256 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#[derive(Debug, Clone, Default, Serialize, Deserialize)]
283pub struct FinalMetrics {
284 #[serde(skip_serializing_if = "Option::is_none")]
286 pub total_prompt_tokens: Option<u64>,
287 #[serde(skip_serializing_if = "Option::is_none")]
289 pub total_completion_tokens: Option<u64>,
290 #[serde(skip_serializing_if = "Option::is_none")]
292 pub total_cached_tokens: Option<u64>,
293 #[serde(skip_serializing_if = "Option::is_none")]
295 total_cost_usd: Option<f64>,
296 #[serde(skip_serializing_if = "Option::is_none")]
298 total_steps: Option<u64>,
299 #[serde(skip_serializing_if = "Option::is_none")]
301 extra: Option<Value>,
302}
303
304pub 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 total_input_tokens: u64,
323 total_output_tokens: u64,
324 total_cached_tokens: u64,
325 num_turns: usize,
326 saw_per_turn_usage: bool,
329 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 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 pub fn set_session_id(&mut self, id: impl Into<String>) {
362 self.session_id = Some(id.into());
363 }
364
365 pub fn process_event(&mut self, event: &ThreadEvent) {
367 self.process_event_at(event, Utc::now());
368 }
369
370 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 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 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 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 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 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 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 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 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 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 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 crate::EventEmitter::emit(&mut builder, &event);
1037 assert_eq!(builder.step_count(), 0); }
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}