1use crate::{
51 Event, MessageId, ReasoningMessageChunkEvent, SubagentRunId, TextMessageChunkEvent,
52 TextMessageStartEvent, ToolCallChunkEvent, ToolCallId, ToolCallStartEvent,
53};
54
55use crate::client::error::{Error, Result};
56
57#[derive(Clone, Copy, Debug, PartialEq, Eq)]
59enum Kind {
60 Text,
61 Tool,
62 Reasoning,
63}
64
65impl Kind {
66 const fn chunk_name(self) -> &'static str {
67 match self {
68 Self::Text => "TEXT_MESSAGE_CHUNK",
69 Self::Tool => "TOOL_CALL_CHUNK",
70 Self::Reasoning => "REASONING_MESSAGE_CHUNK",
71 }
72 }
73
74 const fn id_name(self) -> &'static str {
75 match self {
76 Self::Text | Self::Reasoning => "messageId",
77 Self::Tool => "toolCallId",
78 }
79 }
80
81 const fn noun(self) -> &'static str {
82 match self {
83 Self::Text => "message",
84 Self::Tool => "call",
85 Self::Reasoning => "reasoning message",
86 }
87 }
88}
89
90#[derive(Clone, Debug)]
93struct Open {
94 kind: Kind,
95 id: String,
97 owed: bool,
100 owner: Option<SubagentRunId>,
102}
103
104#[derive(Clone, Debug, Default)]
112pub struct ChunkNormalizer {
113 open: Vec<Open>,
114}
115
116impl ChunkNormalizer {
117 pub fn new() -> Self {
119 Self::default()
120 }
121
122 pub fn is_open(&self) -> bool {
124 !self.open.is_empty()
125 }
126
127 pub fn normalize(&mut self, event: Event, out: &mut Vec<Event>) -> Result<()> {
140 match event {
141 Event::TextMessageChunk(chunk) => self.text_chunk(chunk, out),
142 Event::ToolCallChunk(chunk) => self.tool_chunk(chunk, out),
143 Event::ReasoningMessageChunk(chunk) => self.reasoning_chunk(chunk, out),
144 other => {
145 self.observe(&other, out);
146 out.push(other);
147 Ok(())
148 }
149 }
150 }
151
152 pub fn finish(&mut self, out: &mut Vec<Event>) {
158 self.close_all(out);
159 }
160
161 fn stream_id(
166 &self,
167 kind: Kind,
168 named: Option<&str>,
169 tag: &Option<SubagentRunId>,
170 ) -> Result<String> {
171 if let Some(id) = named {
172 return Ok(id.to_owned());
173 }
174 let missing = |why: String| {
175 Error::protocol(format!(
176 "{} carries no {} and {why}",
177 kind.chunk_name(),
178 kind.id_name()
179 ))
180 };
181 if tag.is_some() {
182 return self.current(kind, tag).ok_or_else(|| {
183 missing(format!(
184 "subagent {:?} has no {} open",
185 tag.as_ref().map(|id| id.as_str()).unwrap_or_default(),
186 kind.noun()
187 ))
188 });
189 }
190 if let Some(id) = self.current(kind, &None) {
191 return Ok(id);
192 }
193 let mut of_kind = self.open.iter().filter(|open| open.kind == kind);
194 match (of_kind.next(), of_kind.next()) {
195 (Some(only), None) => Ok(only.id.clone()),
196 (None, _) => Err(missing(format!("no {} is open", kind.noun()))),
197 (Some(_), Some(_)) => Err(missing(format!(
198 "several subagents have a {} open; attribute the chunk",
199 kind.noun()
200 ))),
201 }
202 }
203
204 fn begin(
207 &mut self,
208 kind: Kind,
209 id: &str,
210 owner: Option<SubagentRunId>,
211 start: Event,
212 out: &mut Vec<Event>,
213 ) {
214 self.close_owner(&owner, out);
215 out.push(start);
216 self.open.push(Open {
217 kind,
218 id: id.to_owned(),
219 owed: true,
220 owner,
221 });
222 }
223
224 fn text_chunk(&mut self, chunk: TextMessageChunkEvent, out: &mut Vec<Event>) -> Result<()> {
225 let id = MessageId::new(self.stream_id(
226 Kind::Text,
227 chunk.message_id.as_ref().map(MessageId::as_str),
228 &chunk.subagent_run_id,
229 )?);
230
231 let owner = self.owner_of(Kind::Text, id.as_str(), &chunk.subagent_run_id);
232 if self.current(Kind::Text, &owner).as_deref() != Some(id.as_str()) {
233 let mut start = TextMessageStartEvent::new(id.clone(), chunk.role.unwrap_or_default());
234 start.name = chunk.name;
235 start.base = chunk.base.clone();
236 start.subagent_run_id = owner.clone();
237 self.begin(Kind::Text, id.as_str(), owner.clone(), start.into(), out);
238 }
239
240 if let Some(delta) = chunk.delta {
241 let mut content = crate::TextMessageContentEvent::new(id, delta);
242 content.base = chunk.base;
243 content.subagent_run_id = owner;
244 out.push(content.into());
245 }
246 Ok(())
247 }
248
249 fn tool_chunk(&mut self, chunk: ToolCallChunkEvent, out: &mut Vec<Event>) -> Result<()> {
250 let id = ToolCallId::new(self.stream_id(
251 Kind::Tool,
252 chunk.tool_call_id.as_ref().map(ToolCallId::as_str),
253 &chunk.subagent_run_id,
254 )?);
255
256 let owner = self.owner_of(Kind::Tool, id.as_str(), &chunk.subagent_run_id);
257 if self.current(Kind::Tool, &owner).as_deref() != Some(id.as_str()) {
258 let Some(name) = chunk.tool_call_name else {
262 return Err(Error::protocol(format!(
263 "TOOL_CALL_CHUNK opens tool call {id:?} without a toolCallName"
264 )));
265 };
266 let mut start = ToolCallStartEvent::new(id.clone(), name);
267 start.parent_message_id = chunk.parent_message_id;
268 start.base = chunk.base.clone();
269 start.subagent_run_id = owner.clone();
270 self.begin(Kind::Tool, id.as_str(), owner.clone(), start.into(), out);
271 }
272
273 if let Some(delta) = chunk.delta {
274 let mut args = crate::ToolCallArgsEvent::new(id, delta);
275 args.base = chunk.base;
276 args.subagent_run_id = owner;
277 out.push(args.into());
278 }
279 Ok(())
280 }
281
282 fn reasoning_chunk(
283 &mut self,
284 chunk: ReasoningMessageChunkEvent,
285 out: &mut Vec<Event>,
286 ) -> Result<()> {
287 let id = MessageId::new(self.stream_id(
288 Kind::Reasoning,
289 chunk.message_id.as_ref().map(MessageId::as_str),
290 &chunk.subagent_run_id,
291 )?);
292
293 let owner = self.owner_of(Kind::Reasoning, id.as_str(), &chunk.subagent_run_id);
294 if self.current(Kind::Reasoning, &owner).as_deref() != Some(id.as_str()) {
295 let mut start = crate::ReasoningMessageStartEvent::new(id.clone());
296 start.base = chunk.base.clone();
297 start.subagent_run_id = owner.clone();
298 self.begin(
299 Kind::Reasoning,
300 id.as_str(),
301 owner.clone(),
302 start.into(),
303 out,
304 );
305 }
306
307 if let Some(delta) = chunk.delta {
308 let mut content = crate::ReasoningMessageContentEvent::new(id, delta);
309 content.base = chunk.base;
310 content.subagent_run_id = owner;
311 out.push(content.into());
312 }
313 Ok(())
314 }
315
316 fn observe(&mut self, event: &Event, out: &mut Vec<Event>) {
324 match event {
325 Event::TextMessageStart(e) => {
326 self.open_explicit(Kind::Text, e.message_id.as_str(), &e.subagent_run_id, out);
327 }
328 Event::TextMessageContent(e) => {
329 self.open_explicit(Kind::Text, e.message_id.as_str(), &e.subagent_run_id, out);
330 }
331 Event::TextMessageEnd(e) => self.close_explicit(Kind::Text, e.message_id.as_str()),
332
333 Event::ToolCallStart(e) => {
334 self.open_explicit(Kind::Tool, e.tool_call_id.as_str(), &e.subagent_run_id, out);
335 }
336 Event::ToolCallArgs(e) => {
337 self.open_explicit(Kind::Tool, e.tool_call_id.as_str(), &e.subagent_run_id, out);
338 }
339 Event::ToolCallEnd(e) => self.close_explicit(Kind::Tool, e.tool_call_id.as_str()),
340 Event::ToolCallResult(e) => {
347 self.close_id(Kind::Tool, e.tool_call_id.as_str(), out);
348 self.close_owner(&e.subagent_run_id, out);
349 }
350
351 Event::ReasoningMessageStart(e) => {
352 self.open_explicit(
353 Kind::Reasoning,
354 e.message_id.as_str(),
355 &e.subagent_run_id,
356 out,
357 );
358 }
359 Event::ReasoningMessageContent(e) => {
360 self.open_explicit(
361 Kind::Reasoning,
362 e.message_id.as_str(),
363 &e.subagent_run_id,
364 out,
365 );
366 }
367 Event::ReasoningMessageEnd(e) => {
368 self.close_explicit(Kind::Reasoning, e.message_id.as_str());
369 }
370 Event::ReasoningEnd(e) => self.close_id(Kind::Reasoning, e.message_id.as_str(), out),
372 Event::RunFinished(_) | Event::RunError(_) => self.close_all(out),
373 _ => {}
374 }
375 }
376
377 fn open_explicit(
384 &mut self,
385 kind: Kind,
386 id: &str,
387 tag: &Option<SubagentRunId>,
388 out: &mut Vec<Event>,
389 ) {
390 let owner = self.owner_of(kind, id, tag);
391 if self.current(kind, &owner).as_deref() != Some(id) {
392 self.close_owner(&owner, out);
393 self.open.push(Open {
394 kind,
395 id: id.to_owned(),
396 owed: false,
397 owner,
398 });
399 }
400 }
401
402 fn close_explicit(&mut self, kind: Kind, id: &str) {
404 self.open
405 .retain(|open| !(open.kind == kind && open.id == id));
406 }
407
408 fn owner_of(&self, kind: Kind, id: &str, tag: &Option<SubagentRunId>) -> Option<SubagentRunId> {
415 if tag.is_some() {
416 return tag.clone();
417 }
418 self.open
419 .iter()
420 .find(|open| open.kind == kind && open.id == id)
421 .and_then(|open| open.owner.clone())
422 }
423
424 fn current(&self, kind: Kind, owner: &Option<SubagentRunId>) -> Option<String> {
426 self.open
427 .iter()
428 .find(|open| &open.owner == owner)
429 .filter(|open| open.kind == kind)
430 .map(|open| open.id.clone())
431 }
432
433 fn close_owner(&mut self, owner: &Option<SubagentRunId>, out: &mut Vec<Event>) {
436 if let Some(index) = self.open.iter().position(|open| &open.owner == owner) {
437 let open = self.open.remove(index);
438 Self::settle(open, out);
439 }
440 }
441
442 fn close_id(&mut self, kind: Kind, id: &str, out: &mut Vec<Event>) {
444 if let Some(index) = self
445 .open
446 .iter()
447 .position(|open| open.kind == kind && open.id == id)
448 {
449 let open = self.open.remove(index);
450 Self::settle(open, out);
451 }
452 }
453
454 fn close_all(&mut self, out: &mut Vec<Event>) {
455 for open in std::mem::take(&mut self.open) {
456 Self::settle(open, out);
457 }
458 }
459
460 fn settle(open: Open, out: &mut Vec<Event>) {
463 if !open.owed {
464 return;
465 }
466 let mut end = match open.kind {
467 Kind::Text => Event::text_message_end(MessageId::new(open.id)),
468 Kind::Tool => Event::tool_call_end(ToolCallId::new(open.id)),
469 Kind::Reasoning => Event::reasoning_message_end(MessageId::new(open.id)),
470 };
471 if let Some(owner) = open.owner {
472 end.set_subagent_run_id(owner);
473 }
474 out.push(end);
475 }
476}
477
478pub fn normalize_all(events: impl IntoIterator<Item = Event>) -> Result<Vec<Event>> {
483 let mut normalizer = ChunkNormalizer::new();
484 let mut out = Vec::new();
485 for event in events {
486 normalizer.normalize(event, &mut out)?;
487 }
488 normalizer.finish(&mut out);
489 Ok(out)
490}