Skip to main content

agentic_core/events/
normalize.rs

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/// Normalize a raw SSE data line into a typed [`EventFrame`].
7///
8/// Expects input in the form `data: {...}` (the `data: ` prefix is required).
9/// Returns `None` for non-data lines, empty lines, and the `data: [DONE]`
10/// sentinel.
11#[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
34/// Extract a typed payload from the JSON body based on the classified event type.
35fn 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}