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::ItemStarted(_)
493 | ThreadEvent::ItemUpdated(_)
494 | ThreadEvent::PlanDelta(_)
495 | ThreadEvent::PlanApprovalRequested(_)
496 | ThreadEvent::PlanApprovalResolved(_)
497 | ThreadEvent::PermissionRequested(_)
498 | ThreadEvent::PermissionResolved(_)
499 | ThreadEvent::Interjected(_)
500 | ThreadEvent::Unknown => {}
501 }
502 for step in self.steps.iter_mut().skip(first_step) {
503 let context = match event {
504 ThreadEvent::ItemCompleted(e) => e.item.context.as_deref(),
505 _ => None,
506 };
507 if let Some(context) = context {
508 step.timestamp = Some(context.timestamp.clone());
509 let extra = step.extra.get_or_insert_with(|| serde_json::json!({}));
510 if let Some(extra) = extra.as_object_mut() {
511 let _ = extra.insert("item_context".into(), serde_json::json!(context));
512 }
513 }
514 match event {
515 ThreadEvent::TurnCompleted(e) => {
516 if let Some(ts) = &e.completed_at {
517 step.timestamp = Some(ts.to_string());
518 }
519 }
520 ThreadEvent::TurnFailed(e) => {
521 if let Some(ts) = &e.completed_at {
522 step.timestamp = Some(ts.to_string());
523 }
524 }
525 ThreadEvent::TurnBlocked(e) => {
526 if let Some(ts) = &e.completed_at {
527 step.timestamp = Some(ts.clone());
528 }
529 }
530 _ => {}
531 }
532 }
533 }
534
535 fn process_item_completed(&mut self, item_id: &str, details: &ThreadItemDetails, ts: &str) {
536 match details {
537 ThreadItemDetails::Decision(d) => {
538 let mut step = Step::agent(self.next_step_id, &d.summary);
539 step.timestamp = Some(ts.to_owned());
540 step.extra = Some(
541 serde_json::json!({"vtcode_item_type": "decision", "public_rationale": d.rationale, "alternatives": d.alternatives, "evidence_ids": d.evidence_ids}),
542 );
543 self.push_step(step);
544 }
545 ThreadItemDetails::AgentMessage(msg) => {
546 let mut step = Step::agent(self.next_step_id, &msg.text);
547 step.timestamp = Some(ts.to_string());
548 self.push_step(step);
549 }
550 ThreadItemDetails::Plan(plan) => {
551 let mut step = Step::agent(self.next_step_id, &plan.text);
552 step.timestamp = Some(ts.to_string());
553 step.extra = Some(serde_json::json!({ "vtcode_item_type": "plan" }));
554 self.push_step(step);
555 }
556 ThreadItemDetails::Reasoning(r) => {
557 let mut step = Step::agent(self.next_step_id, "");
558 step.timestamp = Some(ts.to_string());
559 step.reasoning_content = Some(r.text.clone());
560 step.message = None;
561 self.push_step(step);
562 }
563 ThreadItemDetails::ToolInvocation(inv) => {
564 self.pending_tool_calls.push(PendingToolCall {
566 call_id: item_id.to_string(),
567 tool_call_id: inv.tool_call_id.clone(),
568 tool_name: inv.tool_name.clone(),
569 arguments: inv.arguments.clone(),
570 timestamp: ts.to_string(),
571 });
572 }
573 ThreadItemDetails::ToolOutput(output) => {
574 let pending_idx = self.pending_tool_calls.iter().position(|p| p.call_id == output.call_id);
576
577 let (tool_name, arguments, tool_call_id, inv_ts) = if let Some(idx) = pending_idx {
578 let p = self.pending_tool_calls.remove(idx);
579 (p.tool_name, p.arguments, p.tool_call_id, p.timestamp)
580 } else {
581 ("unknown".to_string(), None, output.tool_call_id.clone(), ts.to_string())
582 };
583
584 let call_id = tool_call_id.clone().unwrap_or_else(|| output.call_id.clone());
585
586 let mut step = Step::agent(self.next_step_id, "");
587 step.timestamp = Some(inv_ts);
588 step.message = None;
589 step.tool_calls = Some(vec![AtifToolCall {
590 tool_call_id: call_id.clone(),
591 function_name: tool_name,
592 arguments,
593 }]);
594
595 let status_suffix = match output.status {
596 ToolCallStatus::Failed => " [FAILED]",
597 ToolCallStatus::InProgress => " [IN_PROGRESS]",
598 ToolCallStatus::Completed => "",
599 };
600 let content = format!("{}{}", output.output, status_suffix);
601 step.observation = Some(Observation {
602 results: vec![ObservationResult { source_call_id: call_id, content }],
603 });
604 self.push_step(step);
605 }
606 ThreadItemDetails::CommandExecution(cmd) => {
607 let call_id = item_id.to_string();
608 let mut step = Step::agent(self.next_step_id, "");
609 step.timestamp = Some(ts.to_string());
610 step.message = None;
611 step.tool_calls = Some(vec![AtifToolCall {
612 tool_call_id: call_id.clone(),
613 function_name: "command_execution".to_string(),
614 arguments: Some(serde_json::json!({
615 "command": cmd.command,
616 "arguments": cmd.arguments,
617 })),
618 }]);
619 step.observation = Some(Observation {
620 results: vec![ObservationResult {
621 source_call_id: call_id,
622 content: cmd.aggregated_output.clone(),
623 }],
624 });
625 if let Some(exit_code) = cmd.exit_code {
626 step.extra = Some(serde_json::json!({ "exit_code": exit_code }));
627 }
628 self.push_step(step);
629 }
630 ThreadItemDetails::McpToolCall(mcp) => {
631 let call_id = item_id.to_string();
632 let mut step = Step::agent(self.next_step_id, "");
633 step.timestamp = Some(ts.to_string());
634 step.message = None;
635 step.tool_calls = Some(vec![AtifToolCall {
636 tool_call_id: call_id.clone(),
637 function_name: mcp.tool_name.clone(),
638 arguments: mcp.arguments.clone(),
639 }]);
640 if let Some(result) = &mcp.result {
641 step.observation = Some(Observation {
642 results: vec![ObservationResult { source_call_id: call_id, content: result.clone() }],
643 });
644 }
645 self.push_step(step);
646 }
647 ThreadItemDetails::FileChange(fc) => {
648 let changes: Vec<String> = fc.changes.iter().map(|c| format!("{}: {:?}", c.path, c.kind)).collect();
649 let msg = format!("file_changes: {}", changes.join(", "));
650 let mut step = Step::system(self.next_step_id, msg);
651 step.timestamp = Some(ts.to_string());
652 self.push_step(step);
653 }
654 ThreadItemDetails::WebSearch(ws) => {
655 let mut step = Step::system(self.next_step_id, format!("web_search: {}", ws.query));
656 step.timestamp = Some(ts.to_string());
657 if let Some(results) = &ws.results {
658 step.observation = Some(Observation {
659 results: results
660 .iter()
661 .enumerate()
662 .map(|(i, r)| ObservationResult {
663 source_call_id: format!("search_{i}"),
664 content: r.clone(),
665 })
666 .collect(),
667 });
668 }
669 self.push_step(step);
670 }
671 ThreadItemDetails::Harness(h) => {
672 let msg = format!("harness: {:?}", h.event);
673 let mut step = Step::system(self.next_step_id, msg);
674 step.timestamp = Some(ts.to_string());
675 let mut extra = serde_json::Map::new();
676 if let Some(m) = &h.message {
677 let _ = extra.insert("harness_message".to_string(), Value::String(m.clone()));
678 }
679 if h.event == crate::HarnessEventKind::BackgroundSubprocessCompleted {
680 for (key, value) in [
681 ("task_id", h.task_id.as_ref()),
682 ("session_id", h.session_id.as_ref()),
683 ("exec_session_id", h.exec_session_id.as_ref()),
684 ("status", h.status.as_ref()),
685 ("transcript_path", h.transcript_path.as_ref()),
686 ("archive_path", h.archive_path.as_ref()),
687 ("error_category", h.error_category.as_ref()),
688 ] {
689 if let Some(value) = value {
690 let _ = extra.insert(key.to_string(), Value::String(value.clone()));
691 }
692 }
693 if let Some(exit_code) = h.exit_code {
694 let _ = extra.insert("exit_code".to_string(), Value::from(exit_code));
695 }
696 }
697 if !extra.is_empty() {
698 step.extra = Some(Value::Object(extra));
699 }
700 self.push_step(step);
701 }
702 ThreadItemDetails::Error(e) => {
703 let mut step = Step::system(self.next_step_id, &e.message);
704 step.timestamp = Some(ts.to_string());
705 self.push_step(step);
706 }
707 }
708 }
709
710 fn push_step(&mut self, step: Step) {
711 self.next_step_id = step.step_id + 1;
712 self.steps.push(step);
713 }
714
715 pub fn finish(self, override_metrics: Option<FinalMetrics>) -> Trajectory {
720 let final_metrics = override_metrics.unwrap_or_else(|| FinalMetrics {
721 total_prompt_tokens: Some(self.total_input_tokens),
722 total_completion_tokens: Some(self.total_output_tokens),
723 total_cached_tokens: if self.total_cached_tokens > 0 {
724 Some(self.total_cached_tokens)
725 } else {
726 None
727 },
728 total_cost_usd: None,
729 total_steps: Some(self.steps.len() as u64),
730 extra: Some(serde_json::json!({ "num_turns": self.num_turns })),
731 });
732
733 Trajectory {
734 schema_version: ATIF_SCHEMA_VERSION.to_string(),
735 session_id: self.session_id.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()),
736 agent: self.agent,
737 steps: self.steps,
738 notes: None,
739 final_metrics: Some(final_metrics),
740 extra: self.completed_at.map(|timestamp| serde_json::json!({"completed_at":timestamp})),
741 }
742 }
743
744 pub fn step_count(&self) -> usize {
746 self.steps.len()
747 }
748}
749
750impl crate::EventEmitter for AtifTrajectoryBuilder {
751 fn emit(&mut self, event: &ThreadEvent) {
752 self.process_event(event);
753 }
754}
755
756#[cfg(test)]
757mod tests {
758 use super::*;
759 use crate::{
760 AgentMessageItem, CompactionMode, CompactionTrigger, HarnessEventItem, HarnessEventKind, ItemCompletedEvent,
761 ThreadCompactBoundaryEvent, ThreadItem, ThreadStartedEvent, ToolInvocationItem, ToolOutputItem,
762 TurnCompletedEvent, TurnStartedEvent, Usage,
763 };
764
765 fn fixed_ts() -> DateTime<Utc> {
766 "2025-01-15T10:30:00Z".parse().unwrap()
767 }
768
769 #[test]
770 fn terminal_timestamps_survive_export_without_synthetic_steps() {
771 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
772 let timestamp = "2026-10-03T01:02:03Z";
773 let blocked = serde_json::from_value(
774 serde_json::json!({"type":"turn.blocked", "message":"Blocked", "completed_at":timestamp}),
775 )
776 .unwrap();
777 builder.process_event_at(&blocked, fixed_ts());
778 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();
779 builder.process_event_at(&completed, fixed_ts());
780 let trajectory = builder.finish(None);
781 assert_eq!(trajectory.steps.len(), 1);
782 assert_eq!(trajectory.steps[0].timestamp.as_deref(), Some(timestamp));
783 assert_eq!(trajectory.extra.as_ref().unwrap()["completed_at"], timestamp);
784 }
785
786 #[test]
787 fn explanation_metadata_and_public_rationale_preserve_recorded_time() {
788 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
789 let context = crate::ItemContext {
790 task_id: "task-1".into(),
791 turn_id: "turn-1".into(),
792 actor_id: "root".into(),
793 parent_actor_id: None,
794 timestamp: "2026-10-03T01:02:03Z".into(),
795 activity: None,
796 };
797 builder.process_event_at(
798 &ThreadEvent::ItemCompleted(ItemCompletedEvent {
799 item: ThreadItem {
800 id: "decision-1".into(),
801 context: Some(Box::new(context.clone())),
802 details: ThreadItemDetails::Decision(Box::new(crate::DecisionItem {
803 summary: "Reuse parser".into(),
804 rationale: "Preserve checks".into(),
805 alternatives: vec!["Replace parser".into()],
806 evidence_ids: vec!["read-1".into()],
807 })),
808 },
809 }),
810 fixed_ts(),
811 );
812 let trajectory = builder.finish(None);
813 assert_eq!(trajectory.steps.len(), 1);
814 let step = &trajectory.steps[0];
815 assert_eq!(step.timestamp.as_deref(), Some(context.timestamp.as_str()));
816 let extra = step.extra.as_ref().unwrap();
817 assert_eq!(extra["public_rationale"], "Preserve checks");
818 assert_eq!(extra["item_context"]["task_id"], "task-1");
819 assert_eq!(extra["evidence_ids"][0], "read-1");
820 }
821
822 #[test]
823 fn trajectory_round_trip() {
824 let trajectory = Trajectory {
825 schema_version: ATIF_SCHEMA_VERSION.to_string(),
826 session_id: "test-session".to_string(),
827 agent: AtifAgent::vtcode(),
828 steps: vec![Step::user(1, "hello")],
829 notes: None,
830 final_metrics: None,
831 extra: None,
832 };
833
834 let json = serde_json::to_string_pretty(&trajectory).unwrap();
835 let restored: Trajectory = serde_json::from_str(&json).unwrap();
836 assert_eq!(restored.schema_version, ATIF_SCHEMA_VERSION);
837 assert_eq!(restored.session_id, "test-session");
838 assert_eq!(restored.steps.len(), 1);
839 }
840
841 #[test]
842 fn builder_thread_started_sets_session_id() {
843 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
844 let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "thread-abc".to_string() });
845 builder.process_event_at(&event, fixed_ts());
846 let trajectory = builder.finish(None);
847 assert_eq!(trajectory.session_id, "thread-abc");
848 }
849
850 #[test]
851 fn builder_agent_message_step() {
852 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
853 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
854 item: ThreadItem {
855 context: None,
856 id: "msg-1".to_string(),
857 details: ThreadItemDetails::AgentMessage(AgentMessageItem { text: "Hello, world!".to_string() }),
858 },
859 });
860 builder.process_event_at(&event, fixed_ts());
861 let trajectory = builder.finish(None);
862
863 assert_eq!(trajectory.steps.len(), 1);
864 let step = &trajectory.steps[0];
865 assert_eq!(step.step_id, 1);
866 assert_eq!(step.source, StepSource::Agent);
867 assert_eq!(step.message.as_deref(), Some("Hello, world!"));
868 }
869
870 #[test]
871 fn background_completion_preserves_identity_in_atif_extra() {
872 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
873 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
874 item: ThreadItem {
875 context: None,
876 id: "background-completion:task-1:exec-1".to_string(),
877 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
878 event: HarnessEventKind::BackgroundSubprocessCompleted,
879 message: Some("completed successfully".to_string()),
880 command: None,
881 path: None,
882 exit_code: Some(0),
883 attempt: None,
884 error_category: Some("background_subprocess".to_string()),
885 duration_ms: None,
886 task_id: Some("task-1".to_string()),
887 session_id: Some("session-1".to_string()),
888 exec_session_id: Some("exec-1".to_string()),
889 status: Some("stopped".to_string()),
890 transcript_path: Some("/tmp/transcript.jsonl".to_string()),
891 archive_path: Some("/tmp/archive.json".to_string()),
892 })),
893 },
894 });
895
896 builder.process_event_at(&event, fixed_ts());
897 let trajectory = builder.finish(None);
898 let step = trajectory.steps.first().expect("background completion step");
899 let extra = step.extra.as_ref().expect("background completion metadata");
900 assert_eq!(step.message.as_deref(), Some("harness: BackgroundSubprocessCompleted"));
901 assert_eq!(extra["harness_message"], "completed successfully");
902 assert_eq!(extra["task_id"], "task-1");
903 assert_eq!(extra["session_id"], "session-1");
904 assert_eq!(extra["exec_session_id"], "exec-1");
905 assert_eq!(extra["status"], "stopped");
906 assert_eq!(extra["exit_code"], 0);
907 assert_eq!(extra["transcript_path"], "/tmp/transcript.jsonl");
908 assert_eq!(extra["archive_path"], "/tmp/archive.json");
909 assert_eq!(extra["error_category"], "background_subprocess");
910 }
911
912 #[test]
913 fn builder_tool_invocation_with_output() {
914 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
915 let ts = fixed_ts();
916
917 let inv_event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
919 item: ThreadItem {
920 context: None,
921 id: "tool_1".to_string(),
922 details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
923 tool_name: "read_file".to_string(),
924 arguments: Some(serde_json::json!({"path": "README.md"})),
925 tool_call_id: Some("tc_0".to_string()),
926 status: ToolCallStatus::Completed,
927 outcome: None,
928 })),
929 },
930 });
931 builder.process_event_at(&inv_event, ts);
932
933 let out_event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
935 item: ThreadItem {
936 context: None,
937 id: "tool_1:output".to_string(),
938 details: ThreadItemDetails::ToolOutput(Box::new(ToolOutputItem {
939 call_id: "tool_1".to_string(),
940 tool_call_id: Some("tc_0".to_string()),
941 spool_path: None,
942 output: "file contents here".to_string(),
943 exit_code: Some(0),
944 status: ToolCallStatus::Completed,
945 })),
946 },
947 });
948 builder.process_event_at(&out_event, ts);
949
950 let trajectory = builder.finish(None);
951 assert_eq!(trajectory.steps.len(), 1);
953 let step = &trajectory.steps[0];
954 assert_eq!(step.source, StepSource::Agent);
955
956 let calls = step.tool_calls.as_ref().unwrap();
957 assert_eq!(calls.len(), 1);
958 assert_eq!(calls[0].function_name, "read_file");
959 assert_eq!(calls[0].tool_call_id, "tc_0");
960
961 let obs = step.observation.as_ref().unwrap();
962 assert_eq!(obs.results.len(), 1);
963 assert_eq!(obs.results[0].content, "file contents here");
964 }
965
966 #[test]
967 fn builder_turn_completed_accumulates_metrics() {
968 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
969 let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
970 completed_at: None,
971 usage: Usage {
972 input_tokens: 500,
973 cached_input_tokens: 100,
974 cache_creation_tokens: 0,
975 output_tokens: 200,
976 },
977 in_progress_exec_sessions: Vec::new(),
978 });
979 builder.process_event_at(&event, fixed_ts());
980
981 let trajectory = builder.finish(None);
982 let fm = trajectory.final_metrics.as_ref().unwrap();
983 assert_eq!(fm.total_prompt_tokens, Some(500));
984 assert_eq!(fm.total_completion_tokens, Some(200));
985 assert_eq!(fm.total_cached_tokens, Some(100));
986 }
987
988 #[test]
989 fn builder_turn_completed_preserves_in_progress_sessions_for_resume() {
990 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
991 let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
992 completed_at: None,
993 usage: Usage::default(),
994 in_progress_exec_sessions: vec!["run-7".to_string()],
995 });
996 builder.process_event_at(&event, fixed_ts());
997
998 let trajectory = builder.finish(None);
999 let step = trajectory.steps.last().expect("turn_completed step");
1000 let extra = step.extra.clone().expect("extra carries resume ids");
1001 assert_eq!(extra["in_progress_exec_sessions"], serde_json::json!(["run-7"]));
1002
1003 let mut empty_builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1005 let empty = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1006 completed_at: None,
1007 usage: Usage::default(),
1008 in_progress_exec_sessions: Vec::new(),
1009 });
1010 empty_builder.process_event_at(&empty, fixed_ts());
1011 let empty_trajectory = empty_builder.finish(None);
1012 assert!(empty_trajectory.steps.last().expect("step").extra.is_none());
1013 }
1014
1015 #[test]
1016 fn step_metrics_from_usage() {
1017 let usage = Usage {
1018 input_tokens: 1000,
1019 cached_input_tokens: 200,
1020 cache_creation_tokens: 50,
1021 output_tokens: 300,
1022 };
1023 let metrics = StepMetrics::from_usage(&usage);
1024 assert_eq!(metrics.prompt_tokens, Some(1000));
1025 assert_eq!(metrics.completion_tokens, Some(300));
1026 assert_eq!(metrics.cached_tokens, Some(200));
1027 assert!(metrics.extra.is_some());
1028 }
1029
1030 #[test]
1031 fn builder_implements_event_emitter() {
1032 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1033 let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "t-1".to_string() });
1034 crate::EventEmitter::emit(&mut builder, &event);
1036 assert_eq!(builder.step_count(), 0); }
1038
1039 #[test]
1040 fn skips_lifecycle_events() {
1041 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1042 builder.process_event(&ThreadEvent::TurnStarted(TurnStartedEvent::default()));
1043 assert_eq!(builder.step_count(), 0);
1044 }
1045
1046 #[test]
1047 fn compact_boundary_with_segment_metadata_preserves_atif_export() {
1048 let mut builder = AtifTrajectoryBuilder::new(AtifAgent::vtcode());
1049 let event = ThreadEvent::ThreadCompactBoundary(Box::new(ThreadCompactBoundaryEvent {
1050 thread_id: "thread-1".to_string(),
1051 trigger: CompactionTrigger::Auto,
1052 mode: CompactionMode::Local,
1053 original_message_count: 12,
1054 compacted_message_count: 5,
1055 history_artifact_path: None,
1056 previous_segment_id: Some("segment-0001".to_string()),
1057 new_segment_id: Some("segment-0002".to_string()),
1058 previous_prefix_hash: Some("prefix-before".to_string()),
1059 new_prefix_hash: Some("prefix-after".to_string()),
1060 previous_catalog_hash: Some("catalog-before".to_string()),
1061 new_catalog_hash: Some("catalog-after".to_string()),
1062 }));
1063
1064 builder.process_event_at(&event, fixed_ts());
1065 let trajectory = builder.finish(None);
1066
1067 assert_eq!(trajectory.steps.len(), 1);
1068 assert_eq!(trajectory.steps[0].source, StepSource::System);
1069 assert_eq!(
1070 trajectory.steps[0].message.as_deref(),
1071 Some("context_compaction: 12 messages -> 5 messages (auto)")
1072 );
1073 }
1074}