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.17.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(Box<ThreadCompletedEvent>),
375 #[serde(rename = "thread.compact_boundary")]
377 ThreadCompactBoundary(Box<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(Box<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(Box<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 #[serde(default, skip_serializing_if = "Option::is_none")]
522 pub completed_at: Option<String>,
523 pub thread_id: String,
525 pub session_id: String,
527 pub subtype: ThreadCompletionSubtype,
529 pub outcome_code: String,
531 #[serde(skip_serializing_if = "Option::is_none")]
533 pub result: Option<String>,
534 #[serde(skip_serializing_if = "Option::is_none")]
536 pub stop_reason: Option<String>,
537 pub usage: Usage,
539 #[serde(skip_serializing_if = "Option::is_none")]
541 pub total_cost_usd: Option<serde_json::Number>,
542 pub num_turns: usize,
544}
545
546#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
547#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
548pub struct ThreadCompactBoundaryEvent {
549 pub thread_id: String,
551 pub trigger: CompactionTrigger,
553 pub mode: CompactionMode,
555 pub original_message_count: usize,
557 pub compacted_message_count: usize,
559 #[serde(skip_serializing_if = "Option::is_none")]
561 pub history_artifact_path: Option<String>,
562 #[serde(skip_serializing_if = "Option::is_none")]
564 pub previous_segment_id: Option<String>,
565 #[serde(skip_serializing_if = "Option::is_none")]
567 pub new_segment_id: Option<String>,
568 #[serde(skip_serializing_if = "Option::is_none")]
570 pub previous_prefix_hash: Option<String>,
571 #[serde(skip_serializing_if = "Option::is_none")]
573 pub new_prefix_hash: Option<String>,
574 #[serde(skip_serializing_if = "Option::is_none")]
576 pub previous_catalog_hash: Option<String>,
577 #[serde(skip_serializing_if = "Option::is_none")]
579 pub new_catalog_hash: Option<String>,
580}
581
582#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
583#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
584#[serde(rename_all = "snake_case")]
585pub enum ContextResetTrigger {
586 PlanApproval,
588 #[serde(other)]
590 Unknown,
591}
592
593#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
594#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
595pub struct ContextResetEvent {
596 pub thread_id: String,
598 pub turn_id: String,
600 pub trigger: ContextResetTrigger,
602 pub plan_preserved: bool,
604 pub previous_context_usage_percent: u8,
606 pub tool_budget_reset: bool,
608}
609
610#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
611#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
612pub struct TurnStartedEvent {
613 #[serde(skip_serializing_if = "Option::is_none")]
617 token_breakdown: Option<Box<TokenBreakdown>>,
618 #[serde(default, skip_serializing_if = "Option::is_none")]
620 pub context: Option<Box<ExecutionContext>>,
621}
622
623impl TurnStartedEvent {
624 pub fn token_breakdown(&self) -> Option<&TokenBreakdown> {
626 self.token_breakdown.as_deref()
627 }
628}
629
630#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
632#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
633#[serde(rename_all = "snake_case")]
634pub enum InputOrigin {
635 User,
636 Correction,
637 PlanApproval,
638 Continuation,
639 Retry,
640}
641
642#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
644#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
645pub struct ExecutionContext {
646 pub task_id: String,
647 pub turn_id: String,
648 pub actor_id: String,
649 #[serde(default, skip_serializing_if = "Option::is_none")]
650 pub parent_actor_id: Option<String>,
651 pub origin: InputOrigin,
652 pub timestamp: String,
653 #[serde(default, skip_serializing_if = "Option::is_none")]
654 pub goal: Option<String>,
655}
656
657#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
659#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
660#[serde(rename_all = "snake_case")]
661pub enum CommandActivity {
662 Inspection,
663 Verification,
664 Mutation,
665}
666
667#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
669#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
670pub struct ItemContext {
671 pub task_id: String,
672 pub turn_id: String,
673 pub actor_id: String,
674 #[serde(default, skip_serializing_if = "Option::is_none")]
675 pub parent_actor_id: Option<String>,
676 pub timestamp: String,
677 #[serde(default, skip_serializing_if = "Option::is_none")]
678 pub activity: Option<CommandActivity>,
679}
680
681#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
683#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
684pub struct TokenBreakdown {
685 system_prompt_tokens: u64,
687 tool_schema_tokens: u64,
689 instruction_file_tokens: u64,
691 message_history_tokens: u64,
693 cache_read_tokens: u64,
695 cache_write_tokens: u64,
697 cache_miss_tokens: u64,
699 #[serde(skip_serializing_if = "Option::is_none")]
701 subagent_bootstrap_tokens: Option<u64>,
702}
703
704pub const MAX_IN_PROGRESS_EXEC_SESSIONS: usize = 4;
708
709#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
710#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
711pub struct TurnCompletedEvent {
712 #[serde(default, skip_serializing_if = "Option::is_none")]
713 pub completed_at: Option<Box<String>>,
714 pub usage: Usage,
716 #[serde(
721 default,
722 skip_serializing_if = "Vec::is_empty",
723 deserialize_with = "deserialize_null_as_default"
724 )]
725 pub in_progress_exec_sessions: Vec<String>,
726}
727
728#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
729#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
730pub struct TurnFailedEvent {
731 #[serde(default, skip_serializing_if = "Option::is_none")]
732 pub completed_at: Option<Box<String>>,
733 pub message: String,
735 #[serde(skip_serializing_if = "Option::is_none")]
737 pub usage: Option<Usage>,
738}
739
740#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
741#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
742pub struct TurnBlockedEvent {
743 #[serde(default, skip_serializing_if = "Option::is_none")]
745 pub completed_at: Option<String>,
746 pub message: String,
748 #[serde(skip_serializing_if = "Option::is_none")]
750 pub last_tool: Option<String>,
751 #[serde(default)]
753 pub blocked_streak: usize,
754 #[serde(default)]
756 pub blocked_total: usize,
757 #[serde(default)]
759 pub consecutive_cap: usize,
760 #[serde(default)]
762 pub total_cap: usize,
763 #[serde(default)]
765 pub recovery_active: bool,
766 #[serde(skip_serializing_if = "Option::is_none")]
768 pub usage: Option<Usage>,
769}
770
771#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
772#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
773pub struct ThreadErrorEvent {
774 pub message: String,
776}
777
778#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
779#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
780pub struct Usage {
781 #[serde(default, deserialize_with = "deserialize_null_as_default")]
783 pub input_tokens: u64,
784 #[serde(default, deserialize_with = "deserialize_null_as_default")]
786 pub cached_input_tokens: u64,
787 #[serde(default, deserialize_with = "deserialize_null_as_default")]
789 pub cache_creation_tokens: u64,
790 #[serde(default, deserialize_with = "deserialize_null_as_default")]
792 pub output_tokens: u64,
793}
794
795pub fn deserialize_null_as_default<'de, D, T>(deserializer: D) -> Result<T, D::Error>
802where
803 D: serde::Deserializer<'de>,
804 T: Deserialize<'de> + Default,
805{
806 Ok(Option::<T>::deserialize(deserializer)?.unwrap_or_default())
807}
808
809impl Usage {
810 #[must_use]
815 fn uncached_input_tokens(&self) -> u64 {
816 self.input_tokens
817 .saturating_sub(self.cached_input_tokens)
818 .saturating_sub(self.cache_creation_tokens)
819 }
820
821 #[must_use]
824 pub fn cache_hit_rate(&self) -> Option<f64> {
825 if self.input_tokens == 0 {
826 return None;
827 }
828 Some(self.cached_input_tokens as f64 / self.input_tokens as f64)
829 }
830
831 #[must_use]
833 pub fn cache_summary(&self) -> String {
834 let total_input = self.input_tokens;
835 if total_input == 0 {
836 return "No input tokens recorded.".to_string();
837 }
838
839 let cached = self.cached_input_tokens;
840 let creation = self.cache_creation_tokens;
841 let uncached = self.uncached_input_tokens();
842 let rate = cached as f64 / total_input as f64 * 100.0;
843 format!(
844 "Cache: {cached} cached / {total_input} total input ({rate:.1}% hit rate), \
845 {creation} cache-creation, {uncached} uncached"
846 )
847 }
848
849 pub fn add(&mut self, other: &Usage) {
851 self.input_tokens = self.input_tokens.saturating_add(other.input_tokens);
852 self.cached_input_tokens = self.cached_input_tokens.saturating_add(other.cached_input_tokens);
853 self.cache_creation_tokens = self.cache_creation_tokens.saturating_add(other.cache_creation_tokens);
854 self.output_tokens = self.output_tokens.saturating_add(other.output_tokens);
855 }
856}
857
858#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
859#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
860pub struct ItemCompletedEvent {
861 pub item: ThreadItem,
863}
864
865#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
866#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
867pub struct ItemStartedEvent {
868 pub item: ThreadItem,
870}
871
872#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
873#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
874pub struct ItemUpdatedEvent {
875 pub item: ThreadItem,
877}
878
879#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
880#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
881pub struct PlanDeltaEvent {
882 pub thread_id: String,
884 pub turn_id: String,
886 pub item_id: String,
888 pub delta: String,
890}
891
892#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
893#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
894pub struct PlanApprovalRequestedEvent {
895 pub thread_id: String,
897 pub turn_id: String,
899 #[serde(skip_serializing_if = "Option::is_none")]
901 pub plan_file: Option<String>,
902}
903
904#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
905#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
906#[serde(rename_all = "snake_case")]
907pub enum PlanApprovalDecision {
908 Execute,
910 AutoAccept,
912 FreshContext,
914 Revise,
916 Cancel,
918 SwitchBuild,
920 SwitchAuto,
922 #[serde(other)]
924 Unknown,
925}
926
927#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
928#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
929pub struct PlanApprovalResolvedEvent {
930 pub thread_id: String,
932 pub turn_id: String,
934 pub decision: PlanApprovalDecision,
936 pub automatic: bool,
938}
939
940#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
941#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
942pub struct ThreadItem {
943 #[serde(default, skip_serializing_if = "Option::is_none")]
944 pub context: Option<Box<ItemContext>>,
945 pub id: String,
947 #[serde(flatten)]
949 pub details: ThreadItemDetails,
950}
951
952#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
953#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
954#[serde(tag = "type", rename_all = "snake_case")]
955pub enum ThreadItemDetails {
956 AgentMessage(AgentMessageItem),
958 Plan(PlanItem),
960 Reasoning(Box<ReasoningItem>),
962 Decision(Box<DecisionItem>),
964 CommandExecution(Box<CommandExecutionItem>),
966 ToolInvocation(Box<ToolInvocationItem>),
968 ToolOutput(Box<ToolOutputItem>),
970 FileChange(Box<FileChangeItem>),
972 McpToolCall(Box<McpToolCallItem>),
974 WebSearch(Box<WebSearchItem>),
976 Harness(Box<HarnessEventItem>),
978 Error(ErrorItem),
980}
981
982#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
983#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
984pub struct AgentMessageItem {
985 pub text: String,
987}
988
989#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
991#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
992pub struct DecisionItem {
993 pub summary: String,
994 pub rationale: String,
995 #[serde(default)]
996 pub alternatives: Vec<String>,
997 #[serde(default)]
998 pub evidence_ids: Vec<String>,
999}
1000
1001#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1002#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1003pub struct PlanItem {
1004 pub text: String,
1006}
1007
1008#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1009#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1010pub struct ReasoningItem {
1011 pub text: String,
1013 #[serde(skip_serializing_if = "Option::is_none")]
1016 pub stage: Option<String>,
1017}
1018
1019#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
1020#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1021#[serde(rename_all = "snake_case")]
1022pub enum CommandExecutionStatus {
1023 #[default]
1025 Completed,
1026 Failed,
1028 InProgress,
1030}
1031
1032#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1033#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1034pub struct CommandExecutionItem {
1035 pub command: String,
1037 #[serde(skip_serializing_if = "Option::is_none")]
1039 pub arguments: Option<Value>,
1040 #[serde(default)]
1042 pub aggregated_output: String,
1043 #[serde(skip_serializing_if = "Option::is_none")]
1045 pub exit_code: Option<i32>,
1046 pub status: CommandExecutionStatus,
1048}
1049
1050#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
1051#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1052#[serde(rename_all = "snake_case")]
1053pub enum ToolCallStatus {
1054 #[default]
1056 Completed,
1057 Failed,
1059 InProgress,
1061}
1062
1063#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
1071#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1072#[serde(rename_all = "snake_case")]
1073pub enum ToolOutcome {
1074 #[default]
1076 Success,
1077 Error,
1079 PermissionRejected,
1081 PermissionCancelled,
1083 Followup,
1085 HookDenied,
1087 InvalidTool,
1089 Cancelled,
1091}
1092
1093impl ToolOutcome {
1094 #[must_use]
1095 pub const fn is_terminal(self) -> bool {
1096 !matches!(self, Self::Followup)
1097 }
1098}
1099
1100#[must_use]
1107#[allow(
1108 clippy::unreachable,
1109 reason = "Intentional compatibility, platform, or test-only suppression."
1110)]
1111pub fn tool_outcome_from_status(status: &ToolCallStatus) -> ToolOutcome {
1112 match status {
1113 ToolCallStatus::Completed => ToolOutcome::Success,
1114 ToolCallStatus::Failed => ToolOutcome::Error,
1115 ToolCallStatus::InProgress => unreachable!("InProgress status passed to completion event"),
1116 }
1117}
1118
1119#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1120#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1121pub struct ToolInvocationItem {
1122 pub tool_name: String,
1124 #[serde(skip_serializing_if = "Option::is_none")]
1126 pub arguments: Option<Value>,
1127 #[serde(skip_serializing_if = "Option::is_none")]
1129 pub tool_call_id: Option<String>,
1130 pub status: ToolCallStatus,
1132 #[serde(skip_serializing_if = "Option::is_none")]
1134 pub outcome: Option<ToolOutcome>,
1135}
1136
1137#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1138#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1139pub struct ToolOutputItem {
1140 pub call_id: String,
1142 #[serde(skip_serializing_if = "Option::is_none")]
1144 pub tool_call_id: Option<String>,
1145 #[serde(skip_serializing_if = "Option::is_none")]
1147 pub spool_path: Option<String>,
1148 #[serde(default)]
1150 pub output: String,
1151 #[serde(skip_serializing_if = "Option::is_none")]
1153 pub exit_code: Option<i32>,
1154 pub status: ToolCallStatus,
1156}
1157
1158#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1159#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1160pub struct FileChangeItem {
1161 #[serde(default, skip_serializing_if = "Option::is_none")]
1163 pub diff_incomplete: Option<bool>,
1164 pub changes: Vec<FileUpdateChange>,
1166 pub status: PatchApplyStatus,
1168 #[serde(default, skip_serializing_if = "Option::is_none")]
1173 pub unified_diff: Option<String>,
1174 #[serde(default, skip_serializing_if = "Option::is_none")]
1176 pub additions: Option<u64>,
1177 #[serde(default, skip_serializing_if = "Option::is_none")]
1179 pub deletions: Option<u64>,
1180}
1181
1182#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1183#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1184pub struct FileUpdateChange {
1185 pub path: String,
1187 pub kind: PatchChangeKind,
1189}
1190
1191#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1192#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1193#[serde(rename_all = "snake_case")]
1194pub enum PatchApplyStatus {
1195 Completed,
1197 Failed,
1199}
1200
1201#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1202#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1203#[serde(rename_all = "snake_case")]
1204pub enum PatchChangeKind {
1205 Add,
1207 Delete,
1209 Update,
1211}
1212
1213#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1214#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1215pub struct McpToolCallItem {
1216 pub tool_name: String,
1218 #[serde(skip_serializing_if = "Option::is_none")]
1220 pub arguments: Option<Value>,
1221 #[serde(skip_serializing_if = "Option::is_none")]
1223 pub result: Option<String>,
1224 #[serde(skip_serializing_if = "Option::is_none")]
1226 pub status: Option<McpToolCallStatus>,
1227}
1228
1229#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1230#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1231#[serde(rename_all = "snake_case")]
1232pub enum McpToolCallStatus {
1233 Started,
1235 Completed,
1237 Failed,
1239}
1240
1241#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1242#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1243pub struct WebSearchItem {
1244 pub query: String,
1246 #[serde(skip_serializing_if = "Option::is_none")]
1248 pub provider: Option<String>,
1249 #[serde(skip_serializing_if = "Option::is_none")]
1251 pub results: Option<Vec<String>>,
1252}
1253
1254#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1255#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1256#[serde(rename_all = "snake_case")]
1257pub enum HarnessEventKind {
1258 PlanningStarted,
1259 PlanningCompleted,
1260 ContinuationStarted,
1261 ContinuationSkipped,
1262 TurnBlocked,
1265 BlockedRecoveryStarted,
1267 BlockedRecoveryFinished,
1269 BlockedHandoffWritten,
1270 BlockedHandoffResolved,
1273 EvaluationStarted,
1274 EvaluationPassed,
1275 EvaluationFailed,
1276 RevisionStarted,
1277 EscalationTriggered,
1278 EscalationBypassed,
1279 VerificationStarted,
1280 VerificationPassed,
1281 VerificationFailed,
1282 ErrorRecovered,
1284 ToolRetryAttempted,
1286 ToolLatencyRecorded,
1288 SnapshotCreated,
1290 SnapshotRestored,
1292 SessionToolLimitIncreased,
1295 ToolLoopLimitIncreased,
1297 BackgroundSubprocessCompleted,
1299 DelegatedAgentStatus,
1301}
1302
1303#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1304#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1305#[serde(rename_all = "snake_case")]
1306pub enum PermissionDecision {
1307 Allow,
1308 Deny,
1309 Cancelled,
1310 Followup,
1311}
1312
1313#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1314#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1315pub struct PermissionRequestedEvent {
1316 pub tool_name: String,
1318}
1319
1320#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1321#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1322pub struct PermissionResolvedEvent {
1323 pub tool_name: String,
1325 pub decision: PermissionDecision,
1327 pub wait_ms: u64,
1329}
1330
1331#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1332#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1333#[serde(rename_all = "snake_case")]
1334pub enum InterjectionSource {
1335 Direct,
1336 Queue,
1337}
1338
1339#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1340#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1341#[serde(rename_all = "snake_case")]
1342pub enum RedirectKind {
1343 Interjection,
1344}
1345
1346#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1347#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1348pub struct InterjectedEvent {
1349 #[serde(default, skip_serializing_if = "Option::is_none")]
1351 pub text: Option<Box<String>>,
1352 pub source: InterjectionSource,
1354 pub image_count: u32,
1356 pub redirect_kind: RedirectKind,
1359}
1360
1361#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1362#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1363pub struct HarnessEventItem {
1364 pub event: HarnessEventKind,
1366 #[serde(skip_serializing_if = "Option::is_none")]
1368 pub message: Option<String>,
1369 #[serde(skip_serializing_if = "Option::is_none")]
1371 pub command: Option<String>,
1372 #[serde(skip_serializing_if = "Option::is_none")]
1374 pub path: Option<String>,
1375 #[serde(skip_serializing_if = "Option::is_none")]
1377 pub exit_code: Option<i32>,
1378 #[serde(skip_serializing_if = "Option::is_none")]
1380 pub attempt: Option<u32>,
1381 #[serde(skip_serializing_if = "Option::is_none")]
1383 pub error_category: Option<String>,
1384 #[serde(skip_serializing_if = "Option::is_none")]
1386 pub duration_ms: Option<u64>,
1387 #[serde(skip_serializing_if = "Option::is_none")]
1389 pub task_id: Option<String>,
1390 #[serde(skip_serializing_if = "Option::is_none")]
1392 pub session_id: Option<String>,
1393 #[serde(skip_serializing_if = "Option::is_none")]
1395 pub exec_session_id: Option<String>,
1396 #[serde(skip_serializing_if = "Option::is_none")]
1398 pub status: Option<String>,
1399 #[serde(skip_serializing_if = "Option::is_none")]
1401 pub transcript_path: Option<String>,
1402 #[serde(skip_serializing_if = "Option::is_none")]
1404 pub archive_path: Option<String>,
1405}
1406
1407#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1408#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1409pub struct ErrorItem {
1410 pub message: String,
1412}
1413
1414#[cfg(test)]
1415mod tests {
1416 use super::*;
1417 use std::error::Error;
1418 use std::mem::size_of;
1419
1420 #[test]
1425 fn thread_event_stays_compact() {
1426 assert!(
1427 size_of::<ThreadEvent>() <= 80,
1428 "ThreadEvent grew to {} bytes; box new large payloads instead of inlining them",
1429 size_of::<ThreadEvent>()
1430 );
1431 }
1432
1433 #[test]
1436 fn boxed_thread_item_details_payloads_stay_boxed() {
1437 assert!(size_of::<Option<Box<CommandExecutionItem>>>() < size_of::<Option<CommandExecutionItem>>());
1438 assert!(size_of::<Option<Box<ToolInvocationItem>>>() < size_of::<Option<ToolInvocationItem>>());
1439 assert!(size_of::<Option<Box<ToolOutputItem>>>() < size_of::<Option<ToolOutputItem>>());
1440 assert!(size_of::<Option<Box<FileChangeItem>>>() < size_of::<Option<FileChangeItem>>());
1441 assert!(size_of::<Option<Box<McpToolCallItem>>>() < size_of::<Option<McpToolCallItem>>());
1442 assert!(size_of::<Option<Box<WebSearchItem>>>() < size_of::<Option<WebSearchItem>>());
1443 assert!(size_of::<Option<Box<HarnessEventItem>>>() < size_of::<Option<HarnessEventItem>>());
1444 }
1445
1446 #[test]
1447 fn file_change_item_optional_diff_fields_round_trip() -> Result<(), Box<dyn Error>> {
1448 let legacy_json = r#"{
1450 "changes": [{"path": "src/main.rs", "kind": "add"}],
1451 "status": "completed"
1452 }"#;
1453 let legacy: FileChangeItem = serde_json::from_str(legacy_json)?;
1454 assert!(legacy.unified_diff.is_none());
1455 assert!(legacy.additions.is_none());
1456 assert!(legacy.deletions.is_none());
1457
1458 let legacy_reserialized = serde_json::to_value(&legacy)?;
1460 assert!(legacy_reserialized.get("unified_diff").is_none());
1461 assert!(legacy_reserialized.get("additions").is_none());
1462 assert!(legacy_reserialized.get("deletions").is_none());
1463
1464 let populated = FileChangeItem {
1466 diff_incomplete: None,
1467 changes: legacy.changes.clone(),
1468 status: PatchApplyStatus::Completed,
1469 unified_diff: Some("diff --git a/x b/x\n".to_string()),
1470 additions: Some(3),
1471 deletions: Some(1),
1472 };
1473 let json = serde_json::to_string(&populated)?;
1474 let restored: FileChangeItem = serde_json::from_str(&json)?;
1475 assert_eq!(restored, populated);
1476 Ok(())
1477 }
1478
1479 #[test]
1480 fn thread_event_round_trip() -> Result<(), Box<dyn Error>> {
1481 let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1482 completed_at: None,
1483 usage: Usage {
1484 input_tokens: 1,
1485 cached_input_tokens: 2,
1486 cache_creation_tokens: 0,
1487 output_tokens: 3,
1488 },
1489 in_progress_exec_sessions: Vec::new(),
1490 });
1491
1492 let json = serde_json::to_string(&event)?;
1493 let restored: ThreadEvent = serde_json::from_str(&json)?;
1494
1495 assert_eq!(restored, event);
1496 Ok(())
1497 }
1498
1499 #[test]
1500 fn turn_blocked_event_round_trip() -> Result<(), Box<dyn Error>> {
1501 let event = ThreadEvent::TurnBlocked(Box::new(TurnBlockedEvent {
1502 completed_at: None,
1503 message: "Blocked tool-call limit reached after 3 consecutive blocked calls.".to_string(),
1504 last_tool: Some("exec_command".to_string()),
1505 blocked_streak: 4,
1506 blocked_total: 4,
1507 consecutive_cap: 3,
1508 total_cap: 6,
1509 recovery_active: false,
1510 usage: None,
1511 }));
1512
1513 let json = serde_json::to_string(&event)?;
1514 assert!(json.contains("turn.blocked"));
1515 let restored: ThreadEvent = serde_json::from_str(&json)?;
1516 assert_eq!(restored, event);
1517
1518 let legacy = serde_json::json!({"type": "turn.blocked", "message": "blocked"});
1520 let parsed: ThreadEvent = serde_json::from_value(legacy)?;
1521 assert!(matches!(parsed, ThreadEvent::TurnBlocked(_)));
1522 Ok(())
1523 }
1524
1525 #[test]
1526 fn turn_completed_in_progress_sessions_default_empty_and_omitted() -> Result<(), Box<dyn Error>> {
1527 let legacy = serde_json::json!({
1529 "type": "turn.completed",
1530 "usage": {"input_tokens": 1, "cached_input_tokens": 0, "cache_creation_tokens": 0, "output_tokens": 2}
1531 });
1532 let parsed: ThreadEvent = serde_json::from_value(legacy)?;
1533 let ThreadEvent::TurnCompleted(completed) = parsed else {
1534 panic!("expected turn.completed");
1535 };
1536 assert!(completed.in_progress_exec_sessions.is_empty());
1537
1538 let json = serde_json::to_value(ThreadEvent::TurnCompleted(completed))?;
1540 assert!(json.get("in_progress_exec_sessions").is_none());
1541
1542 let null_field = serde_json::json!({
1544 "type": "turn.completed",
1545 "usage": {"input_tokens": 0, "cached_input_tokens": 0, "cache_creation_tokens": 0, "output_tokens": 0},
1546 "in_progress_exec_sessions": null
1547 });
1548 let parsed_null: ThreadEvent = serde_json::from_value(null_field)?;
1549 let ThreadEvent::TurnCompleted(null_completed) = parsed_null else {
1550 panic!("expected turn.completed");
1551 };
1552 assert!(null_completed.in_progress_exec_sessions.is_empty());
1553 Ok(())
1554 }
1555
1556 #[test]
1557 fn turn_completed_in_progress_sessions_round_trip_and_bound() -> Result<(), Box<dyn Error>> {
1558 assert_eq!(MAX_IN_PROGRESS_EXEC_SESSIONS, 4);
1559 let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1560 completed_at: None,
1561 usage: Usage::default(),
1562 in_progress_exec_sessions: vec!["run-1".to_string(), "run-2".to_string()],
1563 });
1564 let json = serde_json::to_string(&event)?;
1565 assert!(json.contains("in_progress_exec_sessions"));
1566 let restored: ThreadEvent = serde_json::from_str(&json)?;
1567 assert_eq!(restored, event);
1568 Ok(())
1569 }
1570
1571 #[test]
1572 fn usage_uncached_input_tokens_saturates() {
1573 let usage = Usage {
1574 input_tokens: 1_000,
1575 cached_input_tokens: 800,
1576 cache_creation_tokens: 100,
1577 output_tokens: 50,
1578 };
1579 assert_eq!(usage.uncached_input_tokens(), 100);
1580
1581 let inconsistent = Usage {
1582 input_tokens: 100,
1583 cached_input_tokens: 150,
1584 cache_creation_tokens: 0,
1585 output_tokens: 0,
1586 };
1587 assert_eq!(inconsistent.uncached_input_tokens(), 0);
1588
1589 let inconsistent_with_creation = Usage {
1590 input_tokens: 100,
1591 cached_input_tokens: 80,
1592 cache_creation_tokens: 50,
1593 output_tokens: 0,
1594 };
1595 assert_eq!(inconsistent_with_creation.uncached_input_tokens(), 0);
1596 }
1597
1598 #[test]
1599 fn usage_cache_hit_rate() {
1600 assert_eq!(Usage::default().cache_hit_rate(), None);
1601
1602 let usage = Usage {
1603 input_tokens: 1_000,
1604 cached_input_tokens: 750,
1605 cache_creation_tokens: 0,
1606 output_tokens: 0,
1607 };
1608 let rate = usage.cache_hit_rate().expect("rate");
1609 assert!((rate - 0.75).abs() < f64::EPSILON);
1610 }
1611
1612 #[test]
1613 fn usage_cache_summary_formats() {
1614 assert_eq!(Usage::default().cache_summary(), "No input tokens recorded.");
1615
1616 let usage = Usage {
1617 input_tokens: 1_000,
1618 cached_input_tokens: 800,
1619 cache_creation_tokens: 100,
1620 output_tokens: 50,
1621 };
1622 assert_eq!(
1623 usage.cache_summary(),
1624 "Cache: 800 cached / 1000 total input (80.0% hit rate), 100 cache-creation, 100 uncached"
1625 );
1626 }
1627
1628 #[test]
1629 fn usage_add_accumulates_all_fields_with_saturation() {
1630 let mut total = Usage {
1631 input_tokens: 100,
1632 cached_input_tokens: 20,
1633 cache_creation_tokens: 5,
1634 output_tokens: 10,
1635 };
1636 total.add(&Usage {
1637 input_tokens: 50,
1638 cached_input_tokens: 10,
1639 cache_creation_tokens: 2,
1640 output_tokens: 8,
1641 });
1642
1643 assert_eq!(total.input_tokens, 150);
1644 assert_eq!(total.cached_input_tokens, 30);
1645 assert_eq!(total.cache_creation_tokens, 7);
1646 assert_eq!(total.output_tokens, 18);
1647
1648 let mut saturating = Usage {
1649 input_tokens: u64::MAX,
1650 cached_input_tokens: u64::MAX,
1651 cache_creation_tokens: u64::MAX,
1652 output_tokens: u64::MAX,
1653 };
1654 saturating.add(&Usage {
1655 input_tokens: 1,
1656 cached_input_tokens: 1,
1657 cache_creation_tokens: 1,
1658 output_tokens: 1,
1659 });
1660 assert_eq!(saturating.input_tokens, u64::MAX);
1661 assert_eq!(saturating.cached_input_tokens, u64::MAX);
1662 assert_eq!(saturating.cache_creation_tokens, u64::MAX);
1663 assert_eq!(saturating.output_tokens, u64::MAX);
1664 }
1665
1666 #[test]
1667 fn versioned_event_wraps_schema_version() {
1668 let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "abc".to_string() });
1669
1670 let versioned = VersionedThreadEvent::new(event.clone());
1671
1672 assert_eq!(versioned.schema_version, EVENT_SCHEMA_VERSION);
1673 assert_eq!(versioned.event, event);
1674 assert_eq!(versioned.into_event(), event);
1675 }
1676
1677 #[test]
1678 fn plan_approval_events_round_trip_with_decision() {
1679 let requested = ThreadEvent::PlanApprovalRequested(PlanApprovalRequestedEvent {
1680 thread_id: "thread-1".to_string(),
1681 turn_id: "turn-2".to_string(),
1682 plan_file: Some(".vtcode/plans/change.md".to_string()),
1683 });
1684 let resolved = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1685 thread_id: "thread-1".to_string(),
1686 turn_id: "turn-3".to_string(),
1687 decision: PlanApprovalDecision::AutoAccept,
1688 automatic: false,
1689 });
1690
1691 for event in [requested, resolved] {
1692 let serialized = serde_json::to_string(&event).expect("serialize plan approval event");
1693 let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize plan approval event");
1694 assert_eq!(restored, event);
1695 }
1696 }
1697
1698 #[test]
1699 fn context_reset_event_round_trips_with_handoff_metadata() {
1700 let event = ThreadEvent::ContextReset(ContextResetEvent {
1701 thread_id: "thread-1".to_string(),
1702 turn_id: "turn-3".to_string(),
1703 trigger: ContextResetTrigger::PlanApproval,
1704 plan_preserved: true,
1705 previous_context_usage_percent: 7,
1706 tool_budget_reset: true,
1707 });
1708
1709 let serialized = serde_json::to_string(&event).expect("serialize context reset event");
1710 let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize context reset event");
1711 assert_eq!(restored, event);
1712 assert_eq!(serde_json::to_value(event).expect("wire value")["type"], "context.reset");
1713 }
1714
1715 #[test]
1716 fn plan_approval_decision_uses_stable_wire_names() {
1717 let event = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1718 thread_id: "thread-1".to_string(),
1719 turn_id: "turn-1".to_string(),
1720 decision: PlanApprovalDecision::SwitchBuild,
1721 automatic: false,
1722 });
1723
1724 let serialized = serde_json::to_value(event).expect("serialize plan approval decision");
1725 assert_eq!(serialized["type"], "plan.approval.resolved");
1726 assert_eq!(serialized["decision"], "switch_build");
1727 }
1728
1729 #[test]
1730 fn plan_approval_decision_is_forward_compatible() {
1731 let payload = serde_json::json!({
1732 "type": "plan.approval.resolved",
1733 "thread_id": "thread-1",
1734 "turn_id": "turn-1",
1735 "decision": "future_decision",
1736 "automatic": true,
1737 });
1738 let event: ThreadEvent = serde_json::from_value(payload).expect("future decision should deserialize");
1739 assert!(matches!(
1740 event,
1741 ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1742 decision: PlanApprovalDecision::Unknown,
1743 automatic: true,
1744 ..
1745 })
1746 ));
1747 }
1748
1749 #[cfg(feature = "serde-json")]
1750 #[test]
1751 fn versioned_json_round_trip() -> Result<(), Box<dyn Error>> {
1752 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1753 item: ThreadItem {
1754 context: None,
1755 id: "item-1".to_string(),
1756 details: ThreadItemDetails::AgentMessage(AgentMessageItem { text: "hello".to_string() }),
1757 },
1758 });
1759
1760 let payload = json::versioned_to_string(&event)?;
1761 let restored = json::versioned_from_str(&payload)?;
1762
1763 assert_eq!(restored.schema_version, EVENT_SCHEMA_VERSION);
1764 assert_eq!(restored.event, event);
1765 Ok(())
1766 }
1767
1768 #[test]
1769 fn compaction_trigger_serializes_snake_case_and_round_trips() {
1770 for trigger in [
1771 CompactionTrigger::Manual,
1772 CompactionTrigger::Auto,
1773 CompactionTrigger::Recovery,
1774 CompactionTrigger::ModelSwitch,
1775 CompactionTrigger::Unknown,
1776 ] {
1777 let json = serde_json::to_string(&trigger).unwrap();
1778 assert_eq!(json, format!("\"{}\"", trigger.as_str()));
1779 let restored: CompactionTrigger = serde_json::from_str(&json).unwrap();
1780 assert_eq!(restored, trigger);
1781 }
1782 }
1783
1784 #[test]
1785 fn tool_invocation_round_trip() -> Result<(), Box<dyn Error>> {
1786 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1787 item: ThreadItem {
1788 context: None,
1789 id: "tool_1".to_string(),
1790 details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
1791 tool_name: "read_file".to_string(),
1792 arguments: Some(serde_json::json!({ "path": "README.md" })),
1793 tool_call_id: Some("tool_call_0".to_string()),
1794 status: ToolCallStatus::Completed,
1795 outcome: None,
1796 })),
1797 },
1798 });
1799
1800 let json = serde_json::to_string(&event)?;
1801 let restored: ThreadEvent = serde_json::from_str(&json)?;
1802
1803 assert_eq!(restored, event);
1804 Ok(())
1805 }
1806
1807 #[test]
1808 fn tool_outcome_serializes_snake_case() {
1809 for outcome in [
1810 ToolOutcome::Success,
1811 ToolOutcome::Error,
1812 ToolOutcome::PermissionRejected,
1813 ToolOutcome::PermissionCancelled,
1814 ToolOutcome::Followup,
1815 ToolOutcome::HookDenied,
1816 ToolOutcome::InvalidTool,
1817 ToolOutcome::Cancelled,
1818 ] {
1819 let json = serde_json::to_string(&outcome).unwrap();
1820 let restored: ToolOutcome = serde_json::from_str(&json).unwrap();
1821 assert_eq!(restored, outcome);
1822 }
1823 }
1824
1825 #[test]
1826 fn tool_invocation_outcome_round_trip() -> Result<(), Box<dyn Error>> {
1827 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1828 item: ThreadItem {
1829 context: None,
1830 id: "tool_1".to_string(),
1831 details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
1832 tool_name: "exec_command".to_string(),
1833 arguments: Some(serde_json::json!({ "command": ["pwd"] })),
1834 tool_call_id: Some("tool_call_0".to_string()),
1835 status: ToolCallStatus::Failed,
1836 outcome: Some(ToolOutcome::PermissionRejected),
1837 })),
1838 },
1839 });
1840
1841 let json = serde_json::to_string(&event)?;
1842 let restored: ThreadEvent = serde_json::from_str(&json)?;
1843
1844 assert_eq!(restored, event);
1845 Ok(())
1846 }
1847
1848 #[test]
1849 fn tool_output_round_trip_preserves_raw_tool_call_id() -> Result<(), Box<dyn Error>> {
1850 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1851 item: ThreadItem {
1852 context: None,
1853 id: "tool_1:output".to_string(),
1854 details: ThreadItemDetails::ToolOutput(Box::new(ToolOutputItem {
1855 call_id: "tool_1".to_string(),
1856 tool_call_id: Some("tool_call_0".to_string()),
1857 spool_path: None,
1858 output: "done".to_string(),
1859 exit_code: Some(0),
1860 status: ToolCallStatus::Completed,
1861 })),
1862 },
1863 });
1864
1865 let json = serde_json::to_string(&event)?;
1866 let restored: ThreadEvent = serde_json::from_str(&json)?;
1867
1868 assert_eq!(restored, event);
1869 Ok(())
1870 }
1871
1872 #[test]
1873 fn harness_item_round_trip() -> Result<(), Box<dyn Error>> {
1874 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1875 item: ThreadItem {
1876 context: None,
1877 id: "harness_1".to_string(),
1878 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1879 event: HarnessEventKind::VerificationFailed,
1880 message: Some("cargo check failed".to_string()),
1881 command: Some("cargo check".to_string()),
1882 path: None,
1883 exit_code: Some(101),
1884 attempt: None,
1885 error_category: None,
1886 duration_ms: None,
1887 task_id: None,
1888 session_id: None,
1889 exec_session_id: None,
1890 status: None,
1891 transcript_path: None,
1892 archive_path: None,
1893 })),
1894 },
1895 });
1896
1897 let json = serde_json::to_string(&event)?;
1898 let restored: ThreadEvent = serde_json::from_str(&json)?;
1899
1900 assert_eq!(restored, event);
1901 Ok(())
1902 }
1903
1904 #[test]
1905 fn background_completion_harness_item_preserves_terminal_identity() -> Result<(), Box<dyn Error>> {
1906 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1907 item: ThreadItem {
1908 context: None,
1909 id: "background-completion:task:exec:0".to_string(),
1910 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1911 event: HarnessEventKind::BackgroundSubprocessCompleted,
1912 message: Some("Background subprocess completed successfully".to_string()),
1913 command: None,
1914 path: None,
1915 exit_code: Some(0),
1916 attempt: None,
1917 error_category: None,
1918 duration_ms: None,
1919 task_id: Some("task".to_string()),
1920 session_id: Some("child-session".to_string()),
1921 exec_session_id: Some("exec-session".to_string()),
1922 status: Some("stopped".to_string()),
1923 transcript_path: Some("/tmp/transcript.jsonl".to_string()),
1924 archive_path: Some("/tmp/archive.json".to_string()),
1925 })),
1926 },
1927 });
1928
1929 let value = serde_json::to_value(&event)?;
1930 assert_eq!(value["item"]["event"], "background_subprocess_completed");
1931 assert_eq!(value["item"]["task_id"], "task");
1932 assert_eq!(value["item"]["exec_session_id"], "exec-session");
1933 assert_eq!(value["item"]["status"], "stopped");
1934
1935 let restored: ThreadEvent = serde_json::from_value(value)?;
1936 assert_eq!(restored, event);
1937 Ok(())
1938 }
1939
1940 #[test]
1941 fn blocked_handoff_resolved_uses_stable_wire_name() -> Result<(), Box<dyn Error>> {
1942 let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1943 item: ThreadItem {
1944 context: None,
1945 id: "harness_resolved".to_string(),
1946 details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1947 event: HarnessEventKind::BlockedHandoffResolved,
1948 message: Some("resolved".to_string()),
1949 command: None,
1950 path: None,
1951 exit_code: None,
1952 attempt: None,
1953 error_category: None,
1954 duration_ms: None,
1955 task_id: None,
1956 session_id: None,
1957 exec_session_id: None,
1958 status: None,
1959 transcript_path: None,
1960 archive_path: None,
1961 })),
1962 },
1963 });
1964
1965 let value = serde_json::to_value(&event)?;
1966 assert_eq!(value["item"]["event"], "blocked_handoff_resolved");
1967
1968 let restored: ThreadEvent = serde_json::from_value(value)?;
1969 assert_eq!(restored, event);
1970 Ok(())
1971 }
1972
1973 #[test]
1974 fn thread_completed_round_trip() -> Result<(), Box<dyn Error>> {
1975 let event = ThreadEvent::ThreadCompleted(Box::new(ThreadCompletedEvent {
1976 completed_at: None,
1977 thread_id: "thread-1".to_string(),
1978 session_id: "session-1".to_string(),
1979 subtype: ThreadCompletionSubtype::ErrorMaxBudgetUsd,
1980 outcome_code: "budget_limit_reached".to_string(),
1981 result: None,
1982 stop_reason: Some("max_tokens".to_string()),
1983 usage: Usage {
1984 input_tokens: 10,
1985 cached_input_tokens: 4,
1986 cache_creation_tokens: 2,
1987 output_tokens: 5,
1988 },
1989 total_cost_usd: serde_json::Number::from_f64(1.25),
1990 num_turns: 3,
1991 }));
1992
1993 let json = serde_json::to_string(&event)?;
1994 let restored: ThreadEvent = serde_json::from_str(&json)?;
1995
1996 assert_eq!(restored, event);
1997 Ok(())
1998 }
1999
2000 #[test]
2001 fn compact_boundary_round_trip() -> Result<(), Box<dyn Error>> {
2002 let event = ThreadEvent::ThreadCompactBoundary(Box::new(ThreadCompactBoundaryEvent {
2003 thread_id: "thread-1".to_string(),
2004 trigger: CompactionTrigger::Recovery,
2005 mode: CompactionMode::Provider,
2006 original_message_count: 12,
2007 compacted_message_count: 5,
2008 history_artifact_path: Some("/tmp/history.jsonl".to_string()),
2009 previous_segment_id: Some("segment-0001".to_string()),
2010 new_segment_id: Some("segment-0002".to_string()),
2011 previous_prefix_hash: Some("prefix-before".to_string()),
2012 new_prefix_hash: Some("prefix-after".to_string()),
2013 previous_catalog_hash: Some("catalog-before".to_string()),
2014 new_catalog_hash: Some("catalog-after".to_string()),
2015 }));
2016
2017 let json = serde_json::to_string(&event)?;
2018 let restored: ThreadEvent = serde_json::from_str(&json)?;
2019
2020 assert_eq!(restored, event);
2021 Ok(())
2022 }
2023
2024 #[test]
2025 fn compact_boundary_deserializes_legacy_payload_without_segment_metadata() -> Result<(), Box<dyn Error>> {
2026 let payload = r#"{
2027 "type":"thread.compact_boundary",
2028 "thread_id":"thread-1",
2029 "trigger":"recovery",
2030 "mode":"provider",
2031 "original_message_count":12,
2032 "compacted_message_count":5
2033 }"#;
2034
2035 let restored: ThreadEvent = serde_json::from_str(payload)?;
2036 let ThreadEvent::ThreadCompactBoundary(event) = restored else {
2037 panic!("expected thread.compact_boundary event");
2038 };
2039
2040 assert_eq!(event.thread_id, "thread-1");
2041 assert_eq!(event.history_artifact_path, None);
2042 assert_eq!(event.previous_segment_id, None);
2043 assert_eq!(event.new_segment_id, None);
2044 assert_eq!(event.previous_prefix_hash, None);
2045 assert_eq!(event.new_prefix_hash, None);
2046 assert_eq!(event.previous_catalog_hash, None);
2047 assert_eq!(event.new_catalog_hash, None);
2048 Ok(())
2049 }
2050}