1use serde::{Deserialize, Serialize};
7use std::sync::{Arc, Mutex};
8
9use super::{ContentPart, OutputItem, Response, ResponseId};
10
11#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
16#[serde(tag = "type")]
17pub enum ResponseStreamEvent {
18 #[serde(rename = "response.created")]
23 ResponseCreated {
24 response: Response,
26 },
27
28 #[serde(rename = "response.in_progress")]
30 ResponseInProgress {
31 response: Response,
33 },
34
35 #[serde(rename = "response.completed")]
37 ResponseCompleted {
38 response: Response,
40 },
41
42 #[serde(rename = "response.failed")]
44 ResponseFailed {
45 response: Response,
47 },
48
49 #[serde(rename = "response.incomplete")]
51 ResponseIncomplete {
52 response: Response,
54 },
55
56 #[serde(rename = "response.output_item.added")]
61 OutputItemAdded {
62 response_id: ResponseId,
64 output_index: usize,
66 item: OutputItem,
68 },
69
70 #[serde(rename = "response.output_item.done")]
72 OutputItemDone {
73 response_id: ResponseId,
75 output_index: usize,
77 item: OutputItem,
79 },
80
81 #[serde(rename = "response.content_part.added")]
86 ContentPartAdded {
87 response_id: ResponseId,
89 item_id: String,
91 output_index: usize,
93 content_index: usize,
95 part: ContentPart,
97 },
98
99 #[serde(rename = "response.content_part.done")]
101 ContentPartDone {
102 response_id: ResponseId,
104 item_id: String,
106 output_index: usize,
108 content_index: usize,
110 part: ContentPart,
112 },
113
114 #[serde(rename = "response.output_text.delta")]
119 OutputTextDelta {
120 response_id: ResponseId,
122 item_id: String,
124 output_index: usize,
126 content_index: usize,
128 delta: String,
130 },
131
132 #[serde(rename = "response.output_text.done")]
134 OutputTextDone {
135 response_id: ResponseId,
137 item_id: String,
139 output_index: usize,
141 content_index: usize,
143 text: String,
145 },
146
147 #[serde(rename = "response.function_call_arguments.delta")]
152 FunctionCallArgumentsDelta {
153 response_id: ResponseId,
155 item_id: String,
157 output_index: usize,
159 delta: String,
161 },
162
163 #[serde(rename = "response.function_call_arguments.done")]
165 FunctionCallArgumentsDone {
166 response_id: ResponseId,
168 item_id: String,
170 output_index: usize,
172 arguments: String,
174 },
175
176 #[serde(rename = "response.reasoning.delta")]
181 ReasoningDelta {
182 response_id: ResponseId,
184 item_id: String,
186 output_index: usize,
188 delta: String,
190 },
191
192 #[serde(rename = "response.reasoning.done")]
194 ReasoningDone {
195 response_id: ResponseId,
197 item_id: String,
199 output_index: usize,
201 item: OutputItem,
203 },
204
205 #[serde(rename = "response.custom_event")]
213 CustomEvent {
214 response_id: ResponseId,
216 event_type: String,
218 sequence_number: u64,
220 data: serde_json::Value,
222 },
223 #[serde(other)]
225 Unknown,
226}
227
228impl ResponseStreamEvent {
229 #[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 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 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 fn is_terminal(&self) -> bool {
291 matches!(self, Self::ResponseCompleted { .. } | Self::ResponseFailed { .. } | Self::ResponseIncomplete { .. })
292 }
293
294 pub fn is_unknown(&self) -> bool {
296 matches!(self, Self::Unknown)
297 }
298}
299
300#[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
307pub trait StreamEventEmitter: Send {
309 fn emit(&mut self, event: ResponseStreamEvent);
311
312 fn response_created(&mut self, response: Response) {
314 self.emit(ResponseStreamEvent::ResponseCreated { response });
315 }
316
317 fn response_in_progress(&mut self, response: Response) {
319 self.emit(ResponseStreamEvent::ResponseInProgress { response });
320 }
321
322 fn response_completed(&mut self, response: Response) {
324 self.emit(ResponseStreamEvent::ResponseCompleted { response });
325 }
326
327 fn response_failed(&mut self, response: Response) {
329 self.emit(ResponseStreamEvent::ResponseFailed { response });
330 }
331
332 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 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 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 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#[derive(Debug, Default)]
381pub struct VecStreamEmitter {
382 events: Vec<ResponseStreamEvent>,
383}
384
385impl VecStreamEmitter {
386 pub fn new() -> Self {
388 Self::default()
389 }
390
391 fn events(&self) -> &[ResponseStreamEvent] {
393 &self.events
394 }
395
396 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#[derive(Debug, Clone, Serialize)]
411pub struct SequencedEvent<'a> {
412 sequence_number: u64,
414 #[serde(flatten)]
416 event: &'a ResponseStreamEvent,
417}
418
419impl<'a> SequencedEvent<'a> {
420 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}