Skip to main content

vtcode_exec_events/
lib.rs

1#![allow(
2    missing_docs,
3    dead_code,
4    unused_imports,
5    reason = "Intentional compatibility, platform, or test-only suppression."
6)]
7//! Structured execution telemetry events shared across VT Code crates.
8//!
9//! This crate exposes the serialized schema for thread lifecycle updates,
10//! command execution results, and other timeline artifacts emitted by the
11//! automation runtime. Downstream applications can deserialize these
12//! structures to drive dashboards, logging, or auditing pipelines without
13//! depending on the full `vtcode-core` crate.
14//!
15//! # Agent Trace Support
16//!
17//! This crate implements the [Agent Trace](https://agent-trace.dev/) specification
18//! for tracking AI-generated code attribution. See the [`trace`] module for details.
19
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22
23pub mod atif;
24pub mod trace;
25
26/// Semantic version of the serialized event schema exported by this crate.
27pub const EVENT_SCHEMA_VERSION: &str = "0.14.0";
28
29/// Wraps a [`ThreadEvent`] with schema metadata so downstream consumers can
30/// negotiate compatibility before processing an event stream.
31#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
32#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
33pub struct VersionedThreadEvent {
34    /// Semantic version describing the schema of the nested event payload.
35    schema_version: String,
36    /// Concrete event emitted by the agent runtime.
37    event: ThreadEvent,
38}
39
40impl VersionedThreadEvent {
41    /// Creates a new [`VersionedThreadEvent`] using the current
42    /// [`EVENT_SCHEMA_VERSION`].
43    pub fn new(event: ThreadEvent) -> Self {
44        Self {
45            schema_version: EVENT_SCHEMA_VERSION.to_string(),
46            event,
47        }
48    }
49
50    /// Returns the nested [`ThreadEvent`], consuming the wrapper.
51    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
62/// Sink for processing [`ThreadEvent`] instances.
63pub trait EventEmitter {
64    /// Invoked for each event emitted by the automation runtime.
65    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/// JSON helper utilities for serializing and deserializing thread events.
78#[cfg(feature = "serde-json")]
79pub(crate) mod json {
80    use super::{ThreadEvent, VersionedThreadEvent};
81
82    /// Converts an event into a `serde_json::Value`.
83    pub fn to_value(event: &ThreadEvent) -> serde_json::Result<serde_json::Value> {
84        serde_json::to_value(event)
85    }
86
87    /// Serializes an event into a JSON string.
88    pub(crate) fn to_string(event: &ThreadEvent) -> serde_json::Result<String> {
89        serde_json::to_string(event)
90    }
91
92    /// Deserializes an event from a JSON string.
93    pub fn from_str(payload: &str) -> serde_json::Result<ThreadEvent> {
94        serde_json::from_str(payload)
95    }
96
97    /// Serializes a [`VersionedThreadEvent`] wrapper.
98    pub(crate) fn versioned_to_string(event: &ThreadEvent) -> serde_json::Result<String> {
99        serde_json::to_string(&VersionedThreadEvent::new(event.clone()))
100    }
101
102    /// Deserializes a [`VersionedThreadEvent`] wrapper.
103    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    /// Emits JSON serialized events to the `log` facade at the configured level.
115    #[derive(Debug, Clone)]
116    pub struct LogEmitter {
117        level: Level,
118    }
119
120    impl LogEmitter {
121        /// Creates a new [`LogEmitter`] that logs at the provided [`Level`].
122        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    /// Emits structured events as `tracing` events at the specified level.
157    #[derive(Debug, Clone)]
158    pub struct TracingEmitter {
159        level: Level,
160    }
161
162    impl TracingEmitter {
163        /// Creates a new [`TracingEmitter`] with the provided [`Level`].
164        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    /// Emits [`ThreadEvent`]s as OpenTelemetry spans and span events.
231    ///
232    /// Each `ThreadEvent` is recorded as an OTel span with attributes derived
233    /// from the event payload.  Harness events are attached as span events
234    /// with their own attributes (event kind, message, path, etc.).
235    ///
236    /// # Usage
237    ///
238    /// ```rust,ignore
239    /// // Requires concrete SDK type (e.g. opentelemetry_sdk::trace::SdkTracerProvider)
240    /// # use vtcode_exec_events::OtelEmitter;
241    /// # let tracer = opentelemetry_sdk::trace::SdkTracerProvider::default()
242    /// #     .tracer("vtcode");
243    /// # let mut emitter = OtelEmitter::new(tracer);
244    /// ```
245    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    /// Generates a JSON Schema describing [`ThreadEvent`].
354    pub fn thread_event_schema() -> Schema {
355        schema_for!(ThreadEvent)
356    }
357
358    /// Generates a JSON Schema describing [`VersionedThreadEvent`].
359    pub fn versioned_thread_event_schema() -> Schema {
360        schema_for!(VersionedThreadEvent)
361    }
362}
363
364/// Structured events emitted during autonomous execution.
365#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
366#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
367#[serde(tag = "type")]
368pub enum ThreadEvent {
369    /// Indicates that a new execution thread has started.
370    #[serde(rename = "thread.started")]
371    ThreadStarted(ThreadStartedEvent),
372    /// Indicates that an execution thread has reached a terminal outcome.
373    #[serde(rename = "thread.completed")]
374    ThreadCompleted(Box<ThreadCompletedEvent>),
375    /// Indicates that conversation compaction replaced older history with a boundary.
376    #[serde(rename = "thread.compact_boundary")]
377    ThreadCompactBoundary(Box<ThreadCompactBoundaryEvent>),
378    /// Indicates that the approved plan handoff rebuilt a fresh execution context.
379    #[serde(rename = "context.reset")]
380    ContextReset(ContextResetEvent),
381    /// Marks the beginning of an execution turn.
382    #[serde(rename = "turn.started")]
383    TurnStarted(TurnStartedEvent),
384    /// Marks the completion of an execution turn.
385    #[serde(rename = "turn.completed")]
386    TurnCompleted(TurnCompletedEvent),
387    /// Marks a turn as failed with additional context.
388    #[serde(rename = "turn.failed")]
389    TurnFailed(TurnFailedEvent),
390    /// Marks a turn as blocked before success could be confirmed. Emitted
391    /// alongside `turn.failed` so UI subscribers get a first-class signal
392    /// with the fuse counters and last tool instead of inferring it.
393    #[serde(rename = "turn.blocked")]
394    TurnBlocked(Box<TurnBlockedEvent>),
395    /// Indicates that an item has started processing.
396    #[serde(rename = "item.started")]
397    ItemStarted(ItemStartedEvent),
398    /// Indicates that an item has been updated.
399    #[serde(rename = "item.updated")]
400    ItemUpdated(ItemUpdatedEvent),
401    /// Indicates that an item reached a terminal state.
402    #[serde(rename = "item.completed")]
403    ItemCompleted(ItemCompletedEvent),
404    /// Emitted when a tool requires user permission before execution.
405    #[serde(rename = "permission.requested")]
406    PermissionRequested(PermissionRequestedEvent),
407    /// Emitted when the user resolves a permission prompt.
408    #[serde(rename = "permission.resolved")]
409    PermissionResolved(PermissionResolvedEvent),
410    /// A mid-turn user interjection was merged into the running turn.
411    #[serde(rename = "interjected")]
412    Interjected(InterjectedEvent),
413    /// Streaming delta for a plan item in Planning workflow.
414    #[serde(rename = "plan.delta")]
415    PlanDelta(Box<PlanDeltaEvent>),
416    /// Indicates that a completed plan is waiting for an implementation decision.
417    #[serde(rename = "plan.approval.requested")]
418    PlanApprovalRequested(PlanApprovalRequestedEvent),
419    /// Records the user's or policy's decision about a completed plan.
420    #[serde(rename = "plan.approval.resolved")]
421    PlanApprovalResolved(PlanApprovalResolvedEvent),
422    /// Represents a fatal error.
423    #[serde(rename = "error")]
424    Error(ThreadErrorEvent),
425    /// Catch-all for unknown event types added in newer schema versions.
426    /// Preserves forward compatibility when older binaries read newer event streams.
427    #[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    /// Unique identifier for the thread that was started.
435    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    /// Catch-all for unknown completion subtypes added in newer schema versions.
448    #[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    /// Compaction triggered by a mid-session switch of the main model or
477    /// provider, so the newly selected model starts from a clean summary.
478    ModelSwitch,
479    /// Catch-all for unknown triggers added in newer schema versions.
480    #[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    /// Catch-all for unknown modes added in newer schema versions.
503    #[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    /// Stable thread identifier for the session.
521    pub thread_id: String,
522    /// Stable session identifier for the runtime that produced the thread.
523    pub session_id: String,
524    /// Coarse result category aligned with SDK-style terminal states.
525    pub subtype: ThreadCompletionSubtype,
526    /// VT Code-specific detailed outcome code.
527    pub outcome_code: String,
528    /// Final assistant result text when the thread completed successfully.
529    #[serde(skip_serializing_if = "Option::is_none")]
530    pub result: Option<String>,
531    /// Provider stop reason or VT Code terminal reason when available.
532    #[serde(skip_serializing_if = "Option::is_none")]
533    pub stop_reason: Option<String>,
534    /// Aggregated token usage across the thread.
535    pub usage: Usage,
536    /// Optional estimated total API cost for the thread.
537    #[serde(skip_serializing_if = "Option::is_none")]
538    pub total_cost_usd: Option<serde_json::Number>,
539    /// Number of turns executed before completion.
540    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    /// Stable thread identifier for the session.
547    pub thread_id: String,
548    /// Whether compaction was triggered manually or automatically.
549    pub trigger: CompactionTrigger,
550    /// Whether the compaction boundary came from provider-native or local compaction.
551    pub mode: CompactionMode,
552    /// Number of messages before compaction.
553    pub original_message_count: usize,
554    /// Number of messages after compaction.
555    pub compacted_message_count: usize,
556    /// Optional persisted artifact containing the archived compaction summary/history.
557    #[serde(skip_serializing_if = "Option::is_none")]
558    pub history_artifact_path: Option<String>,
559    /// Segment identifier that contained the request prefix before compaction.
560    #[serde(skip_serializing_if = "Option::is_none")]
561    pub previous_segment_id: Option<String>,
562    /// Segment identifier created after compaction.
563    #[serde(skip_serializing_if = "Option::is_none")]
564    pub new_segment_id: Option<String>,
565    /// Hash of the immutable request prefix before compaction.
566    #[serde(skip_serializing_if = "Option::is_none")]
567    pub previous_prefix_hash: Option<String>,
568    /// Hash of the immutable request prefix after compaction.
569    #[serde(skip_serializing_if = "Option::is_none")]
570    pub new_prefix_hash: Option<String>,
571    /// Hash of the ordered tool catalog before compaction.
572    #[serde(skip_serializing_if = "Option::is_none")]
573    pub previous_catalog_hash: Option<String>,
574    /// Hash of the ordered tool catalog after compaction.
575    #[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    /// The user selected the fresh-context plan approval path.
584    PlanApproval,
585    /// Catch-all for triggers introduced by newer schema versions.
586    #[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    /// Stable thread identifier for the session.
594    pub thread_id: String,
595    /// Identifier of the turn that approved the plan.
596    pub turn_id: String,
597    /// What initiated the context reset.
598    pub trigger: ContextResetTrigger,
599    /// Whether the approved plan and task tracker survived the reset.
600    pub plan_preserved: bool,
601    /// Context pressure reported before the reset, expressed as a percentage.
602    pub previous_context_usage_percent: u8,
603    /// Whether the per-turn and per-session tool budgets were reset.
604    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    /// Optional decomposition of the assembled first-request prefix so
611    /// downstream consumers can attribute token overhead without inventing
612    /// parallel event types.
613    #[serde(skip_serializing_if = "Option::is_none")]
614    token_breakdown: Option<TokenBreakdown>,
615}
616
617/// Per-request token-budget breakdown for the assembled first-request prefix.
618#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
619#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
620pub struct TokenBreakdown {
621    /// System prompt text tokens.
622    system_prompt_tokens: u64,
623    /// On-wire tool schema tokens.
624    tool_schema_tokens: u64,
625    /// Instruction file tokens included in the prompt.
626    instruction_file_tokens: u64,
627    /// Message history text tokens.
628    message_history_tokens: u64,
629    /// Cache read tokens (served from prior turns).
630    cache_read_tokens: u64,
631    /// Cache write tokens (new cache entries created this turn).
632    cache_write_tokens: u64,
633    /// Tokens that missed cache (neither read nor written).
634    cache_miss_tokens: u64,
635    /// Subagent bootstrap tokens, if this turn spawned a child agent.
636    #[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    /// Token usage summary for the completed turn.
644    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    /// Human-readable explanation describing why the turn failed.
651    pub message: String,
652    /// Optional token usage that was consumed before the failure occurred.
653    #[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    /// Human-readable explanation describing why the turn was blocked.
661    pub message: String,
662    /// Display label of the last blocked tool call, when known.
663    #[serde(skip_serializing_if = "Option::is_none")]
664    pub last_tool: Option<String>,
665    /// Consecutive blocked tool calls observed this turn.
666    #[serde(default)]
667    pub blocked_streak: usize,
668    /// Total blocked tool calls observed this turn.
669    #[serde(default)]
670    pub blocked_total: usize,
671    /// Consecutive cap that was enforced.
672    #[serde(default)]
673    pub consecutive_cap: usize,
674    /// Total cap that was enforced.
675    #[serde(default)]
676    pub total_cap: usize,
677    /// Whether the fuse tripped while a tool-free recovery pass was active.
678    #[serde(default)]
679    pub recovery_active: bool,
680    /// Optional token usage that was consumed before the block occurred.
681    #[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    /// Fatal error message associated with the thread.
689    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    /// Number of prompt tokens processed during the turn.
696    #[serde(default, deserialize_with = "deserialize_null_as_default")]
697    pub input_tokens: u64,
698    /// Number of cached prompt tokens reused from previous turns.
699    #[serde(default, deserialize_with = "deserialize_null_as_default")]
700    pub cached_input_tokens: u64,
701    /// Number of cache-creation tokens charged during the turn.
702    #[serde(default, deserialize_with = "deserialize_null_as_default")]
703    pub cache_creation_tokens: u64,
704    /// Number of completion tokens generated by the model.
705    #[serde(default, deserialize_with = "deserialize_null_as_default")]
706    pub output_tokens: u64,
707}
708
709/// Serde helper that accepts explicit `null` as `T::default()` for
710/// backward-compatible checkpoint/diagnostics payloads. Pair with
711/// `#[serde(default, deserialize_with = "deserialize_null_as_default")]` so
712/// both missing and `null` fields degrade to the default instead of failing
713/// deserialization. Reused by downstream crates (e.g. `vtcode-core`
714/// snapshots) so the null-tolerance rule cannot drift between copies.
715pub fn deserialize_null_as_default<'de, D, T>(deserializer: D) -> Result<T, D::Error>
716where
717    D: serde::Deserializer<'de>,
718    T: Deserialize<'de> + Default,
719{
720    Ok(Option::<T>::deserialize(deserializer)?.unwrap_or_default())
721}
722
723impl Usage {
724    /// Number of input tokens billed at the full input rate: neither served
725    /// from cache nor written to it. `input_tokens` is the total prompt token
726    /// count (uncached + cached + cache-creation), so both cached and
727    /// cache-creation tokens are subtracted out here.
728    #[must_use]
729    fn uncached_input_tokens(&self) -> u64 {
730        self.input_tokens
731            .saturating_sub(self.cached_input_tokens)
732            .saturating_sub(self.cache_creation_tokens)
733    }
734
735    /// Cache hit rate as a fraction (0.0 to 1.0): cached input over total input.
736    /// Returns `None` when no input tokens were recorded.
737    #[must_use]
738    pub fn cache_hit_rate(&self) -> Option<f64> {
739        if self.input_tokens == 0 {
740            return None;
741        }
742        Some(self.cached_input_tokens as f64 / self.input_tokens as f64)
743    }
744
745    /// Human-readable summary of prompt cache efficiency.
746    #[must_use]
747    pub fn cache_summary(&self) -> String {
748        let total_input = self.input_tokens;
749        if total_input == 0 {
750            return "No input tokens recorded.".to_string();
751        }
752
753        let cached = self.cached_input_tokens;
754        let creation = self.cache_creation_tokens;
755        let uncached = self.uncached_input_tokens();
756        let rate = cached as f64 / total_input as f64 * 100.0;
757        format!(
758            "Cache: {cached} cached / {total_input} total input ({rate:.1}% hit rate), \
759             {creation} cache-creation, {uncached} uncached"
760        )
761    }
762
763    /// Accumulate another usage sample into this one.
764    pub fn add(&mut self, other: &Usage) {
765        self.input_tokens = self.input_tokens.saturating_add(other.input_tokens);
766        self.cached_input_tokens = self.cached_input_tokens.saturating_add(other.cached_input_tokens);
767        self.cache_creation_tokens = self.cache_creation_tokens.saturating_add(other.cache_creation_tokens);
768        self.output_tokens = self.output_tokens.saturating_add(other.output_tokens);
769    }
770}
771
772#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
773#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
774pub struct ItemCompletedEvent {
775    /// Snapshot of the thread item that completed.
776    pub item: ThreadItem,
777}
778
779#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
780#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
781pub struct ItemStartedEvent {
782    /// Snapshot of the thread item that began processing.
783    pub item: ThreadItem,
784}
785
786#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
787#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
788pub struct ItemUpdatedEvent {
789    /// Snapshot of the thread item after it was updated.
790    pub item: ThreadItem,
791}
792
793#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
794#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
795pub struct PlanDeltaEvent {
796    /// Identifier of the thread emitting this plan delta.
797    pub thread_id: String,
798    /// Identifier of the current turn.
799    pub turn_id: String,
800    /// Identifier of the plan item receiving the delta.
801    pub item_id: String,
802    /// Incremental plan text chunk.
803    pub delta: String,
804}
805
806#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
807#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
808pub struct PlanApprovalRequestedEvent {
809    /// Identifier of the thread emitting the approval request.
810    pub thread_id: String,
811    /// Identifier of the turn that produced the plan.
812    pub turn_id: String,
813    /// Plan file associated with the approval request, when available.
814    #[serde(skip_serializing_if = "Option::is_none")]
815    pub plan_file: Option<String>,
816}
817
818#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
819#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
820#[serde(rename_all = "snake_case")]
821pub enum PlanApprovalDecision {
822    /// Execute with normal per-edit approval prompts.
823    Execute,
824    /// Execute with automatic edit approval enabled.
825    AutoAccept,
826    /// Execute the plan after rebuilding a fresh context.
827    FreshContext,
828    /// Keep planning and revise the proposed plan.
829    Revise,
830    /// Dismiss the approval request without implementing.
831    Cancel,
832    /// Hand the plan to the build primary agent.
833    SwitchBuild,
834    /// Hand the plan to the auto primary agent.
835    SwitchAuto,
836    /// Catch-all for decisions added in newer schema versions.
837    #[serde(other)]
838    Unknown,
839}
840
841#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
842#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
843pub struct PlanApprovalResolvedEvent {
844    /// Identifier of the thread emitting the approval decision.
845    pub thread_id: String,
846    /// Identifier of the turn in which the decision was made.
847    pub turn_id: String,
848    /// Decision selected by the user or active execution policy.
849    pub decision: PlanApprovalDecision,
850    /// Whether the decision came from policy rather than an interactive user action.
851    pub automatic: bool,
852}
853
854#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
855#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
856pub struct ThreadItem {
857    /// Stable identifier associated with the item.
858    pub id: String,
859    /// Embedded event details for the item type.
860    #[serde(flatten)]
861    pub details: ThreadItemDetails,
862}
863
864#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
865#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
866#[serde(tag = "type", rename_all = "snake_case")]
867pub enum ThreadItemDetails {
868    /// Message authored by the agent.
869    AgentMessage(AgentMessageItem),
870    /// Structured plan content authored by the agent in Planning workflow.
871    Plan(PlanItem),
872    /// Free-form reasoning text produced during a turn.
873    Reasoning(ReasoningItem),
874    /// Command execution lifecycle update for an actual shell/PTY process.
875    CommandExecution(Box<CommandExecutionItem>),
876    /// Tool invocation lifecycle update.
877    ToolInvocation(Box<ToolInvocationItem>),
878    /// Tool output lifecycle update tied to a tool invocation.
879    ToolOutput(Box<ToolOutputItem>),
880    /// File change summary associated with the turn.
881    FileChange(Box<FileChangeItem>),
882    /// MCP tool invocation status.
883    McpToolCall(Box<McpToolCallItem>),
884    /// Web search event emitted by a registered search provider.
885    WebSearch(Box<WebSearchItem>),
886    /// Harness-managed continuation or verification lifecycle event.
887    Harness(Box<HarnessEventItem>),
888    /// General error captured for auditing.
889    Error(ErrorItem),
890}
891
892#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
893#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
894pub struct AgentMessageItem {
895    /// Textual content of the agent message.
896    pub text: String,
897}
898
899#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
900#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
901pub struct PlanItem {
902    /// Plan markdown content.
903    pub text: String,
904}
905
906#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
907#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
908pub struct ReasoningItem {
909    /// Free-form reasoning content captured during planning.
910    pub text: String,
911    /// Optional stage of reasoning (e.g., "analysis", "plan", "verification",
912    /// or the bounded evidence-only "diagnosis" stage).
913    #[serde(skip_serializing_if = "Option::is_none")]
914    pub stage: Option<String>,
915}
916
917#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
918#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
919#[serde(rename_all = "snake_case")]
920pub enum CommandExecutionStatus {
921    /// Command finished successfully.
922    #[default]
923    Completed,
924    /// Command failed (non-zero exit code or runtime error).
925    Failed,
926    /// Command is still running and may emit additional output.
927    InProgress,
928}
929
930#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
931#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
932pub struct CommandExecutionItem {
933    /// Tool or command identifier executed by the runner.
934    pub command: String,
935    /// Arguments passed to the tool invocation, when available.
936    #[serde(skip_serializing_if = "Option::is_none")]
937    pub arguments: Option<Value>,
938    /// Aggregated output emitted by the command.
939    #[serde(default)]
940    pub aggregated_output: String,
941    /// Exit code reported by the process, when available.
942    #[serde(skip_serializing_if = "Option::is_none")]
943    pub exit_code: Option<i32>,
944    /// Current status of the command execution.
945    pub status: CommandExecutionStatus,
946}
947
948#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Default)]
949#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
950#[serde(rename_all = "snake_case")]
951pub enum ToolCallStatus {
952    /// Tool finished successfully.
953    #[default]
954    Completed,
955    /// Tool failed.
956    Failed,
957    /// Tool is still running and may emit additional output.
958    InProgress,
959}
960
961/// Fine-grained outcome of a tool invocation lifecycle.
962///
963/// Mirrors the outcome taxonomy used by the runtime: `status` remains the
964/// coarse lifecycle signal (`Completed` / `Failed` / `InProgress`), while
965/// `outcome` captures *why* the invocation terminated. Consumers that only
966/// need success/failure can continue to read `status`; analytics and the UI
967/// layer use `outcome` for richer classification.
968#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, Default)]
969#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
970#[serde(rename_all = "snake_case")]
971pub enum ToolOutcome {
972    /// Tool executed and returned a result.
973    #[default]
974    Success,
975    /// Tool executed but returned an error.
976    Error,
977    /// User rejected the permission prompt.
978    PermissionRejected,
979    /// User cancelled the permission prompt (e.g. Ctrl+C / Esc).
980    PermissionCancelled,
981    /// User provided a followup message instead of approving.
982    Followup,
983    /// A user-configured hook blocked execution.
984    HookDenied,
985    /// Tool not found or arguments couldn't be parsed.
986    InvalidTool,
987    /// Tool was running when the turn was cancelled.
988    Cancelled,
989}
990
991impl ToolOutcome {
992    #[must_use]
993    pub const fn is_terminal(self) -> bool {
994        !matches!(self, Self::Followup)
995    }
996}
997
998/// Map a terminal [`ToolCallStatus`] to its corresponding [`ToolOutcome`].
999///
1000/// # Panics
1001///
1002/// Panics if `status` is [`ToolCallStatus::InProgress`], which is a non-terminal
1003/// state and must never be passed to a completion-event emitter.
1004#[must_use]
1005#[allow(
1006    clippy::unreachable,
1007    reason = "Intentional compatibility, platform, or test-only suppression."
1008)]
1009pub fn tool_outcome_from_status(status: &ToolCallStatus) -> ToolOutcome {
1010    match status {
1011        ToolCallStatus::Completed => ToolOutcome::Success,
1012        ToolCallStatus::Failed => ToolOutcome::Error,
1013        ToolCallStatus::InProgress => unreachable!("InProgress status passed to completion event"),
1014    }
1015}
1016
1017#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1018#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1019pub struct ToolInvocationItem {
1020    /// Name of the invoked tool.
1021    pub tool_name: String,
1022    /// Structured arguments passed to the tool.
1023    #[serde(skip_serializing_if = "Option::is_none")]
1024    pub arguments: Option<Value>,
1025    /// Raw model-emitted tool call identifier, when available.
1026    #[serde(skip_serializing_if = "Option::is_none")]
1027    pub tool_call_id: Option<String>,
1028    /// Current lifecycle status of the invocation.
1029    pub status: ToolCallStatus,
1030    /// Fine-grained outcome of the invocation lifecycle.
1031    #[serde(skip_serializing_if = "Option::is_none")]
1032    pub outcome: Option<ToolOutcome>,
1033}
1034
1035#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1036#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1037pub struct ToolOutputItem {
1038    /// Identifier of the related harness invocation item.
1039    pub call_id: String,
1040    /// Raw model-emitted tool call identifier, when available.
1041    #[serde(skip_serializing_if = "Option::is_none")]
1042    pub tool_call_id: Option<String>,
1043    /// Canonical spool file path when the full output was written to disk.
1044    #[serde(skip_serializing_if = "Option::is_none")]
1045    pub spool_path: Option<String>,
1046    /// Aggregated output emitted by the tool.
1047    #[serde(default)]
1048    pub output: String,
1049    /// Exit code reported by the tool, when available.
1050    #[serde(skip_serializing_if = "Option::is_none")]
1051    pub exit_code: Option<i32>,
1052    /// Current lifecycle status of the output item.
1053    pub status: ToolCallStatus,
1054}
1055
1056#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1057#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1058pub struct FileChangeItem {
1059    /// List of individual file updates included in the change set.
1060    pub changes: Vec<FileUpdateChange>,
1061    /// Whether the patch application succeeded.
1062    pub status: PatchApplyStatus,
1063    /// Optional precomputed unified diff for the change set.
1064    ///
1065    /// Populated by the turn diff tracker so consumers can render per-change
1066    /// previews without recomputation. Absent in older events.
1067    #[serde(default, skip_serializing_if = "Option::is_none")]
1068    pub unified_diff: Option<String>,
1069    /// Optional added-line count for the change set.
1070    #[serde(default, skip_serializing_if = "Option::is_none")]
1071    pub additions: Option<u64>,
1072    /// Optional deleted-line count for the change set.
1073    #[serde(default, skip_serializing_if = "Option::is_none")]
1074    pub deletions: Option<u64>,
1075}
1076
1077#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1078#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1079pub struct FileUpdateChange {
1080    /// Path of the file that was updated.
1081    pub path: String,
1082    /// Type of change applied to the file.
1083    pub kind: PatchChangeKind,
1084}
1085
1086#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1087#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1088#[serde(rename_all = "snake_case")]
1089pub enum PatchApplyStatus {
1090    /// Patch successfully applied.
1091    Completed,
1092    /// Patch application failed.
1093    Failed,
1094}
1095
1096#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1097#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1098#[serde(rename_all = "snake_case")]
1099pub enum PatchChangeKind {
1100    /// File addition.
1101    Add,
1102    /// File deletion.
1103    Delete,
1104    /// File update in place.
1105    Update,
1106}
1107
1108#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1109#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1110pub struct McpToolCallItem {
1111    /// Name of the MCP tool invoked by the agent.
1112    pub tool_name: String,
1113    /// Arguments passed to the tool invocation, if any.
1114    #[serde(skip_serializing_if = "Option::is_none")]
1115    pub arguments: Option<Value>,
1116    /// Result payload returned by the tool, if captured.
1117    #[serde(skip_serializing_if = "Option::is_none")]
1118    pub result: Option<String>,
1119    /// Lifecycle status for the tool call.
1120    #[serde(skip_serializing_if = "Option::is_none")]
1121    pub status: Option<McpToolCallStatus>,
1122}
1123
1124#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1125#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1126#[serde(rename_all = "snake_case")]
1127pub enum McpToolCallStatus {
1128    /// Tool invocation has started.
1129    Started,
1130    /// Tool invocation completed successfully.
1131    Completed,
1132    /// Tool invocation failed.
1133    Failed,
1134}
1135
1136#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1137#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1138pub struct WebSearchItem {
1139    /// Query that triggered the search.
1140    pub query: String,
1141    /// Search provider identifier, when known.
1142    #[serde(skip_serializing_if = "Option::is_none")]
1143    pub provider: Option<String>,
1144    /// Optional raw search results captured for auditing.
1145    #[serde(skip_serializing_if = "Option::is_none")]
1146    pub results: Option<Vec<String>>,
1147}
1148
1149#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1150#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1151#[serde(rename_all = "snake_case")]
1152pub enum HarnessEventKind {
1153    PlanningStarted,
1154    PlanningCompleted,
1155    ContinuationStarted,
1156    ContinuationSkipped,
1157    /// A turn was blocked before success could be confirmed. Carries the fuse
1158    /// counters so UI layers can render without correlating multiple events.
1159    TurnBlocked,
1160    /// A bounded tool-free recovery pass was scheduled after blocked calls.
1161    BlockedRecoveryStarted,
1162    /// A bounded tool-free recovery pass finished.
1163    BlockedRecoveryFinished,
1164    BlockedHandoffWritten,
1165    /// The owning session resolved its archived blocked handoff and removed
1166    /// the live recovery pointer.
1167    BlockedHandoffResolved,
1168    EvaluationStarted,
1169    EvaluationPassed,
1170    EvaluationFailed,
1171    RevisionStarted,
1172    EscalationTriggered,
1173    EscalationBypassed,
1174    VerificationStarted,
1175    VerificationPassed,
1176    VerificationFailed,
1177    /// Agent recovered from a transient error (e.g. after retry succeeded).
1178    ErrorRecovered,
1179    /// A transient tool failure triggered an automatic retry attempt.
1180    ToolRetryAttempted,
1181    /// Latency record for a tool execution, emitted on turn completion.
1182    ToolLatencyRecorded,
1183    /// A checkpoint snapshot was created for the current turn.
1184    SnapshotCreated,
1185    /// A checkpoint snapshot was restored (rewind operation).
1186    SnapshotRestored,
1187    /// The user granted additional session tool-call capacity and the
1188    /// pending call will be retried in the same turn.
1189    SessionToolLimitIncreased,
1190    /// The user granted additional tool-loop capacity for the current turn.
1191    ToolLoopLimitIncreased,
1192}
1193
1194#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1195#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1196#[serde(rename_all = "snake_case")]
1197pub enum PermissionDecision {
1198    Allow,
1199    Deny,
1200    Cancelled,
1201    Followup,
1202}
1203
1204#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1205#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1206pub struct PermissionRequestedEvent {
1207    /// Name of the tool that requires permission.
1208    pub tool_name: String,
1209}
1210
1211#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1212#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1213pub struct PermissionResolvedEvent {
1214    /// Name of the tool that was permitted or denied.
1215    pub tool_name: String,
1216    /// User's decision on the permission prompt.
1217    pub decision: PermissionDecision,
1218    /// Wall-clock time the prompt was visible, in milliseconds.
1219    pub wait_ms: u64,
1220}
1221
1222#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1223#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1224#[serde(rename_all = "snake_case")]
1225pub enum InterjectionSource {
1226    Direct,
1227    Queue,
1228}
1229
1230#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
1231#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1232#[serde(rename_all = "snake_case")]
1233pub enum RedirectKind {
1234    Interjection,
1235}
1236
1237#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1238#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1239pub struct InterjectedEvent {
1240    /// How the interjection reached the running turn.
1241    pub source: InterjectionSource,
1242    /// Number of image attachments that accompanied the interjection.
1243    pub image_count: u32,
1244    /// Always `Interjection` for this event; carried so the shared
1245    /// `redirect_kind` field is queryable uniformly across redirect events.
1246    pub redirect_kind: RedirectKind,
1247}
1248
1249#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1250#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1251pub struct HarnessEventItem {
1252    /// Specific harness event emitted by the runtime.
1253    pub event: HarnessEventKind,
1254    /// Optional human-readable message associated with the event.
1255    #[serde(skip_serializing_if = "Option::is_none")]
1256    pub message: Option<String>,
1257    /// Optional verification command associated with the event.
1258    #[serde(skip_serializing_if = "Option::is_none")]
1259    pub command: Option<String>,
1260    /// Optional artifact path associated with the event.
1261    #[serde(skip_serializing_if = "Option::is_none")]
1262    pub path: Option<String>,
1263    /// Optional exit code associated with verification results.
1264    #[serde(skip_serializing_if = "Option::is_none")]
1265    pub exit_code: Option<i32>,
1266    /// Retry/recovery attempt number (1-indexed). Only set for retry-related events.
1267    #[serde(skip_serializing_if = "Option::is_none")]
1268    pub attempt: Option<u32>,
1269    /// Canonical error category for retry/recovery events.
1270    #[serde(skip_serializing_if = "Option::is_none")]
1271    pub error_category: Option<String>,
1272    /// Latency in milliseconds for tool-execution latency events.
1273    #[serde(skip_serializing_if = "Option::is_none")]
1274    pub duration_ms: Option<u64>,
1275}
1276
1277#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1278#[cfg_attr(feature = "schema-export", derive(schemars::JsonSchema))]
1279pub struct ErrorItem {
1280    /// Error message displayed to the user or logs.
1281    pub message: String,
1282}
1283
1284#[cfg(test)]
1285mod tests {
1286    use super::*;
1287    use std::error::Error;
1288    use std::mem::size_of;
1289
1290    /// `ThreadEvent` is pushed into `Vec`s per streaming delta and accumulated
1291    /// for whole sessions. Large sparse payloads must stay boxed so the enum
1292    /// does not balloon from alignment/discriminant padding (see
1293    /// docs/development/rust-performance-principles.md, "Enum footprint").
1294    #[test]
1295    fn thread_event_stays_compact() {
1296        assert!(
1297            size_of::<ThreadEvent>() <= 80,
1298            "ThreadEvent grew to {} bytes; box new large payloads instead of inlining them",
1299            size_of::<ThreadEvent>()
1300        );
1301    }
1302
1303    /// Boxing only pays off while the inline (unboxed) payload is larger than
1304    /// a pointer. Guard each boxed variant against accidental unboxing.
1305    #[test]
1306    fn boxed_thread_item_details_payloads_stay_boxed() {
1307        assert!(size_of::<Option<Box<CommandExecutionItem>>>() < size_of::<Option<CommandExecutionItem>>());
1308        assert!(size_of::<Option<Box<ToolInvocationItem>>>() < size_of::<Option<ToolInvocationItem>>());
1309        assert!(size_of::<Option<Box<ToolOutputItem>>>() < size_of::<Option<ToolOutputItem>>());
1310        assert!(size_of::<Option<Box<FileChangeItem>>>() < size_of::<Option<FileChangeItem>>());
1311        assert!(size_of::<Option<Box<McpToolCallItem>>>() < size_of::<Option<McpToolCallItem>>());
1312        assert!(size_of::<Option<Box<WebSearchItem>>>() < size_of::<Option<WebSearchItem>>());
1313        assert!(size_of::<Option<Box<HarnessEventItem>>>() < size_of::<Option<HarnessEventItem>>());
1314    }
1315
1316    #[test]
1317    fn file_change_item_optional_diff_fields_round_trip() -> Result<(), Box<dyn Error>> {
1318        // Legacy payload without the new optional fields must deserialize.
1319        let legacy_json = r#"{
1320            "changes": [{"path": "src/main.rs", "kind": "add"}],
1321            "status": "completed"
1322        }"#;
1323        let legacy: FileChangeItem = serde_json::from_str(legacy_json)?;
1324        assert!(legacy.unified_diff.is_none());
1325        assert!(legacy.additions.is_none());
1326        assert!(legacy.deletions.is_none());
1327
1328        // New fields are omitted from output when unset.
1329        let legacy_reserialized = serde_json::to_value(&legacy)?;
1330        assert!(legacy_reserialized.get("unified_diff").is_none());
1331        assert!(legacy_reserialized.get("additions").is_none());
1332        assert!(legacy_reserialized.get("deletions").is_none());
1333
1334        // Populated fields survive a round trip.
1335        let populated = FileChangeItem {
1336            changes: legacy.changes.clone(),
1337            status: PatchApplyStatus::Completed,
1338            unified_diff: Some("diff --git a/x b/x\n".to_string()),
1339            additions: Some(3),
1340            deletions: Some(1),
1341        };
1342        let json = serde_json::to_string(&populated)?;
1343        let restored: FileChangeItem = serde_json::from_str(&json)?;
1344        assert_eq!(restored, populated);
1345        Ok(())
1346    }
1347
1348    #[test]
1349    fn thread_event_round_trip() -> Result<(), Box<dyn Error>> {
1350        let event = ThreadEvent::TurnCompleted(TurnCompletedEvent {
1351            usage: Usage {
1352                input_tokens: 1,
1353                cached_input_tokens: 2,
1354                cache_creation_tokens: 0,
1355                output_tokens: 3,
1356            },
1357        });
1358
1359        let json = serde_json::to_string(&event)?;
1360        let restored: ThreadEvent = serde_json::from_str(&json)?;
1361
1362        assert_eq!(restored, event);
1363        Ok(())
1364    }
1365
1366    #[test]
1367    fn turn_blocked_event_round_trip() -> Result<(), Box<dyn Error>> {
1368        let event = ThreadEvent::TurnBlocked(Box::new(TurnBlockedEvent {
1369            message: "Blocked tool-call limit reached after 3 consecutive blocked calls.".to_string(),
1370            last_tool: Some("exec_command".to_string()),
1371            blocked_streak: 4,
1372            blocked_total: 4,
1373            consecutive_cap: 3,
1374            total_cap: 6,
1375            recovery_active: false,
1376            usage: None,
1377        }));
1378
1379        let json = serde_json::to_string(&event)?;
1380        assert!(json.contains("turn.blocked"));
1381        let restored: ThreadEvent = serde_json::from_str(&json)?;
1382        assert_eq!(restored, event);
1383
1384        // Legacy payloads without new counters still parse via defaults.
1385        let legacy = serde_json::json!({"type": "turn.blocked", "message": "blocked"});
1386        let parsed: ThreadEvent = serde_json::from_value(legacy)?;
1387        assert!(matches!(parsed, ThreadEvent::TurnBlocked(_)));
1388        Ok(())
1389    }
1390
1391    #[test]
1392    fn usage_uncached_input_tokens_saturates() {
1393        let usage = Usage {
1394            input_tokens: 1_000,
1395            cached_input_tokens: 800,
1396            cache_creation_tokens: 100,
1397            output_tokens: 50,
1398        };
1399        assert_eq!(usage.uncached_input_tokens(), 100);
1400
1401        let inconsistent = Usage {
1402            input_tokens: 100,
1403            cached_input_tokens: 150,
1404            cache_creation_tokens: 0,
1405            output_tokens: 0,
1406        };
1407        assert_eq!(inconsistent.uncached_input_tokens(), 0);
1408
1409        let inconsistent_with_creation = Usage {
1410            input_tokens: 100,
1411            cached_input_tokens: 80,
1412            cache_creation_tokens: 50,
1413            output_tokens: 0,
1414        };
1415        assert_eq!(inconsistent_with_creation.uncached_input_tokens(), 0);
1416    }
1417
1418    #[test]
1419    fn usage_cache_hit_rate() {
1420        assert_eq!(Usage::default().cache_hit_rate(), None);
1421
1422        let usage = Usage {
1423            input_tokens: 1_000,
1424            cached_input_tokens: 750,
1425            cache_creation_tokens: 0,
1426            output_tokens: 0,
1427        };
1428        let rate = usage.cache_hit_rate().expect("rate");
1429        assert!((rate - 0.75).abs() < f64::EPSILON);
1430    }
1431
1432    #[test]
1433    fn usage_cache_summary_formats() {
1434        assert_eq!(Usage::default().cache_summary(), "No input tokens recorded.");
1435
1436        let usage = Usage {
1437            input_tokens: 1_000,
1438            cached_input_tokens: 800,
1439            cache_creation_tokens: 100,
1440            output_tokens: 50,
1441        };
1442        assert_eq!(
1443            usage.cache_summary(),
1444            "Cache: 800 cached / 1000 total input (80.0% hit rate), 100 cache-creation, 100 uncached"
1445        );
1446    }
1447
1448    #[test]
1449    fn usage_add_accumulates_all_fields_with_saturation() {
1450        let mut total = Usage {
1451            input_tokens: 100,
1452            cached_input_tokens: 20,
1453            cache_creation_tokens: 5,
1454            output_tokens: 10,
1455        };
1456        total.add(&Usage {
1457            input_tokens: 50,
1458            cached_input_tokens: 10,
1459            cache_creation_tokens: 2,
1460            output_tokens: 8,
1461        });
1462
1463        assert_eq!(total.input_tokens, 150);
1464        assert_eq!(total.cached_input_tokens, 30);
1465        assert_eq!(total.cache_creation_tokens, 7);
1466        assert_eq!(total.output_tokens, 18);
1467
1468        let mut saturating = Usage {
1469            input_tokens: u64::MAX,
1470            cached_input_tokens: u64::MAX,
1471            cache_creation_tokens: u64::MAX,
1472            output_tokens: u64::MAX,
1473        };
1474        saturating.add(&Usage {
1475            input_tokens: 1,
1476            cached_input_tokens: 1,
1477            cache_creation_tokens: 1,
1478            output_tokens: 1,
1479        });
1480        assert_eq!(saturating.input_tokens, u64::MAX);
1481        assert_eq!(saturating.cached_input_tokens, u64::MAX);
1482        assert_eq!(saturating.cache_creation_tokens, u64::MAX);
1483        assert_eq!(saturating.output_tokens, u64::MAX);
1484    }
1485
1486    #[test]
1487    fn versioned_event_wraps_schema_version() {
1488        let event = ThreadEvent::ThreadStarted(ThreadStartedEvent { thread_id: "abc".to_string() });
1489
1490        let versioned = VersionedThreadEvent::new(event.clone());
1491
1492        assert_eq!(versioned.schema_version, EVENT_SCHEMA_VERSION);
1493        assert_eq!(versioned.event, event);
1494        assert_eq!(versioned.into_event(), event);
1495    }
1496
1497    #[test]
1498    fn plan_approval_events_round_trip_with_decision() {
1499        let requested = ThreadEvent::PlanApprovalRequested(PlanApprovalRequestedEvent {
1500            thread_id: "thread-1".to_string(),
1501            turn_id: "turn-2".to_string(),
1502            plan_file: Some(".vtcode/plans/change.md".to_string()),
1503        });
1504        let resolved = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1505            thread_id: "thread-1".to_string(),
1506            turn_id: "turn-3".to_string(),
1507            decision: PlanApprovalDecision::AutoAccept,
1508            automatic: false,
1509        });
1510
1511        for event in [requested, resolved] {
1512            let serialized = serde_json::to_string(&event).expect("serialize plan approval event");
1513            let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize plan approval event");
1514            assert_eq!(restored, event);
1515        }
1516    }
1517
1518    #[test]
1519    fn context_reset_event_round_trips_with_handoff_metadata() {
1520        let event = ThreadEvent::ContextReset(ContextResetEvent {
1521            thread_id: "thread-1".to_string(),
1522            turn_id: "turn-3".to_string(),
1523            trigger: ContextResetTrigger::PlanApproval,
1524            plan_preserved: true,
1525            previous_context_usage_percent: 7,
1526            tool_budget_reset: true,
1527        });
1528
1529        let serialized = serde_json::to_string(&event).expect("serialize context reset event");
1530        let restored: ThreadEvent = serde_json::from_str(&serialized).expect("deserialize context reset event");
1531        assert_eq!(restored, event);
1532        assert_eq!(serde_json::to_value(event).expect("wire value")["type"], "context.reset");
1533    }
1534
1535    #[test]
1536    fn plan_approval_decision_uses_stable_wire_names() {
1537        let event = ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1538            thread_id: "thread-1".to_string(),
1539            turn_id: "turn-1".to_string(),
1540            decision: PlanApprovalDecision::SwitchBuild,
1541            automatic: false,
1542        });
1543
1544        let serialized = serde_json::to_value(event).expect("serialize plan approval decision");
1545        assert_eq!(serialized["type"], "plan.approval.resolved");
1546        assert_eq!(serialized["decision"], "switch_build");
1547    }
1548
1549    #[test]
1550    fn plan_approval_decision_is_forward_compatible() {
1551        let payload = serde_json::json!({
1552            "type": "plan.approval.resolved",
1553            "thread_id": "thread-1",
1554            "turn_id": "turn-1",
1555            "decision": "future_decision",
1556            "automatic": true,
1557        });
1558        let event: ThreadEvent = serde_json::from_value(payload).expect("future decision should deserialize");
1559        assert!(matches!(
1560            event,
1561            ThreadEvent::PlanApprovalResolved(PlanApprovalResolvedEvent {
1562                decision: PlanApprovalDecision::Unknown,
1563                automatic: true,
1564                ..
1565            })
1566        ));
1567    }
1568
1569    #[cfg(feature = "serde-json")]
1570    #[test]
1571    fn versioned_json_round_trip() -> Result<(), Box<dyn Error>> {
1572        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1573            item: ThreadItem {
1574                id: "item-1".to_string(),
1575                details: ThreadItemDetails::AgentMessage(AgentMessageItem { text: "hello".to_string() }),
1576            },
1577        });
1578
1579        let payload = json::versioned_to_string(&event)?;
1580        let restored = json::versioned_from_str(&payload)?;
1581
1582        assert_eq!(restored.schema_version, EVENT_SCHEMA_VERSION);
1583        assert_eq!(restored.event, event);
1584        Ok(())
1585    }
1586
1587    #[test]
1588    fn compaction_trigger_serializes_snake_case_and_round_trips() {
1589        for trigger in [
1590            CompactionTrigger::Manual,
1591            CompactionTrigger::Auto,
1592            CompactionTrigger::Recovery,
1593            CompactionTrigger::ModelSwitch,
1594            CompactionTrigger::Unknown,
1595        ] {
1596            let json = serde_json::to_string(&trigger).unwrap();
1597            assert_eq!(json, format!("\"{}\"", trigger.as_str()));
1598            let restored: CompactionTrigger = serde_json::from_str(&json).unwrap();
1599            assert_eq!(restored, trigger);
1600        }
1601    }
1602
1603    #[test]
1604    fn tool_invocation_round_trip() -> Result<(), Box<dyn Error>> {
1605        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1606            item: ThreadItem {
1607                id: "tool_1".to_string(),
1608                details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
1609                    tool_name: "read_file".to_string(),
1610                    arguments: Some(serde_json::json!({ "path": "README.md" })),
1611                    tool_call_id: Some("tool_call_0".to_string()),
1612                    status: ToolCallStatus::Completed,
1613                    outcome: None,
1614                })),
1615            },
1616        });
1617
1618        let json = serde_json::to_string(&event)?;
1619        let restored: ThreadEvent = serde_json::from_str(&json)?;
1620
1621        assert_eq!(restored, event);
1622        Ok(())
1623    }
1624
1625    #[test]
1626    fn tool_outcome_serializes_snake_case() {
1627        for outcome in [
1628            ToolOutcome::Success,
1629            ToolOutcome::Error,
1630            ToolOutcome::PermissionRejected,
1631            ToolOutcome::PermissionCancelled,
1632            ToolOutcome::Followup,
1633            ToolOutcome::HookDenied,
1634            ToolOutcome::InvalidTool,
1635            ToolOutcome::Cancelled,
1636        ] {
1637            let json = serde_json::to_string(&outcome).unwrap();
1638            let restored: ToolOutcome = serde_json::from_str(&json).unwrap();
1639            assert_eq!(restored, outcome);
1640        }
1641    }
1642
1643    #[test]
1644    fn tool_invocation_outcome_round_trip() -> Result<(), Box<dyn Error>> {
1645        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1646            item: ThreadItem {
1647                id: "tool_1".to_string(),
1648                details: ThreadItemDetails::ToolInvocation(Box::new(ToolInvocationItem {
1649                    tool_name: "exec_command".to_string(),
1650                    arguments: Some(serde_json::json!({ "command": ["pwd"] })),
1651                    tool_call_id: Some("tool_call_0".to_string()),
1652                    status: ToolCallStatus::Failed,
1653                    outcome: Some(ToolOutcome::PermissionRejected),
1654                })),
1655            },
1656        });
1657
1658        let json = serde_json::to_string(&event)?;
1659        let restored: ThreadEvent = serde_json::from_str(&json)?;
1660
1661        assert_eq!(restored, event);
1662        Ok(())
1663    }
1664
1665    #[test]
1666    fn tool_output_round_trip_preserves_raw_tool_call_id() -> Result<(), Box<dyn Error>> {
1667        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1668            item: ThreadItem {
1669                id: "tool_1:output".to_string(),
1670                details: ThreadItemDetails::ToolOutput(Box::new(ToolOutputItem {
1671                    call_id: "tool_1".to_string(),
1672                    tool_call_id: Some("tool_call_0".to_string()),
1673                    spool_path: None,
1674                    output: "done".to_string(),
1675                    exit_code: Some(0),
1676                    status: ToolCallStatus::Completed,
1677                })),
1678            },
1679        });
1680
1681        let json = serde_json::to_string(&event)?;
1682        let restored: ThreadEvent = serde_json::from_str(&json)?;
1683
1684        assert_eq!(restored, event);
1685        Ok(())
1686    }
1687
1688    #[test]
1689    fn harness_item_round_trip() -> Result<(), Box<dyn Error>> {
1690        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1691            item: ThreadItem {
1692                id: "harness_1".to_string(),
1693                details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1694                    event: HarnessEventKind::VerificationFailed,
1695                    message: Some("cargo check failed".to_string()),
1696                    command: Some("cargo check".to_string()),
1697                    path: None,
1698                    exit_code: Some(101),
1699                    attempt: None,
1700                    error_category: None,
1701                    duration_ms: None,
1702                })),
1703            },
1704        });
1705
1706        let json = serde_json::to_string(&event)?;
1707        let restored: ThreadEvent = serde_json::from_str(&json)?;
1708
1709        assert_eq!(restored, event);
1710        Ok(())
1711    }
1712
1713    #[test]
1714    fn blocked_handoff_resolved_uses_stable_wire_name() -> Result<(), Box<dyn Error>> {
1715        let event = ThreadEvent::ItemCompleted(ItemCompletedEvent {
1716            item: ThreadItem {
1717                id: "harness_resolved".to_string(),
1718                details: ThreadItemDetails::Harness(Box::new(HarnessEventItem {
1719                    event: HarnessEventKind::BlockedHandoffResolved,
1720                    message: Some("resolved".to_string()),
1721                    command: None,
1722                    path: None,
1723                    exit_code: None,
1724                    attempt: None,
1725                    error_category: None,
1726                    duration_ms: None,
1727                })),
1728            },
1729        });
1730
1731        let value = serde_json::to_value(&event)?;
1732        assert_eq!(value["item"]["event"], "blocked_handoff_resolved");
1733
1734        let restored: ThreadEvent = serde_json::from_value(value)?;
1735        assert_eq!(restored, event);
1736        Ok(())
1737    }
1738
1739    #[test]
1740    fn thread_completed_round_trip() -> Result<(), Box<dyn Error>> {
1741        let event = ThreadEvent::ThreadCompleted(Box::new(ThreadCompletedEvent {
1742            thread_id: "thread-1".to_string(),
1743            session_id: "session-1".to_string(),
1744            subtype: ThreadCompletionSubtype::ErrorMaxBudgetUsd,
1745            outcome_code: "budget_limit_reached".to_string(),
1746            result: None,
1747            stop_reason: Some("max_tokens".to_string()),
1748            usage: Usage {
1749                input_tokens: 10,
1750                cached_input_tokens: 4,
1751                cache_creation_tokens: 2,
1752                output_tokens: 5,
1753            },
1754            total_cost_usd: serde_json::Number::from_f64(1.25),
1755            num_turns: 3,
1756        }));
1757
1758        let json = serde_json::to_string(&event)?;
1759        let restored: ThreadEvent = serde_json::from_str(&json)?;
1760
1761        assert_eq!(restored, event);
1762        Ok(())
1763    }
1764
1765    #[test]
1766    fn compact_boundary_round_trip() -> Result<(), Box<dyn Error>> {
1767        let event = ThreadEvent::ThreadCompactBoundary(Box::new(ThreadCompactBoundaryEvent {
1768            thread_id: "thread-1".to_string(),
1769            trigger: CompactionTrigger::Recovery,
1770            mode: CompactionMode::Provider,
1771            original_message_count: 12,
1772            compacted_message_count: 5,
1773            history_artifact_path: Some("/tmp/history.jsonl".to_string()),
1774            previous_segment_id: Some("segment-0001".to_string()),
1775            new_segment_id: Some("segment-0002".to_string()),
1776            previous_prefix_hash: Some("prefix-before".to_string()),
1777            new_prefix_hash: Some("prefix-after".to_string()),
1778            previous_catalog_hash: Some("catalog-before".to_string()),
1779            new_catalog_hash: Some("catalog-after".to_string()),
1780        }));
1781
1782        let json = serde_json::to_string(&event)?;
1783        let restored: ThreadEvent = serde_json::from_str(&json)?;
1784
1785        assert_eq!(restored, event);
1786        Ok(())
1787    }
1788
1789    #[test]
1790    fn compact_boundary_deserializes_legacy_payload_without_segment_metadata() -> Result<(), Box<dyn Error>> {
1791        let payload = r#"{
1792            "type":"thread.compact_boundary",
1793            "thread_id":"thread-1",
1794            "trigger":"recovery",
1795            "mode":"provider",
1796            "original_message_count":12,
1797            "compacted_message_count":5
1798        }"#;
1799
1800        let restored: ThreadEvent = serde_json::from_str(payload)?;
1801        let ThreadEvent::ThreadCompactBoundary(event) = restored else {
1802            panic!("expected thread.compact_boundary event");
1803        };
1804
1805        assert_eq!(event.thread_id, "thread-1");
1806        assert_eq!(event.history_artifact_path, None);
1807        assert_eq!(event.previous_segment_id, None);
1808        assert_eq!(event.new_segment_id, None);
1809        assert_eq!(event.previous_prefix_hash, None);
1810        assert_eq!(event.new_prefix_hash, None);
1811        assert_eq!(event.previous_catalog_hash, None);
1812        assert_eq!(event.new_catalog_hash, None);
1813        Ok(())
1814    }
1815}