1#![allow(
2 missing_docs,
3 dead_code,
4 unused_imports,
5 reason = "Intentional compatibility, platform, or test-only suppression."
6)]
7use serde::{Deserialize, Serialize};
21use serde_json::Value;
22
23pub mod atif;
24pub mod trace;
25
26pub const EVENT_SCHEMA_VERSION: &str = "0.13.0";
28
29#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
32#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
33pub struct VersionedThreadEvent {
34 schema_version: String,
36 event: ThreadEvent,
38}
39
40impl VersionedThreadEvent {
41 pub fn new(event: ThreadEvent) -> Self {
44 Self {
45 schema_version: EVENT_SCHEMA_VERSION.to_string(),
46 event,
47 }
48 }
49
50 pub fn into_event(self) -> ThreadEvent {
52 self.event
53 }
54}
55
56impl From<ThreadEvent> for VersionedThreadEvent {
57 fn from(event: ThreadEvent) -> Self {
58 Self::new(event)
59 }
60}
61
62pub trait EventEmitter {
64 fn emit(&mut self, event: &ThreadEvent);
66}
67
68impl<F> EventEmitter for F
69where
70 F: FnMut(&ThreadEvent),
71{
72 fn emit(&mut self, event: &ThreadEvent) {
73 self(event);
74 }
75}
76
77#[cfg(feature = "serde-json")]
79pub(crate) mod json {
80 use super::{ThreadEvent, VersionedThreadEvent};
81
82 pub fn to_value(event: &ThreadEvent) -> serde_json::Result<serde_json::Value> {
84 serde_json::to_value(event)
85 }
86
87 pub(crate) fn to_string(event: &ThreadEvent) -> serde_json::Result<String> {
89 serde_json::to_string(event)
90 }
91
92 pub fn from_str(payload: &str) -> serde_json::Result<ThreadEvent> {
94 serde_json::from_str(payload)
95 }
96
97 pub(crate) fn versioned_to_string(event: &ThreadEvent) -> serde_json::Result<String> {
99 serde_json::to_string(&VersionedThreadEvent::new(event.clone()))
100 }
101
102 pub(crate) fn versioned_from_str(payload: &str) -> serde_json::Result<VersionedThreadEvent> {
104 serde_json::from_str(payload)
105 }
106}
107
108#[cfg(feature = "telemetry-log")]
109mod log_support {
110 use log::Level;
111
112 use super::{EventEmitter, ThreadEvent, json};
113
114 #[derive(Debug, Clone)]
116 pub struct LogEmitter {
117 level: Level,
118 }
119
120 impl LogEmitter {
121 pub fn new(level: Level) -> Self {
123 Self { level }
124 }
125 }
126
127 impl Default for LogEmitter {
128 fn default() -> Self {
129 Self { level: Level::Info }
130 }
131 }
132
133 impl EventEmitter for LogEmitter {
134 fn emit(&mut self, event: &ThreadEvent) {
135 if log::log_enabled!(self.level) {
136 match json::to_string(event) {
137 Ok(serialized) => log::log!(self.level, "{serialized}"),
138 Err(err) => log::log!(self.level, "failed to serialize vtcode exec event for logging: {err}"),
139 }
140 }
141 }
142 }
143
144 pub use LogEmitter as PublicLogEmitter;
145}
146
147#[cfg(feature = "telemetry-log")]
148pub use log_support::PublicLogEmitter as LogEmitter;
149
150#[cfg(feature = "telemetry-tracing")]
151mod tracing_support {
152 use tracing::Level;
153
154 use super::{EVENT_SCHEMA_VERSION, EventEmitter, ThreadEvent, VersionedThreadEvent};
155
156 #[derive(Debug, Clone)]
158 pub struct TracingEmitter {
159 level: Level,
160 }
161
162 impl TracingEmitter {
163 pub fn new(level: Level) -> Self {
165 Self { level }
166 }
167 }
168
169 impl Default for TracingEmitter {
170 fn default() -> Self {
171 Self { level: Level::INFO }
172 }
173 }
174
175 impl EventEmitter for TracingEmitter {
176 fn emit(&mut self, event: &ThreadEvent) {
177 match self.level {
178 Level::TRACE => tracing::event!(
179 target: "vtcode_exec_events",
180 Level::TRACE,
181 schema_version = EVENT_SCHEMA_VERSION,
182 event = ?VersionedThreadEvent::new(event.clone()),
183 "vtcode_exec_event"
184 ),
185 Level::DEBUG => tracing::event!(
186 target: "vtcode_exec_events",
187 Level::DEBUG,
188 schema_version = EVENT_SCHEMA_VERSION,
189 event = ?VersionedThreadEvent::new(event.clone()),
190 "vtcode_exec_event"
191 ),
192 Level::INFO => tracing::event!(
193 target: "vtcode_exec_events",
194 Level::INFO,
195 schema_version = EVENT_SCHEMA_VERSION,
196 event = ?VersionedThreadEvent::new(event.clone()),
197 "vtcode_exec_event"
198 ),
199 Level::WARN => tracing::event!(
200 target: "vtcode_exec_events",
201 Level::WARN,
202 schema_version = EVENT_SCHEMA_VERSION,
203 event = ?VersionedThreadEvent::new(event.clone()),
204 "vtcode_exec_event"
205 ),
206 Level::ERROR => tracing::event!(
207 target: "vtcode_exec_events",
208 Level::ERROR,
209 schema_version = EVENT_SCHEMA_VERSION,
210 event = ?VersionedThreadEvent::new(event.clone()),
211 "vtcode_exec_event"
212 ),
213 }
214 }
215 }
216
217 pub use TracingEmitter as PublicTracingEmitter;
218}
219
220#[cfg(feature = "telemetry-tracing")]
221pub use tracing_support::PublicTracingEmitter as TracingEmitter;
222
223#[cfg(feature = "telemetry-otel")]
224mod otel_support {
225 use opentelemetry::KeyValue;
226 use opentelemetry::trace::{Span, Status, Tracer};
227
228 use super::{EventEmitter, ThreadEvent, ThreadItemDetails};
229
230 pub struct OtelEmitter<T: Tracer> {
246 tracer: T,
247 }
248
249 impl<T: Tracer> OtelEmitter<T> {
250 pub fn new(tracer: T) -> Self {
251 Self { tracer }
252 }
253 }
254
255 impl<T: Tracer> EventEmitter for OtelEmitter<T> {
256 fn emit(&mut self, event: &ThreadEvent) {
257 let span_name = match event {
258 ThreadEvent::ThreadStarted(_) => "thread.started",
259 ThreadEvent::ThreadCompleted(_) => "thread.completed",
260 ThreadEvent::ContextReset(_) => "context.reset",
261 ThreadEvent::TurnStarted(_) => "turn.started",
262 ThreadEvent::TurnCompleted(_) => "turn.completed",
263 ThreadEvent::TurnFailed(_) => "turn.failed",
264 ThreadEvent::ItemStarted(_) => "item.started",
265 ThreadEvent::ItemUpdated(_) => "item.updated",
266 ThreadEvent::ItemCompleted(_) => "item.completed",
267 ThreadEvent::Error(_) => "error",
268 _ => "event",
269 };
270
271 let mut span = self.tracer.start(span_name);
272
273 match event {
274 ThreadEvent::ThreadStarted(e) => {
275 span.set_attribute(KeyValue::new("thread_id", e.thread_id.clone()));
276 }
277 ThreadEvent::ThreadCompleted(e) => {
278 if let Some(ref cost) = e.total_cost_usd {
279 span.set_attribute(KeyValue::new("total_cost_usd", cost.as_f64().unwrap_or(0.0)));
280 }
281 span.set_attribute(KeyValue::new(
282 "input_tokens",
283 i64::try_from(e.usage.input_tokens).unwrap_or(i64::MAX),
284 ));
285 span.set_attribute(KeyValue::new(
286 "output_tokens",
287 i64::try_from(e.usage.output_tokens).unwrap_or(i64::MAX),
288 ));
289 span.set_attribute(KeyValue::new("completion_subtype", e.subtype.as_str().to_string()));
290 }
291 ThreadEvent::ContextReset(e) => {
292 span.set_attribute(KeyValue::new("thread_id", e.thread_id.clone()));
293 span.set_attribute(KeyValue::new("turn_id", e.turn_id.clone()));
294 span.set_attribute(KeyValue::new("plan_preserved", e.plan_preserved));
295 span.set_attribute(KeyValue::new(
296 "previous_context_usage_percent",
297 e.previous_context_usage_percent as i64,
298 ));
299 span.set_attribute(KeyValue::new("tool_budget_reset", e.tool_budget_reset));
300 }
301 ThreadEvent::TurnCompleted(e) => {
302 span.set_attribute(KeyValue::new(
303 "turn_input_tokens",
304 i64::try_from(e.usage.input_tokens).unwrap_or(i64::MAX),
305 ));
306 span.set_attribute(KeyValue::new(
307 "turn_output_tokens",
308 i64::try_from(e.usage.output_tokens).unwrap_or(i64::MAX),
309 ));
310 }
311 ThreadEvent::ItemCompleted(e) => {
312 if let ThreadItemDetails::Harness(harness) = &e.item.details {
313 span.set_attribute(KeyValue::new("harness_event", format!("{:?}", harness.event)));
314 if let Some(ref msg) = harness.message {
315 span.set_attribute(KeyValue::new("harness_message", msg.clone()));
316 }
317 if let Some(ref path) = harness.path {
318 span.set_attribute(KeyValue::new("harness_path", path.clone()));
319 }
320 if let Some(dur) = harness.duration_ms {
321 span.set_attribute(KeyValue::new("duration_ms", i64::try_from(dur).unwrap_or(i64::MAX)));
322 }
323 let mut event_attrs = vec![KeyValue::new("event_kind", format!("{:?}", harness.event))];
324 if let Some(ref msg) = harness.message {
325 event_attrs.push(KeyValue::new("message", msg.clone()));
326 }
327 span.add_event("harness_event", event_attrs);
328 }
329 }
330 ThreadEvent::Error(e) => {
331 span.set_status(Status::Error { description: e.message.clone().into() });
332 span.set_attribute(KeyValue::new("error_message", e.message.clone()));
333 }
334 _ => {}
335 }
336
337 span.end();
338 }
339 }
340
341 pub use OtelEmitter as PublicOtelEmitter;
342}
343
344#[cfg(feature = "telemetry-otel")]
345pub use otel_support::PublicOtelEmitter as OtelEmitter;
346
347#[cfg(feature = "schema-export")]
348pub mod schema {
349 use schemars::{Schema, schema_for};
350
351 use super::{ThreadEvent, VersionedThreadEvent};
352
353 pub fn thread_event_schema() -> Schema {
355 schema_for!(ThreadEvent)
356 }
357
358 pub fn versioned_thread_event_schema() -> Schema {
360 schema_for!(VersionedThreadEvent)
361 }
362}
363
364#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
366#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
367#[serde(tag = "type")]
368pub enum ThreadEvent {
369 #[serde(rename = "thread.started")]
371 ThreadStarted(ThreadStartedEvent),
372 #[serde(rename = "thread.completed")]
374 ThreadCompleted(ThreadCompletedEvent),
375 #[serde(rename = "thread.compact_boundary")]
377 ThreadCompactBoundary(ThreadCompactBoundaryEvent),
378 #[serde(rename = "context.reset")]
380 ContextReset(ContextResetEvent),
381 #[serde(rename = "turn.started")]
383 TurnStarted(TurnStartedEvent),
384 #[serde(rename = "turn.completed")]
386 TurnCompleted(TurnCompletedEvent),
387 #[serde(rename = "turn.failed")]
389 TurnFailed(TurnFailedEvent),
390 #[serde(rename = "turn.blocked")]
394 TurnBlocked(TurnBlockedEvent),
395 #[serde(rename = "item.started")]
397 ItemStarted(ItemStartedEvent),
398 #[serde(rename = "item.updated")]
400 ItemUpdated(ItemUpdatedEvent),
401 #[serde(rename = "item.completed")]
403 ItemCompleted(ItemCompletedEvent),
404 #[serde(rename = "permission.requested")]
406 PermissionRequested(PermissionRequestedEvent),
407 #[serde(rename = "permission.resolved")]
409 PermissionResolved(PermissionResolvedEvent),
410 #[serde(rename = "interjected")]
412 Interjected(InterjectedEvent),
413 #[serde(rename = "plan.delta")]
415 PlanDelta(PlanDeltaEvent),
416 #[serde(rename = "plan.approval.requested")]
418 PlanApprovalRequested(PlanApprovalRequestedEvent),
419 #[serde(rename = "plan.approval.resolved")]
421 PlanApprovalResolved(PlanApprovalResolvedEvent),
422 #[serde(rename = "error")]
424 Error(ThreadErrorEvent),
425 #[serde(other)]
428 Unknown,
429}
430
431#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
432#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
433pub struct ThreadStartedEvent {
434 pub thread_id: String,
436}
437
438#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
439#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
440#[serde(rename_all = "snake_case")]
441pub enum ThreadCompletionSubtype {
442 Success,
443 ErrorMaxTurns,
444 ErrorMaxBudgetUsd,
445 ErrorDuringExecution,
446 Cancelled,
447 #[serde(other)]
449 Unknown,
450}
451
452impl ThreadCompletionSubtype {
453 pub const fn as_str(&self) -> &'static str {
454 match self {
455 Self::Success => "success",
456 Self::ErrorMaxTurns => "error_max_turns",
457 Self::ErrorMaxBudgetUsd => "error_max_budget_usd",
458 Self::ErrorDuringExecution => "error_during_execution",
459 Self::Cancelled => "cancelled",
460 Self::Unknown => "unknown",
461 }
462 }
463
464 pub const fn is_success(self) -> bool {
465 matches!(self, Self::Success)
466 }
467}
468
469#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
470#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
471#[serde(rename_all = "snake_case")]
472pub enum CompactionTrigger {
473 Manual,
474 Auto,
475 Recovery,
476 ModelSwitch,
479 #[serde(other)]
481 Unknown,
482}
483
484impl CompactionTrigger {
485 pub const fn as_str(self) -> &'static str {
486 match self {
487 Self::Manual => "manual",
488 Self::Auto => "auto",
489 Self::Recovery => "recovery",
490 Self::ModelSwitch => "model_switch",
491 Self::Unknown => "unknown",
492 }
493 }
494}
495
496#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
497#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
498#[serde(rename_all = "snake_case")]
499pub enum CompactionMode {
500 Provider,
501 Local,
502 #[serde(other)]
504 Unknown,
505}
506
507impl CompactionMode {
508 pub const fn as_str(self) -> &'static str {
509 match self {
510 Self::Provider => "provider",
511 Self::Local => "local",
512 Self::Unknown => "unknown",
513 }
514 }
515}
516
517#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
518#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
519pub struct ThreadCompletedEvent {
520 pub thread_id: String,
522 pub session_id: String,
524 pub subtype: ThreadCompletionSubtype,
526 pub outcome_code: String,
528 #[serde(skip_serializing_if = "Option::is_none")]
530 pub result: Option<String>,
531 #[serde(skip_serializing_if = "Option::is_none")]
533 pub stop_reason: Option<String>,
534 pub usage: Usage,
536 #[serde(skip_serializing_if = "Option::is_none")]
538 pub total_cost_usd: Option<serde_json::Number>,
539 pub num_turns: usize,
541}
542
543#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
544#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
545pub struct ThreadCompactBoundaryEvent {
546 pub thread_id: String,
548 pub trigger: CompactionTrigger,
550 pub mode: CompactionMode,
552 pub original_message_count: usize,
554 pub compacted_message_count: usize,
556 #[serde(skip_serializing_if = "Option::is_none")]
558 pub history_artifact_path: Option<String>,
559 #[serde(skip_serializing_if = "Option::is_none")]
561 pub previous_segment_id: Option<String>,
562 #[serde(skip_serializing_if = "Option::is_none")]
564 pub new_segment_id: Option<String>,
565 #[serde(skip_serializing_if = "Option::is_none")]
567 pub previous_prefix_hash: Option<String>,
568 #[serde(skip_serializing_if = "Option::is_none")]
570 pub new_prefix_hash: Option<String>,
571 #[serde(skip_serializing_if = "Option::is_none")]
573 pub previous_catalog_hash: Option<String>,
574 #[serde(skip_serializing_if = "Option::is_none")]
576 pub new_catalog_hash: Option<String>,
577}
578
579#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
580#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
581#[serde(rename_all = "snake_case")]
582pub enum ContextResetTrigger {
583 PlanApproval,
585 #[serde(other)]
587 Unknown,
588}
589
590#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
591#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
592pub struct ContextResetEvent {
593 pub thread_id: String,
595 pub turn_id: String,
597 pub trigger: ContextResetTrigger,
599 pub plan_preserved: bool,
601 pub previous_context_usage_percent: u8,
603 pub tool_budget_reset: bool,
605}
606
607#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
608#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
609pub struct TurnStartedEvent {
610 #[serde(skip_serializing_if = "Option::is_none")]
614 token_breakdown: Option<TokenBreakdown>,
615}
616
617#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
619#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
620pub struct TokenBreakdown {
621 system_prompt_tokens: u64,
623 tool_schema_tokens: u64,
625 instruction_file_tokens: u64,
627 message_history_tokens: u64,
629 cache_read_tokens: u64,
631 cache_write_tokens: u64,
633 cache_miss_tokens: u64,
635 #[serde(skip_serializing_if = "Option::is_none")]
637 subagent_bootstrap_tokens: Option<u64>,
638}
639
640#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
641#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
642pub struct TurnCompletedEvent {
643 pub usage: Usage,
645}
646
647#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
648#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
649pub struct TurnFailedEvent {
650 pub message: String,
652 #[serde(skip_serializing_if = "Option::is_none")]
654 pub usage: Option<Usage>,
655}
656
657#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
658#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
659pub struct TurnBlockedEvent {
660 pub message: String,
662 #[serde(skip_serializing_if = "Option::is_none")]
664 pub last_tool: Option<String>,
665 #[serde(default)]
667 pub blocked_streak: usize,
668 #[serde(default)]
670 pub blocked_total: usize,
671 #[serde(default)]
673 pub consecutive_cap: usize,
674 #[serde(default)]
676 pub total_cap: usize,
677 #[serde(default)]
679 pub recovery_active: bool,
680 #[serde(skip_serializing_if = "Option::is_none")]
682 pub usage: Option<Usage>,
683}
684
685#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
686#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
687pub struct ThreadErrorEvent {
688 pub message: String,
690}
691
692#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
693#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
694pub struct Usage {
695 pub input_tokens: u64,
697 pub cached_input_tokens: u64,
699 pub cache_creation_tokens: u64,
701 pub output_tokens: u64,
703}
704
705impl Usage {
706 #[must_use]
711 fn uncached_input_tokens(&self) -> u64 {
712 self.input_tokens
713 .saturating_sub(self.cached_input_tokens)
714 .saturating_sub(self.cache_creation_tokens)
715 }
716
717 #[must_use]
720 pub fn cache_hit_rate(&self) -> Option<f64> {
721 if self.input_tokens == 0 {
722 return None;
723 }
724 Some(self.cached_input_tokens as f64 / self.input_tokens as f64)
725 }
726
727 #[must_use]
729 pub fn cache_summary(&self) -> String {
730 let total_input = self.input_tokens;
731 if total_input == 0 {
732 return "No input tokens recorded.".to_string();
733 }
734
735 let cached = self.cached_input_tokens;
736 let creation = self.cache_creation_tokens;
737 let uncached = self.uncached_input_tokens();
738 let rate = cached as f64 / total_input as f64 * 100.0;
739 format!(
740 "Cache: {cached} cached / {total_input} total input ({rate:.1}% hit rate), \
741 {creation} cache-creation, {uncached} uncached"
742 )
743 }
744
745 pub fn add(&mut self, other: &Usage) {
747 self.input_tokens = self.input_tokens.saturating_add(other.input_tokens);
748 self.cached_input_tokens = self.cached_input_tokens.saturating_add(other.cached_input_tokens);
749 self.cache_creation_tokens = self.cache_creation_tokens.saturating_add(other.cache_creation_tokens);
750 self.output_tokens = self.output_tokens.saturating_add(other.output_tokens);
751 }
752}
753
754#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
755#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
756pub struct ItemCompletedEvent {
757 pub item: ThreadItem,
759}
760
761#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
762#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
763pub struct ItemStartedEvent {
764 pub item: ThreadItem,
766}
767
768#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
769#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
770pub struct ItemUpdatedEvent {
771 pub item: ThreadItem,
773}
774
775#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
776#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
777pub struct PlanDeltaEvent {
778 pub thread_id: String,
780 pub turn_id: String,
782 pub item_id: String,
784 pub delta: String,
786}
787
788#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
789#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
790pub struct PlanApprovalRequestedEvent {
791 pub thread_id: String,
793 pub turn_id: String,
795 #[serde(skip_serializing_if = "Option::is_none")]
797 pub plan_file: Option<String>,
798}
799
800#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
801#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
802#[serde(rename_all = "snake_case")]
803pub enum PlanApprovalDecision {
804 Execute,
806 AutoAccept,
808 FreshContext,
810 Revise,
812 Cancel,
814 SwitchBuild,
816 SwitchAuto,
818 #[serde(other)]
820 Unknown,
821}
822
823#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
824#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
825pub struct PlanApprovalResolvedEvent {
826 pub thread_id: String,
828 pub turn_id: String,
830 pub decision: PlanApprovalDecision,
832 pub automatic: bool,
834}
835
836#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
837#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
838pub struct ThreadItem {
839 pub id: String,
841 #[serde(flatten)]
843 pub details: ThreadItemDetails,
844}
845
846#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
847#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
848#[serde(tag = "type", rename_all = "snake_case")]
849pub enum ThreadItemDetails {
850 AgentMessage(AgentMessageItem),
852 Plan(PlanItem),
854 Reasoning(ReasoningItem),
856 CommandExecution(Box<CommandExecutionItem>),
858 ToolInvocation(ToolInvocationItem),
860 ToolOutput(ToolOutputItem),
862 FileChange(Box<FileChangeItem>),
864 McpToolCall(McpToolCallItem),
866 WebSearch(WebSearchItem),
868 Harness(HarnessEventItem),
870 Error(ErrorItem),
872}
873
874#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
875#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
876pub struct AgentMessageItem {
877 pub text: String,
879}
880
881#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
882#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
883pub struct PlanItem {
884 pub text: String,
886}
887
888#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
889#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
890pub struct ReasoningItem {
891 pub text: String,
893 #[serde(skip_serializing_if = "Option::is_none")]
896 pub stage: Option<String>,
897}
898
899#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
900#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
901#[serde(rename_all = "snake_case")]
902pub enum CommandExecutionStatus {
903 #[default]
905 Completed,
906 Failed,
908 InProgress,
910}
911
912#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
913#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
914pub struct CommandExecutionItem {
915 pub command: String,
917 #[serde(skip_serializing_if = "Option::is_none")]
919 pub arguments: Option<Value>,
920 #[serde(default)]
922 pub aggregated_output: String,
923 #[serde(skip_serializing_if = "Option::is_none")]
925 pub exit_code: Option<i32>,
926 pub status: CommandExecutionStatus,
928}
929
930#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
931#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
932#[serde(rename_all = "snake_case")]
933pub enum ToolCallStatus {
934 #[default]
936 Completed,
937 Failed,
939 InProgress,
941}
942
943#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
951#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
952#[serde(rename_all = "snake_case")]
953pub enum ToolOutcome {
954 #[default]
956 Success,
957 Error,
959 PermissionRejected,
961 PermissionCancelled,
963 Followup,
965 HookDenied,
967 InvalidTool,
969 Cancelled,
971}
972
973impl ToolOutcome {
974 #[must_use]
975 pub const fn is_terminal(self) -> bool {
976 !matches!(self, Self::Followup)
977 }
978}
979
980#[must_use]
987#[allow(
988 clippy::unreachable,
989 reason = "Intentional compatibility, platform, or test-only suppression."
990)]
991pub fn tool_outcome_from_status(status: &ToolCallStatus) -> ToolOutcome {
992 match status {
993 ToolCallStatus::Completed => ToolOutcome::Success,
994 ToolCallStatus::Failed => ToolOutcome::Error,
995 ToolCallStatus::InProgress => unreachable!("InProgress status passed to completion event"),
996 }
997}
998
999#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1000#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1001pub struct ToolInvocationItem {
1002 pub tool_name: String,
1004 #[serde(skip_serializing_if = "Option::is_none")]
1006 pub arguments: Option<Value>,
1007 #[serde(skip_serializing_if = "Option::is_none")]
1009 pub tool_call_id: Option<String>,
1010 pub status: ToolCallStatus,
1012 #[serde(skip_serializing_if = "Option::is_none")]
1014 pub outcome: Option<ToolOutcome>,
1015}
1016
1017#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1018#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1019pub struct ToolOutputItem {
1020 pub call_id: String,
1022 #[serde(skip_serializing_if = "Option::is_none")]
1024 pub tool_call_id: Option<String>,
1025 #[serde(skip_serializing_if = "Option::is_none")]
1027 pub spool_path: Option<String>,
1028 #[serde(default)]
1030 pub output: String,
1031 #[serde(skip_serializing_if = "Option::is_none")]
1033 pub exit_code: Option<i32>,
1034 pub status: ToolCallStatus,
1036}
1037
1038#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1039#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1040pub struct FileChangeItem {
1041 pub changes: Vec<FileUpdateChange>,
1043 pub status: PatchApplyStatus,
1045}
1046
1047#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1048#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1049pub struct FileUpdateChange {
1050 pub path: String,
1052 pub kind: PatchChangeKind,
1054}
1055
1056#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1057#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1058#[serde(rename_all = "snake_case")]
1059pub enum PatchApplyStatus {
1060 Completed,
1062 Failed,
1064}
1065
1066#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1067#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1068#[serde(rename_all = "snake_case")]
1069pub enum PatchChangeKind {
1070 Add,
1072 Delete,
1074 Update,
1076}
1077
1078#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1079#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1080pub struct McpToolCallItem {
1081 pub tool_name: String,
1083 #[serde(skip_serializing_if = "Option::is_none")]
1085 pub arguments: Option<Value>,
1086 #[serde(skip_serializing_if = "Option::is_none")]
1088 pub result: Option<String>,
1089 #[serde(skip_serializing_if = "Option::is_none")]
1091 pub status: Option<McpToolCallStatus>,
1092}
1093
1094#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1095#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1096#[serde(rename_all = "snake_case")]
1097pub enum McpToolCallStatus {
1098 Started,
1100 Completed,
1102 Failed,
1104}
1105
1106#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1107#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1108pub struct WebSearchItem {
1109 pub query: String,
1111 #[serde(skip_serializing_if = "Option::is_none")]
1113 pub provider: Option<String>,
1114 #[serde(skip_serializing_if = "Option::is_none")]
1116 pub results: Option<Vec<String>>,
1117}
1118
1119#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1120#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1121#[serde(rename_all = "snake_case")]
1122pub enum HarnessEventKind {
1123 PlanningStarted,
1124 PlanningCompleted,
1125 ContinuationStarted,
1126 ContinuationSkipped,
1127 TurnBlocked,
1130 BlockedRecoveryStarted,
1132 BlockedRecoveryFinished,
1134 BlockedHandoffWritten,
1135 BlockedHandoffResolved,
1138 EvaluationStarted,
1139 EvaluationPassed,
1140 EvaluationFailed,
1141 RevisionStarted,
1142 EscalationTriggered,
1143 EscalationBypassed,
1144 VerificationStarted,
1145 VerificationPassed,
1146 VerificationFailed,
1147 ErrorRecovered,
1149 ToolRetryAttempted,
1151 ToolLatencyRecorded,
1153 SnapshotCreated,
1155 SnapshotRestored,
1157}
1158
1159#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1160#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1161#[serde(rename_all = "snake_case")]
1162pub enum PermissionDecision {
1163 Allow,
1164 Deny,
1165 Cancelled,
1166 Followup,
1167}
1168
1169#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1170#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1171pub struct PermissionRequestedEvent {
1172 pub tool_name: String,
1174}
1175
1176#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1177#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1178pub struct PermissionResolvedEvent {
1179 pub tool_name: String,
1181 pub decision: PermissionDecision,
1183 pub wait_ms: u64,
1185}
1186
1187#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1188#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1189#[serde(rename_all = "snake_case")]
1190pub enum InterjectionSource {
1191 Direct,
1192 Queue,
1193}
1194
1195#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1196#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1197#[serde(rename_all = "snake_case")]
1198pub enum RedirectKind {
1199 Interjection,
1200}
1201
1202#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1203#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1204pub struct InterjectedEvent {
1205 pub source: InterjectionSource,
1207 pub image_count: u32,
1209 pub redirect_kind: RedirectKind,
1212}
1213
1214#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1215#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1216pub struct HarnessEventItem {
1217 pub event: HarnessEventKind,
1219 #[serde(skip_serializing_if = "Option::is_none")]
1221 pub message: Option<String>,
1222 #[serde(skip_serializing_if = "Option::is_none")]
1224 pub command: Option<String>,
1225 #[serde(skip_serializing_if = "Option::is_none")]
1227 pub path: Option<String>,
1228 #[serde(skip_serializing_if = "Option::is_none")]
1230 pub exit_code: Option<i32>,
1231 #[serde(skip_serializing_if = "Option::is_none")]
1233 pub attempt: Option<u32>,
1234 #[serde(skip_serializing_if = "Option::is_none")]
1236 pub error_category: Option<String>,
1237 #[serde(skip_serializing_if = "Option::is_none")]
1239 pub duration_ms: Option<u64>,
1240}
1241
1242#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1243#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1244pub struct ErrorItem {
1245 pub message: String,
1247}
1248
1249#[cfg(test)]
1250mod tests {
1251 use super::*;
1252 use std::error::Error;
1253
1254 #[test]
1255 fn thread_event_round_trip() -> Result<(), Box<dyn Error>> {
1256 let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1257 usage: Usage {
1258 input_tokens: 1,
1259 cached_input_tokens: 2,
1260 cache_creation_tokens: 0,
1261 output_tokens: 3,
1262 },
1263 });
1264
1265 let json = serde_json::to_string(&event)?;
1266 let restored: ThreadEvent = serde_json::from_str(&json)?;
1267
1268 assert_eq!(restored, event);
1269 Ok(())
1270 }
1271
1272 #[test]
1273 fn turn_blocked_event_round_trip() -> Result<(), Box<dyn Error>> {
1274 let event = ThreadEvent::TurnBlocked(TurnBlockedEvent {
1275 message: "Blocked tool-call limit reached after 3 consecutive blocked calls.".to_string(),
1276 last_tool: Some("exec_command".to_string()),
1277 blocked_streak: 4,
1278 blocked_total: 4,
1279 consecutive_cap: 3,
1280 total_cap: 6,
1281 recovery_active: false,
1282 usage: None,
1283 });
1284
1285 let json = serde_json::to_string(&event)?;
1286 assert!(json.contains("turn.blocked"));
1287 let restored: ThreadEvent = serde_json::from_str(&json)?;
1288 assert_eq!(restored, event);
1289
1290 let legacy = serde_json::json!({"type": "turn.blocked", "message": "blocked"});
1292 let parsed: ThreadEvent = serde_json::from_value(legacy)?;
1293 assert!(matches!(parsed, ThreadEvent::TurnBlocked(_)));
1294 Ok(())
1295 }
1296
1297 #[test]
1298 fn usage_uncached_input_tokens_saturates() {
1299 let usage = Usage {
1300 input_tokens: 1_000,
1301 cached_input_tokens: 800,
1302 cache_creation_tokens: 100,
1303 output_tokens: 50,
1304 };
1305 assert_eq!(usage.uncached_input_tokens(), 100);
1306
1307 let inconsistent = Usage {
1308 input_tokens: 100,
1309 cached_input_tokens: 150,
1310 cache_creation_tokens: 0,
1311 output_tokens: 0,
1312 };
1313 assert_eq!(inconsistent.uncached_input_tokens(), 0);
1314
1315 let inconsistent_with_creation = Usage {
1316 input_tokens: 100,
1317 cached_input_tokens: 80,
1318 cache_creation_tokens: 50,
1319 output_tokens: 0,
1320 };
1321 assert_eq!(inconsistent_with_creation.uncached_input_tokens(), 0);
1322 }
1323
1324 #[test]
1325 fn usage_cache_hit_rate() {
1326 assert_eq!(Usage::default().cache_hit_rate(), None);
1327
1328 let usage = Usage {
1329 input_tokens: 1_000,
1330 cached_input_tokens: 750,
1331 cache_creation_tokens: 0,
1332 output_tokens: 0,
1333 };
1334 let rate = usage.cache_hit_rate().expect("rate");
1335 assert!((rate - 0.75).abs() < f64::EPSILON);
1336 }
1337
1338 #[test]
1339 fn usage_cache_summary_formats() {
1340 assert_eq!(Usage::default().cache_summary(), "No input tokens recorded.");
1341
1342 let usage = Usage {
1343 input_tokens: 1_000,
1344 cached_input_tokens: 800,
1345 cache_creation_tokens: 100,
1346 output_tokens: 50,
1347 };
1348 assert_eq!(
1349 usage.cache_summary(),
1350 "Cache: 800 cached / 1000 total input (80.0% hit rate), 100 cache-creation, 100 uncached"
1351 );
1352 }
1353
1354 #[test]
1355 fn usage_add_accumulates_all_fields_with_saturation() {
1356 let mut total = Usage {
1357 input_tokens: 100,
1358 cached_input_tokens: 20,
1359 cache_creation_tokens: 5,
1360 output_tokens: 10,
1361 };
1362 total.add(&Usage {
1363 input_tokens: 50,
1364 cached_input_tokens: 10,
1365 cache_creation_tokens: 2,
1366 output_tokens: 8,
1367 });
1368
1369 assert_eq!(total.input_tokens, 150);
1370 assert_eq!(total.cached_input_tokens, 30);
1371 assert_eq!(total.cache_creation_tokens, 7);
1372 assert_eq!(total.output_tokens, 18);
1373
1374 let mut saturating = Usage {
1375 input_tokens: u64::MAX,
1376 cached_input_tokens: u64::MAX,
1377 cache_creation_tokens: u64::MAX,
1378 output_tokens: u64::MAX,
1379 };
1380 saturating.add(&Usage {
1381 input_tokens: 1,
1382 cached_input_tokens: 1,
1383 cache_creation_tokens: 1,
1384 output_tokens: 1,
1385 });
1386 assert_eq!(saturating.input_tokens, u64::MAX);
1387 assert_eq!(saturating.cached_input_tokens, u64::MAX);
1388 assert_eq!(saturating.cache_creation_tokens, u64::MAX);
1389 assert_eq!(saturating.output_tokens, u64::MAX);
1390 }
1391
1392 #[test]
1393 fn versioned_event_wraps_schema_version() {
1394 let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "abc".to_string() });
1395
1396 let versioned = VersionedThreadEvent::new(event.clone());
1397
1398 assert_eq!(versioned.schema_version, EVENT_SCHEMA_VERSION);
1399 assert_eq!(versioned.event, event);
1400 assert_eq!(versioned.into_event(), event);
1401 }
1402
1403 #[test]
1404 fn plan_approval_events_round_trip_with_decision() {
1405 let requested = ThreadEvent::PlanApprovalRequested(PlanApprovalRequestedEvent {
1406 thread_id: "thread-1".to_string(),
1407 turn_id: "turn-2".to_string(),
1408 plan_file: Some(".vtcode/plans/change.md".to_string()),
1409 });
1410 let resolved = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1411 thread_id: "thread-1".to_string(),
1412 turn_id: "turn-3".to_string(),
1413 decision: PlanApprovalDecision::AutoAccept,
1414 automatic: false,
1415 });
1416
1417 for event in [requested, resolved] {
1418 let serialized = serde_json::to_string(&event).expect("serialize plan approval event");
1419 let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize plan approval event");
1420 assert_eq!(restored, event);
1421 }
1422 }
1423
1424 #[test]
1425 fn context_reset_event_round_trips_with_handoff_metadata() {
1426 let event = ThreadEvent::ContextReset(ContextResetEvent {
1427 thread_id: "thread-1".to_string(),
1428 turn_id: "turn-3".to_string(),
1429 trigger: ContextResetTrigger::PlanApproval,
1430 plan_preserved: true,
1431 previous_context_usage_percent: 7,
1432 tool_budget_reset: true,
1433 });
1434
1435 let serialized = serde_json::to_string(&event).expect("serialize context reset event");
1436 let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize context reset event");
1437 assert_eq!(restored, event);
1438 assert_eq!(serde_json::to_value(event).expect("wire value")["type"], "context.reset");
1439 }
1440
1441 #[test]
1442 fn plan_approval_decision_uses_stable_wire_names() {
1443 let event = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1444 thread_id: "thread-1".to_string(),
1445 turn_id: "turn-1".to_string(),
1446 decision: PlanApprovalDecision::SwitchBuild,
1447 automatic: false,
1448 });
1449
1450 let serialized = serde_json::to_value(event).expect("serialize plan approval decision");
1451 assert_eq!(serialized["type"], "plan.approval.resolved");
1452 assert_eq!(serialized["decision"], "switch_build");
1453 }
1454
1455 #[test]
1456 fn plan_approval_decision_is_forward_compatible() {
1457 let payload = serde_json::json!({
1458 "type": "plan.approval.resolved",
1459 "thread_id": "thread-1",
1460 "turn_id": "turn-1",
1461 "decision": "future_decision",
1462 "automatic": true,
1463 });
1464 let event: ThreadEvent = serde_json::from_value(payload).expect("future decision should deserialize");
1465 assert!(matches!(
1466 event,
1467 ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1468 decision: PlanApprovalDecision::Unknown,
1469 automatic: true,
1470 ..
1471 })
1472 ));
1473 }
1474
1475 #[cfg(feature = "serde-json")]
1476 #[test]
1477 fn versioned_json_round_trip() -> Result<(), Box<dyn Error>> {
1478 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1479 item: ThreadItem {
1480 id: "item-1".to_string(),
1481 details: ThreadItemDetails::AgentMessage(AgentMessageItem { text: "hello".to_string() }),
1482 },
1483 });
1484
1485 let payload = json::versioned_to_string(&event)?;
1486 let restored = json::versioned_from_str(&payload)?;
1487
1488 assert_eq!(restored.schema_version, EVENT_SCHEMA_VERSION);
1489 assert_eq!(restored.event, event);
1490 Ok(())
1491 }
1492
1493 #[test]
1494 fn compaction_trigger_serializes_snake_case_and_round_trips() {
1495 for trigger in [
1496 CompactionTrigger::Manual,
1497 CompactionTrigger::Auto,
1498 CompactionTrigger::Recovery,
1499 CompactionTrigger::ModelSwitch,
1500 CompactionTrigger::Unknown,
1501 ] {
1502 let json = serde_json::to_string(&trigger).unwrap();
1503 assert_eq!(json, format!("\"{}\"", trigger.as_str()));
1504 let restored: CompactionTrigger = serde_json::from_str(&json).unwrap();
1505 assert_eq!(restored, trigger);
1506 }
1507 }
1508
1509 #[test]
1510 fn tool_invocation_round_trip() -> Result<(), Box<dyn Error>> {
1511 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1512 item: ThreadItem {
1513 id: "tool_1".to_string(),
1514 details: ThreadItemDetails::ToolInvocation(ToolInvocationItem {
1515 tool_name: "read_file".to_string(),
1516 arguments: Some(serde_json::json!({ "path": "README.md" })),
1517 tool_call_id: Some("tool_call_0".to_string()),
1518 status: ToolCallStatus::Completed,
1519 outcome: None,
1520 }),
1521 },
1522 });
1523
1524 let json = serde_json::to_string(&event)?;
1525 let restored: ThreadEvent = serde_json::from_str(&json)?;
1526
1527 assert_eq!(restored, event);
1528 Ok(())
1529 }
1530
1531 #[test]
1532 fn tool_outcome_serializes_snake_case() {
1533 for outcome in [
1534 ToolOutcome::Success,
1535 ToolOutcome::Error,
1536 ToolOutcome::PermissionRejected,
1537 ToolOutcome::PermissionCancelled,
1538 ToolOutcome::Followup,
1539 ToolOutcome::HookDenied,
1540 ToolOutcome::InvalidTool,
1541 ToolOutcome::Cancelled,
1542 ] {
1543 let json = serde_json::to_string(&outcome).unwrap();
1544 let restored: ToolOutcome = serde_json::from_str(&json).unwrap();
1545 assert_eq!(restored, outcome);
1546 }
1547 }
1548
1549 #[test]
1550 fn tool_invocation_outcome_round_trip() -> Result<(), Box<dyn Error>> {
1551 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1552 item: ThreadItem {
1553 id: "tool_1".to_string(),
1554 details: ThreadItemDetails::ToolInvocation(ToolInvocationItem {
1555 tool_name: "exec_command".to_string(),
1556 arguments: Some(serde_json::json!({ "command": ["pwd"] })),
1557 tool_call_id: Some("tool_call_0".to_string()),
1558 status: ToolCallStatus::Failed,
1559 outcome: Some(ToolOutcome::PermissionRejected),
1560 }),
1561 },
1562 });
1563
1564 let json = serde_json::to_string(&event)?;
1565 let restored: ThreadEvent = serde_json::from_str(&json)?;
1566
1567 assert_eq!(restored, event);
1568 Ok(())
1569 }
1570
1571 #[test]
1572 fn tool_output_round_trip_preserves_raw_tool_call_id() -> Result<(), Box<dyn Error>> {
1573 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1574 item: ThreadItem {
1575 id: "tool_1:output".to_string(),
1576 details: ThreadItemDetails::ToolOutput(ToolOutputItem {
1577 call_id: "tool_1".to_string(),
1578 tool_call_id: Some("tool_call_0".to_string()),
1579 spool_path: None,
1580 output: "done".to_string(),
1581 exit_code: Some(0),
1582 status: ToolCallStatus::Completed,
1583 }),
1584 },
1585 });
1586
1587 let json = serde_json::to_string(&event)?;
1588 let restored: ThreadEvent = serde_json::from_str(&json)?;
1589
1590 assert_eq!(restored, event);
1591 Ok(())
1592 }
1593
1594 #[test]
1595 fn harness_item_round_trip() -> Result<(), Box<dyn Error>> {
1596 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1597 item: ThreadItem {
1598 id: "harness_1".to_string(),
1599 details: ThreadItemDetails::Harness(HarnessEventItem {
1600 event: HarnessEventKind::VerificationFailed,
1601 message: Some("cargo check failed".to_string()),
1602 command: Some("cargo check".to_string()),
1603 path: None,
1604 exit_code: Some(101),
1605 attempt: None,
1606 error_category: None,
1607 duration_ms: None,
1608 }),
1609 },
1610 });
1611
1612 let json = serde_json::to_string(&event)?;
1613 let restored: ThreadEvent = serde_json::from_str(&json)?;
1614
1615 assert_eq!(restored, event);
1616 Ok(())
1617 }
1618
1619 #[test]
1620 fn blocked_handoff_resolved_uses_stable_wire_name() -> Result<(), Box<dyn Error>> {
1621 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1622 item: ThreadItem {
1623 id: "harness_resolved".to_string(),
1624 details: ThreadItemDetails::Harness(HarnessEventItem {
1625 event: HarnessEventKind::BlockedHandoffResolved,
1626 message: Some("resolved".to_string()),
1627 command: None,
1628 path: None,
1629 exit_code: None,
1630 attempt: None,
1631 error_category: None,
1632 duration_ms: None,
1633 }),
1634 },
1635 });
1636
1637 let value = serde_json::to_value(&event)?;
1638 assert_eq!(value["item"]["event"], "blocked_handoff_resolved");
1639
1640 let restored: ThreadEvent = serde_json::from_value(value)?;
1641 assert_eq!(restored, event);
1642 Ok(())
1643 }
1644
1645 #[test]
1646 fn thread_completed_round_trip() -> Result<(), Box<dyn Error>> {
1647 let event = ThreadEvent::ThreadCompleted(ThreadCompletedEvent {
1648 thread_id: "thread-1".to_string(),
1649 session_id: "session-1".to_string(),
1650 subtype: ThreadCompletionSubtype::ErrorMaxBudgetUsd,
1651 outcome_code: "budget_limit_reached".to_string(),
1652 result: None,
1653 stop_reason: Some("max_tokens".to_string()),
1654 usage: Usage {
1655 input_tokens: 10,
1656 cached_input_tokens: 4,
1657 cache_creation_tokens: 2,
1658 output_tokens: 5,
1659 },
1660 total_cost_usd: serde_json::Number::from_f64(1.25),
1661 num_turns: 3,
1662 });
1663
1664 let json = serde_json::to_string(&event)?;
1665 let restored: ThreadEvent = serde_json::from_str(&json)?;
1666
1667 assert_eq!(restored, event);
1668 Ok(())
1669 }
1670
1671 #[test]
1672 fn compact_boundary_round_trip() -> Result<(), Box<dyn Error>> {
1673 let event = ThreadEvent::ThreadCompactBoundary(ThreadCompactBoundaryEvent {
1674 thread_id: "thread-1".to_string(),
1675 trigger: CompactionTrigger::Recovery,
1676 mode: CompactionMode::Provider,
1677 original_message_count: 12,
1678 compacted_message_count: 5,
1679 history_artifact_path: Some("/tmp/history.jsonl".to_string()),
1680 previous_segment_id: Some("segment-0001".to_string()),
1681 new_segment_id: Some("segment-0002".to_string()),
1682 previous_prefix_hash: Some("prefix-before".to_string()),
1683 new_prefix_hash: Some("prefix-after".to_string()),
1684 previous_catalog_hash: Some("catalog-before".to_string()),
1685 new_catalog_hash: Some("catalog-after".to_string()),
1686 });
1687
1688 let json = serde_json::to_string(&event)?;
1689 let restored: ThreadEvent = serde_json::from_str(&json)?;
1690
1691 assert_eq!(restored, event);
1692 Ok(())
1693 }
1694
1695 #[test]
1696 fn compact_boundary_deserializes_legacy_payload_without_segment_metadata() -> Result<(), Box<dyn Error>> {
1697 let payload = r#"{
1698 "type":"thread.compact_boundary",
1699 "thread_id":"thread-1",
1700 "trigger":"recovery",
1701 "mode":"provider",
1702 "original_message_count":12,
1703 "compacted_message_count":5
1704 }"#;
1705
1706 let restored: ThreadEvent = serde_json::from_str(payload)?;
1707 let ThreadEvent::ThreadCompactBoundary(event) = restored else {
1708 panic!("expected thread.compact_boundary event");
1709 };
1710
1711 assert_eq!(event.thread_id, "thread-1");
1712 assert_eq!(event.history_artifact_path, None);
1713 assert_eq!(event.previous_segment_id, None);
1714 assert_eq!(event.new_segment_id, None);
1715 assert_eq!(event.previous_prefix_hash, None);
1716 assert_eq!(event.new_prefix_hash, None);
1717 assert_eq!(event.previous_catalog_hash, None);
1718 assert_eq!(event.new_catalog_hash, None);
1719 Ok(())
1720 }
1721}