1use crate::event::{EventKind, OutputStream};
19
20#[derive(Debug, Clone, Copy, Default)]
22pub struct NormalizeOptions {
23 pub include_thinking: bool,
26}
27
28#[derive(Debug, PartialEq)]
30pub enum Normalized {
31 Event(EventKind),
33 Dropped,
36 Unknown,
38}
39
40pub fn normalize_frame(frame: &serde_json::Value, opts: NormalizeOptions) -> Normalized {
42 let ty = match frame.get("type").and_then(|v| v.as_str()) {
43 Some(t) => t,
44 None => return Normalized::Unknown,
45 };
46
47 match ty {
48 "AssistantTextDelta" => Normalized::Event(EventKind::Output {
50 stream: OutputStream::Stdout,
51 chunk: str_field(frame, "text"),
52 }),
53 "AssistantThinkingDelta" => {
55 if opts.include_thinking {
56 Normalized::Event(EventKind::Output {
57 stream: OutputStream::Thinking,
58 chunk: str_field(frame, "text"),
59 })
60 } else {
61 Normalized::Dropped
62 }
63 }
64 "ToolCallStart" => Normalized::Event(EventKind::ToolStarted {
67 id: str_field(frame, "tool_use_id"),
68 name: str_field(frame, "name"),
69 input: serde_json::Value::Null,
70 }),
71 "ToolCallEnd" => Normalized::Event(EventKind::ToolFinished {
76 id: str_field(frame, "tool_use_id"),
77 name: String::new(),
78 ok: !frame
79 .get("is_error")
80 .and_then(|v| v.as_bool())
81 .unwrap_or(false),
82 }),
83 "RunEnd" => Normalized::Event(EventKind::AgentFinished {
85 outcome: stop_label(frame.get("stopped_for")),
86 cost_usd: 0.0,
88 turns: frame
89 .get("turn_count")
90 .and_then(|v| v.as_u64())
91 .unwrap_or(0) as u32,
92 }),
93 "TurnStart" | "ToolCallInputDelta" | "TurnEnd" => Normalized::Dropped,
97 _ => Normalized::Unknown,
98 }
99}
100
101fn str_field(frame: &serde_json::Value, key: &str) -> String {
102 frame
103 .get(key)
104 .and_then(|v| v.as_str())
105 .unwrap_or_default()
106 .to_string()
107}
108
109fn stop_label(v: Option<&serde_json::Value>) -> String {
114 match v {
115 Some(serde_json::Value::String(s)) => s.clone(),
116 Some(serde_json::Value::Object(m)) => m.keys().next().cloned().unwrap_or_default(),
117 _ => String::new(),
118 }
119}
120
121#[cfg(test)]
122mod tests {
123 use super::*;
124 use serde_json::json;
125
126 fn norm(v: serde_json::Value) -> Normalized {
127 normalize_frame(&v, NormalizeOptions::default())
128 }
129
130 #[test]
131 fn assistant_text_delta_maps_to_stdout_output() {
132 let f = json!({
134 "type": "AssistantTextDelta",
135 "turn_index": 0, "content_block_index": 0, "text": "hello "
136 });
137 assert_eq!(
138 norm(f),
139 Normalized::Event(EventKind::Output {
140 stream: OutputStream::Stdout,
141 chunk: "hello ".into()
142 })
143 );
144 }
145
146 #[test]
147 fn thinking_delta_dropped_by_default_but_included_on_request() {
148 let f = json!({
149 "type": "AssistantThinkingDelta",
150 "turn_index": 0, "content_block_index": 0, "text": "hmm"
151 });
152 assert_eq!(norm(f.clone()), Normalized::Dropped);
153 let opts = NormalizeOptions {
154 include_thinking: true,
155 };
156 assert_eq!(
157 normalize_frame(&f, opts),
158 Normalized::Event(EventKind::Output {
159 stream: OutputStream::Thinking,
160 chunk: "hmm".into()
161 })
162 );
163 }
164
165 #[test]
166 fn tool_call_start_maps_to_tool_started() {
167 assert_eq!(
168 norm(json!({
169 "type": "ToolCallStart",
170 "turn_index": 0, "tool_use_id": "tu_1", "name": "Read"
171 })),
172 Normalized::Event(EventKind::ToolStarted {
173 id: "tu_1".into(),
174 name: "Read".into(),
175 input: serde_json::Value::Null,
176 })
177 );
178 }
179
180 #[test]
181 fn tool_call_end_ok_is_inverse_of_is_error() {
182 assert_eq!(
183 norm(json!({
184 "type": "ToolCallEnd",
185 "turn_index": 0, "tool_use_id": "tu_1", "is_error": true, "content": []
186 })),
187 Normalized::Event(EventKind::ToolFinished {
188 id: "tu_1".into(),
189 name: String::new(),
190 ok: false
191 })
192 );
193 assert_eq!(
194 norm(json!({
195 "type": "ToolCallEnd",
196 "turn_index": 0, "tool_use_id": "tu_1", "is_error": false, "content": []
197 })),
198 Normalized::Event(EventKind::ToolFinished {
199 id: "tu_1".into(),
200 name: String::new(),
201 ok: true
202 })
203 );
204 }
205
206 #[test]
207 fn run_end_is_terminal_and_carries_turns_and_outcome() {
208 let f = json!({
211 "type": "RunEnd",
212 "final_messages": [],
213 "total_usage": {"input_tokens": 0, "output_tokens": 0},
214 "turn_count": 3,
215 "stopped_for": "EndOfTurn"
216 });
217 assert_eq!(
218 norm(f),
219 Normalized::Event(EventKind::AgentFinished {
220 outcome: "EndOfTurn".into(),
221 cost_usd: 0.0,
222 turns: 3
223 })
224 );
225 }
226
227 #[test]
228 fn run_end_outcome_from_data_carrying_stop_condition() {
229 let f = json!({
230 "type": "RunEnd",
231 "final_messages": [], "total_usage": {},
232 "turn_count": 10,
233 "stopped_for": {"MaxTurnsReached": 10}
234 });
235 assert_eq!(
236 norm(f),
237 Normalized::Event(EventKind::AgentFinished {
238 outcome: "MaxTurnsReached".into(),
239 cost_usd: 0.0,
240 turns: 10
241 })
242 );
243 }
244
245 #[test]
246 fn turn_boundary_and_tool_input_frames_are_dropped_not_unknown() {
247 for f in [
248 json!({"type": "TurnStart", "turn_index": 0, "message_id": "m1", "model": "x"}),
249 json!({"type": "ToolCallInputDelta", "turn_index": 0, "tool_use_id": "tu_1", "partial_json": "{"}),
250 json!({"type": "TurnEnd", "turn_index": 0}),
251 ] {
252 assert_eq!(norm(f), Normalized::Dropped);
253 }
254 }
255
256 #[test]
257 fn unknown_or_typeless_frames_are_unknown() {
258 assert_eq!(norm(json!({"type": "FutureThing"})), Normalized::Unknown);
259 assert_eq!(norm(json!({"no_type": true})), Normalized::Unknown);
260 }
261
262 #[test]
268 fn recognizes_every_known_turnevent_type() {
269 let opts = NormalizeOptions {
270 include_thinking: true,
271 };
272 for ty in [
273 "TurnStart",
274 "AssistantTextDelta",
275 "AssistantThinkingDelta",
276 "ToolCallStart",
277 "ToolCallInputDelta",
278 "ToolCallEnd",
279 "TurnEnd",
280 "RunEnd",
281 ] {
282 let got = normalize_frame(&json!({ "type": ty }), opts);
283 assert_ne!(
284 got,
285 Normalized::Unknown,
286 "caliban TurnEvent::{ty} is not recognized by normalize_frame"
287 );
288 }
289 }
290}