Skip to main content

vtcode_llm/open_responses/
events.rs

1//! Semantic streaming events for Open Responses.
2//!
3//! Streaming is modeled as a series of semantic events, not raw text deltas.
4//! Events describe meaningful transitions like state changes or content deltas.
5
6use serde::{Deserialize, Serialize};
7use std::sync::{Arc, Mutex};
8
9use super::{ContentPart, OutputItem, Response, ResponseId};
10
11/// Semantic streaming events per the Open Responses specification.
12///
13/// These events describe meaningful transitions during response generation,
14/// enabling predictable, provider-agnostic streaming clients.
15#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
16#[serde(tag = "type")]
17pub enum ResponseStreamEvent {
18    // ============================================================
19    // Response lifecycle events
20    // ============================================================
21    /// Initial response creation event.
22    #[serde(rename = "response.created")]
23    ResponseCreated {
24        /// The response object with initial state.
25        response: Response,
26    },
27
28    /// Response has started processing.
29    #[serde(rename = "response.in_progress")]
30    ResponseInProgress {
31        /// The response object.
32        response: Response,
33    },
34
35    /// Response completed successfully.
36    #[serde(rename = "response.completed")]
37    ResponseCompleted {
38        /// The final response object.
39        response: Response,
40    },
41
42    /// Response failed with an error.
43    #[serde(rename = "response.failed")]
44    ResponseFailed {
45        /// The response object with error details.
46        response: Response,
47    },
48
49    /// Response is incomplete (e.g., token limit reached).
50    #[serde(rename = "response.incomplete")]
51    ResponseIncomplete {
52        /// The response object with incomplete details.
53        response: Response,
54    },
55
56    // ============================================================
57    // Output item events
58    // ============================================================
59    /// New output item added to the response.
60    #[serde(rename = "response.output_item.added")]
61    OutputItemAdded {
62        /// ID of the containing response.
63        response_id: ResponseId,
64        /// Index of the item in the output array.
65        output_index: usize,
66        /// The output item being added.
67        item: OutputItem,
68    },
69
70    /// Output item is complete.
71    #[serde(rename = "response.output_item.done")]
72    OutputItemDone {
73        /// ID of the containing response.
74        response_id: ResponseId,
75        /// Index of the item in the output array.
76        output_index: usize,
77        /// The completed output item.
78        item: OutputItem,
79    },
80
81    // ============================================================
82    // Content part events
83    // ============================================================
84    /// New content part added to an output item.
85    #[serde(rename = "response.content_part.added")]
86    ContentPartAdded {
87        /// ID of the containing response.
88        response_id: ResponseId,
89        /// ID of the containing output item.
90        item_id: String,
91        /// Index of the item in the output array.
92        output_index: usize,
93        /// Index of the content part within the item.
94        content_index: usize,
95        /// The content part being added.
96        part: ContentPart,
97    },
98
99    /// Content part is complete.
100    #[serde(rename = "response.content_part.done")]
101    ContentPartDone {
102        /// ID of the containing response.
103        response_id: ResponseId,
104        /// ID of the containing output item.
105        item_id: String,
106        /// Index of the item in the output array.
107        output_index: usize,
108        /// Index of the content part within the item.
109        content_index: usize,
110        /// The completed content part.
111        part: ContentPart,
112    },
113
114    // ============================================================
115    // Text streaming events
116    // ============================================================
117    /// Text content delta for incremental streaming.
118    #[serde(rename = "response.output_text.delta")]
119    OutputTextDelta {
120        /// ID of the containing response.
121        response_id: ResponseId,
122        /// ID of the containing output item.
123        item_id: String,
124        /// Index of the item in the output array.
125        output_index: usize,
126        /// Index of the content part within the item.
127        content_index: usize,
128        /// The text delta to append.
129        delta: String,
130    },
131
132    /// Text content is complete.
133    #[serde(rename = "response.output_text.done")]
134    OutputTextDone {
135        /// ID of the containing response.
136        response_id: ResponseId,
137        /// ID of the containing output item.
138        item_id: String,
139        /// Index of the item in the output array.
140        output_index: usize,
141        /// Index of the content part within the item.
142        content_index: usize,
143        /// The complete text content.
144        text: String,
145    },
146
147    // ============================================================
148    // Function call streaming events
149    // ============================================================
150    /// Function call arguments delta.
151    #[serde(rename = "response.function_call_arguments.delta")]
152    FunctionCallArgumentsDelta {
153        /// ID of the containing response.
154        response_id: ResponseId,
155        /// ID of the function call item.
156        item_id: String,
157        /// Index of the item in the output array.
158        output_index: usize,
159        /// The arguments delta to append.
160        delta: String,
161    },
162
163    /// Function call arguments are complete.
164    #[serde(rename = "response.function_call_arguments.done")]
165    FunctionCallArgumentsDone {
166        /// ID of the containing response.
167        response_id: ResponseId,
168        /// ID of the function call item.
169        item_id: String,
170        /// Index of the item in the output array.
171        output_index: usize,
172        /// The complete arguments JSON string.
173        arguments: String,
174    },
175
176    // ============================================================
177    // Reasoning events
178    // ============================================================
179    /// Reasoning content delta.
180    #[serde(rename = "response.reasoning.delta")]
181    ReasoningDelta {
182        /// ID of the containing response.
183        response_id: ResponseId,
184        /// ID of the reasoning item.
185        item_id: String,
186        /// Index of the item in the output array.
187        output_index: usize,
188        /// The reasoning delta to append.
189        delta: String,
190    },
191
192    /// Reasoning content is complete.
193    #[serde(rename = "response.reasoning.done")]
194    ReasoningDone {
195        /// ID of the containing response.
196        response_id: ResponseId,
197        /// ID of the reasoning item.
198        item_id: String,
199        /// Index of the item in the output array.
200        output_index: usize,
201        /// The reasoning item with complete content.
202        item: OutputItem,
203    },
204
205    // ============================================================
206    // Extension events
207    // ============================================================
208    /// Custom/extension streaming event.
209    ///
210    /// Custom event types must be prefixed with the implementor slug
211    /// (e.g., `vtcode.trace_event`).
212    #[serde(rename = "response.custom_event")]
213    CustomEvent {
214        /// ID of the containing response.
215        response_id: ResponseId,
216        /// Custom event type (must be prefixed, e.g., `vtcode.telemetry`).
217        event_type: String,
218        /// Sequence number for ordering.
219        sequence_number: u64,
220        /// Custom event data.
221        data: serde_json::Value,
222    },
223    /// Catch-all for unknown streaming event types added by the Open Responses spec.
224    #[serde(other)]
225    Unknown,
226}
227
228impl ResponseStreamEvent {
229    /// Returns the response ID associated with this event.
230    #[must_use]
231    pub fn response_id(&self) -> Option<&ResponseId> {
232        match self {
233            Self::ResponseCreated { response, .. }
234            | Self::ResponseInProgress { response, .. }
235            | Self::ResponseCompleted { response, .. }
236            | Self::ResponseFailed { response, .. }
237            | Self::ResponseIncomplete { response, .. } => Some(&response.id),
238
239            Self::OutputItemAdded { response_id, .. }
240            | Self::OutputItemDone { response_id, .. }
241            | Self::ContentPartAdded { response_id, .. }
242            | Self::ContentPartDone { response_id, .. }
243            | Self::OutputTextDelta { response_id, .. }
244            | Self::OutputTextDone { response_id, .. }
245            | Self::FunctionCallArgumentsDelta { response_id, .. }
246            | Self::FunctionCallArgumentsDone { response_id, .. }
247            | Self::ReasoningDelta { response_id, .. }
248            | Self::ReasoningDone { response_id, .. }
249            | Self::CustomEvent { response_id, .. } => Some(response_id),
250            Self::Unknown => None,
251        }
252    }
253
254    /// Returns the event type name.
255    pub fn event_type(&self) -> &'static str {
256        match self {
257            Self::ResponseCreated { .. } => "response.created",
258            Self::ResponseInProgress { .. } => "response.in_progress",
259            Self::ResponseCompleted { .. } => "response.completed",
260            Self::ResponseFailed { .. } => "response.failed",
261            Self::ResponseIncomplete { .. } => "response.incomplete",
262            Self::OutputItemAdded { .. } => "response.output_item.added",
263            Self::OutputItemDone { .. } => "response.output_item.done",
264            Self::ContentPartAdded { .. } => "response.content_part.added",
265            Self::ContentPartDone { .. } => "response.content_part.done",
266            Self::OutputTextDelta { .. } => "response.output_text.delta",
267            Self::OutputTextDone { .. } => "response.output_text.done",
268            Self::FunctionCallArgumentsDelta { .. } => "response.function_call_arguments.delta",
269            Self::FunctionCallArgumentsDone { .. } => "response.function_call_arguments.done",
270            Self::ReasoningDelta { .. } => "response.reasoning.delta",
271            Self::ReasoningDone { .. } => "response.reasoning.done",
272            Self::CustomEvent { .. } => "response.custom_event",
273            Self::Unknown => "unknown",
274        }
275    }
276
277    /// Returns true if this is a response lifecycle event.
278    pub fn is_response_event(&self) -> bool {
279        matches!(
280            self,
281            Self::ResponseCreated { .. }
282                | Self::ResponseInProgress { .. }
283                | Self::ResponseCompleted { .. }
284                | Self::ResponseFailed { .. }
285                | Self::ResponseIncomplete { .. }
286        )
287    }
288
289    /// Returns true if this is a terminal event.
290    fn is_terminal(&self) -> bool {
291        matches!(self, Self::ResponseCompleted { .. } | Self::ResponseFailed { .. } | Self::ResponseIncomplete { .. })
292    }
293
294    /// Returns true if this is an unknown/unsupported event type.
295    pub fn is_unknown(&self) -> bool {
296        matches!(self, Self::Unknown)
297    }
298}
299
300/// Callback type for streaming events.
301#[expect(
302    dead_code,
303    reason = "Intentional compatibility, platform, test, or API-shape suppression."
304)]
305pub type StreamEventCallback = Arc<Mutex<Box<dyn FnMut(&ResponseStreamEvent) + Send>>>;
306
307/// Trait for emitting Open Responses streaming events.
308pub trait StreamEventEmitter: Send {
309    /// Emit a streaming event.
310    fn emit(&mut self, event: ResponseStreamEvent);
311
312    /// Emit a response created event.
313    fn response_created(&mut self, response: Response) {
314        self.emit(ResponseStreamEvent::ResponseCreated { response });
315    }
316
317    /// Emit a response in progress event.
318    fn response_in_progress(&mut self, response: Response) {
319        self.emit(ResponseStreamEvent::ResponseInProgress { response });
320    }
321
322    /// Emit a response completed event.
323    fn response_completed(&mut self, response: Response) {
324        self.emit(ResponseStreamEvent::ResponseCompleted { response });
325    }
326
327    /// Emit a response failed event.
328    fn response_failed(&mut self, response: Response) {
329        self.emit(ResponseStreamEvent::ResponseFailed { response });
330    }
331
332    /// Emit an output item added event.
333    fn output_item_added(&mut self, response_id: &ResponseId, output_index: usize, item: OutputItem) {
334        self.emit(ResponseStreamEvent::OutputItemAdded {
335            response_id: response_id.clone(),
336            output_index,
337            item,
338        });
339    }
340
341    /// Emit an output item done event.
342    fn output_item_done(&mut self, response_id: &ResponseId, output_index: usize, item: OutputItem) {
343        self.emit(ResponseStreamEvent::OutputItemDone {
344            response_id: response_id.clone(),
345            output_index,
346            item,
347        });
348    }
349
350    /// Emit a text delta event.
351    fn output_text_delta(
352        &mut self,
353        response_id: &ResponseId,
354        item_id: &str,
355        output_index: usize,
356        content_index: usize,
357        delta: &str,
358    ) {
359        self.emit(ResponseStreamEvent::OutputTextDelta {
360            response_id: response_id.clone(),
361            item_id: item_id.to_string(),
362            output_index,
363            content_index,
364            delta: delta.to_string(),
365        });
366    }
367
368    /// Emit a reasoning delta event.
369    fn reasoning_delta(&mut self, response_id: &ResponseId, item_id: &str, output_index: usize, delta: &str) {
370        self.emit(ResponseStreamEvent::ReasoningDelta {
371            response_id: response_id.clone(),
372            item_id: item_id.to_string(),
373            output_index,
374            delta: delta.to_string(),
375        });
376    }
377}
378
379/// Vector-based event emitter for collecting events.
380#[derive(Debug, Default)]
381pub struct VecStreamEmitter {
382    events: Vec<ResponseStreamEvent>,
383}
384
385impl VecStreamEmitter {
386    /// Creates a new vector emitter.
387    pub fn new() -> Self {
388        Self::default()
389    }
390
391    /// Returns the collected events.
392    fn events(&self) -> &[ResponseStreamEvent] {
393        &self.events
394    }
395
396    /// Consumes and returns the collected events.
397    pub fn into_events(self) -> Vec<ResponseStreamEvent> {
398        self.events
399    }
400}
401
402impl StreamEventEmitter for VecStreamEmitter {
403    fn emit(&mut self, event: ResponseStreamEvent) {
404        self.events.push(event);
405    }
406}
407
408/// Wrapper for streaming events with sequence number for ordering.
409/// Used when serializing events for SSE transport.
410#[derive(Debug, Clone, Serialize)]
411pub struct SequencedEvent<'a> {
412    /// Monotonically increasing sequence number within the stream.
413    sequence_number: u64,
414    /// The underlying event.
415    #[serde(flatten)]
416    event: &'a ResponseStreamEvent,
417}
418
419impl<'a> SequencedEvent<'a> {
420    /// Creates a new sequenced event.
421    pub fn new(sequence_number: u64, event: &'a ResponseStreamEvent) -> Self {
422        Self { sequence_number, event }
423    }
424}
425
426#[cfg(test)]
427mod tests {
428    use super::*;
429
430    #[test]
431    fn test_event_type() {
432        let response = Response::new("resp_1", "gpt-5");
433        let event = ResponseStreamEvent::ResponseCreated { response };
434        assert_eq!(event.event_type(), "response.created");
435    }
436
437    #[test]
438    fn test_terminal_events() {
439        let response = Response::new("resp_1", "gpt-5");
440        let created = ResponseStreamEvent::ResponseCreated { response: response.clone() };
441        let completed = ResponseStreamEvent::ResponseCompleted { response };
442        assert!(!created.is_terminal());
443        assert!(completed.is_terminal());
444    }
445
446    #[test]
447    fn test_vec_emitter() {
448        let mut emitter = VecStreamEmitter::new();
449        let response = Response::new("resp_1", "gpt-5");
450        emitter.response_created(response);
451        assert_eq!(emitter.events().len(), 1);
452    }
453}