1use serde_json::Value;
2
3use super::types::{EventFrame, EventPayload, SSEEventType, SSEItemType, WireEvent};
4use crate::utils::common::{deserialize_from_str_opt, deserialize_from_value_opt};
5
6#[must_use]
12pub fn normalize_sse_line(line: &str) -> Option<EventFrame> {
13 let data_str = line.strip_prefix("data: ")?;
14 if data_str == "[DONE]" {
15 return None;
16 }
17
18 let json: Value = deserialize_from_str_opt(data_str)?;
19 let event_type = json
20 .get("type")
21 .and_then(Value::as_str)
22 .map_or(SSEEventType::Other, SSEEventType::from);
23
24 let payload = extract_payload(event_type, &json);
25 let wire: WireEvent = deserialize_from_value_opt(json)?;
26
27 Some(EventFrame {
28 event_type,
29 payload,
30 wire,
31 })
32}
33
34fn extract_payload(event_type: SSEEventType, json: &Value) -> EventPayload {
36 match event_type {
37 SSEEventType::ResponseCreated
38 | SSEEventType::ResponseInProgress
39 | SSEEventType::ResponseCompleted
40 | SSEEventType::ResponseFailed
41 | SSEEventType::ResponseIncomplete => extract_response_payload(json),
42
43 SSEEventType::OutputItemAdded => extract_output_item_added(json),
44 SSEEventType::OutputItemDone => extract_output_item_done(json),
45
46 SSEEventType::OutputTextDelta => extract_text_delta(json),
47 SSEEventType::OutputTextDone => extract_text_done(json),
48
49 SSEEventType::FunctionCallArgumentsDelta => extract_fn_call_args_delta(json),
50 SSEEventType::FunctionCallArgumentsDone => extract_fn_call_args_done(json),
51 SSEEventType::CustomToolCallInputDelta => extract_custom_tool_call_input_delta(json),
52 SSEEventType::CustomToolCallInputDone => extract_custom_tool_call_input_done(json),
53
54 SSEEventType::ReasoningTextDelta | SSEEventType::ReasoningSummaryTextDelta => extract_reasoning_delta(json),
55 SSEEventType::ReasoningTextDone | SSEEventType::ReasoningSummaryTextDone => extract_reasoning_done(json),
56
57 SSEEventType::ContentPartAdded
58 | SSEEventType::ContentPartDone
59 | SSEEventType::ReasoningPartAdded
60 | SSEEventType::ReasoningPartDone
61 | SSEEventType::FileSearchCallSearching
62 | SSEEventType::FileSearchCallCompleted
63 | SSEEventType::WebSearchCallInProgress
64 | SSEEventType::WebSearchCallSearching
65 | SSEEventType::WebSearchCallCompleted
66 | SSEEventType::McpCallInProgress
67 | SSEEventType::McpCallArgumentsDelta
68 | SSEEventType::McpCallArgumentsDone
69 | SSEEventType::McpCallCompleted
70 | SSEEventType::McpCallFailed
71 | SSEEventType::McpListToolsInProgress
72 | SSEEventType::McpListToolsCompleted
73 | SSEEventType::McpListToolsFailed
74 | SSEEventType::Other => EventPayload::Raw(json.clone()),
75 }
76}
77
78fn json_str(json: &Value, key: &str) -> String {
79 json[key].as_str().unwrap_or_default().to_string()
80}
81
82fn json_str_opt(json: &Value, key: &str) -> Option<String> {
83 json[key].as_str().map(ToString::to_string)
84}
85
86fn json_u32(json: &Value, key: &str) -> u32 {
87 u32::try_from(json[key].as_u64().unwrap_or(0)).unwrap_or(u32::MAX)
88}
89
90fn extract_response_payload(json: &Value) -> EventPayload {
91 let response = &json["response"];
92 EventPayload::Response {
93 id: json_str(response, "id"),
94 status: json_str(response, "status"),
95 usage: response
96 .get("usage")
97 .filter(|v| !v.is_null())
98 .and_then(|v| deserialize_from_value_opt(v.clone())),
99 }
100}
101
102fn extract_output_item_added(json: &Value) -> EventPayload {
103 let item = &json["item"];
104 EventPayload::OutputItemAdded {
105 item_id: json_str(item, "id"),
106 item_type: SSEItemType::from(json_str(item, "type")),
107 output_index: json_u32(json, "output_index"),
108 name: json_str_opt(item, "name"),
109 namespace: json_str_opt(item, "namespace"),
110 call_id: json_str_opt(item, "call_id"),
111 }
112}
113
114fn extract_output_item_done(json: &Value) -> EventPayload {
115 let item = &json["item"];
116 EventPayload::OutputItemDone {
117 item_id: json_str(item, "id"),
118 item_type: SSEItemType::from(json_str(item, "type")),
119 output_index: json_u32(json, "output_index"),
120 item: item.clone(),
121 }
122}
123
124fn extract_text_delta(json: &Value) -> EventPayload {
125 EventPayload::TextDelta {
126 delta: json_str(json, "delta"),
127 item_id: json_str(json, "item_id"),
128 output_index: json_u32(json, "output_index"),
129 content_index: json_u32(json, "content_index"),
130 }
131}
132
133fn extract_text_done(json: &Value) -> EventPayload {
134 EventPayload::TextDone {
135 text: json_str(json, "text"),
136 item_id: json_str(json, "item_id"),
137 output_index: json_u32(json, "output_index"),
138 }
139}
140
141fn extract_fn_call_args_delta(json: &Value) -> EventPayload {
142 EventPayload::FunctionCallArgsDelta {
143 delta: json_str(json, "delta"),
144 call_id: json_str_opt(json, "call_id"),
145 item_id: json_str(json, "item_id"),
146 output_index: json_u32(json, "output_index"),
147 }
148}
149
150fn extract_fn_call_args_done(json: &Value) -> EventPayload {
151 EventPayload::FunctionCallArgsDone {
152 arguments: json_str(json, "arguments"),
153 call_id: json_str_opt(json, "call_id"),
154 item_id: json_str(json, "item_id"),
155 name: json_str(json, "name"),
156 output_index: json_u32(json, "output_index"),
157 }
158}
159
160fn extract_custom_tool_call_input_delta(json: &Value) -> EventPayload {
161 EventPayload::CustomToolCallInputDelta {
162 delta: json_str(json, "delta"),
163 item_id: json_str(json, "item_id"),
164 output_index: json_u32(json, "output_index"),
165 }
166}
167
168fn extract_custom_tool_call_input_done(json: &Value) -> EventPayload {
169 EventPayload::CustomToolCallInputDone {
170 input: json_str(json, "input"),
171 item_id: json_str(json, "item_id"),
172 output_index: json_u32(json, "output_index"),
173 }
174}
175
176fn extract_reasoning_delta(json: &Value) -> EventPayload {
177 EventPayload::ReasoningDelta {
178 delta: json_str(json, "delta"),
179 item_id: json_str(json, "item_id"),
180 }
181}
182
183fn extract_reasoning_done(json: &Value) -> EventPayload {
184 EventPayload::ReasoningDone {
185 text: json_str(json, "text"),
186 item_id: json_str(json, "item_id"),
187 }
188}