1mod budget_events;
2mod payload;
3
4use std::collections::{BTreeMap, VecDeque};
5use std::sync::{Arc, Mutex};
6use std::thread::ThreadId;
7
8use serde_json::Value;
9use tokio::sync::broadcast;
10
11use crate::events::{AgentErrorPayload, ApprovalAction, RunEvent, RunEventPayload, ToolStatus};
12use crate::result::RunResult;
13use crate::run_handle::{active_sub_run_ids, SharedRunResult};
14
15use payload::{agent_status, completion_reason_from_payload};
16
17const TRUSTED_STREAM_RECEIPT_KEY: &str = "_vv_agent_stream_receipt";
18const TRUSTED_STREAM_SEQUENCE_KEY: &str = "_vv_agent_stream_sequence";
19const MAX_PENDING_STREAM_RECEIPTS: usize = 256;
20const CANONICAL_STREAM_IDENTITY_FIELDS: &[&str] = &[
21 "agent_name",
22 "child_run_id",
23 "child_session_id",
24 "parent_run_id",
25 "parent_tool_call_id",
26 "run_id",
27 "session_id",
28 "sub_agent_name",
29 "task_id",
30 "trace_id",
31];
32const ASSISTANT_DELTA_FIELDS: &[&str] = &[
33 "content_chars",
34 "content_delta",
35 "delta",
36 "estimated_tokens",
37 "event",
38];
39const REASONING_DELTA_FIELDS: &[&str] = &[
40 "estimated_tokens",
41 "event",
42 "reasoning_chars",
43 "reasoning_delta",
44];
45const TOOL_STREAM_FIELDS: &[&str] = &[
46 "arguments_chars",
47 "estimated_tokens",
48 "event",
49 "function_name",
50 "tool_call_id",
51 "tool_call_index",
52];
53
54#[derive(Debug)]
55struct TrustedStreamReceipt {
56 marker: String,
57 sequence: u64,
58 fingerprint: String,
59 thread_id: ThreadId,
60}
61
62#[derive(Debug, Default)]
63struct TrustedStreamReceipts {
64 pending: VecDeque<TrustedStreamReceipt>,
65}
66
67#[doc(hidden)]
68#[derive(Clone, Debug)]
69pub struct RuntimeEventContext {
70 run_id: String,
71 trace_id: String,
72 agent_name: String,
73 session_id: Option<String>,
74 input: String,
75 trusted_stream_receipts: Arc<Mutex<TrustedStreamReceipts>>,
76}
77
78impl RuntimeEventContext {
79 pub fn new(
80 run_id: impl Into<String>,
81 trace_id: impl Into<String>,
82 agent_name: impl Into<String>,
83 session_id: Option<String>,
84 input: impl Into<String>,
85 ) -> Self {
86 Self {
87 run_id: run_id.into(),
88 trace_id: trace_id.into(),
89 agent_name: agent_name.into(),
90 session_id,
91 input: input.into(),
92 trusted_stream_receipts: Arc::new(Mutex::new(TrustedStreamReceipts::default())),
93 }
94 }
95
96 #[doc(hidden)]
97 pub fn map_stream_payload(&self, payload: &BTreeMap<String, Value>) -> Option<RunEvent> {
98 map_stream_event(payload, self)
99 }
100
101 fn attach(&self, event: RunEvent) -> RunEvent {
102 if event.session_id().is_some() {
103 return event;
104 }
105 match &self.session_id {
106 Some(session_id) => event.with_session_id(session_id),
107 None => event,
108 }
109 }
110
111 fn register_trusted_stream_receipt(
112 &self,
113 payload: &BTreeMap<String, Value>,
114 canonical: &BTreeMap<String, Value>,
115 ) -> bool {
116 let Some(marker) = payload
117 .get(TRUSTED_STREAM_RECEIPT_KEY)
118 .and_then(Value::as_str)
119 .filter(|marker| valid_stream_receipt(marker))
120 else {
121 return false;
122 };
123 let Some(sequence) = payload
124 .get(TRUSTED_STREAM_SEQUENCE_KEY)
125 .and_then(Value::as_u64)
126 .filter(|sequence| *sequence > 0)
127 else {
128 return false;
129 };
130 let Some(fingerprint) = canonical_stream_fingerprint(canonical) else {
131 return false;
132 };
133 let Ok(mut receipts) = self.trusted_stream_receipts.lock() else {
134 return false;
135 };
136 if receipts
137 .pending
138 .iter()
139 .any(|receipt| receipt.marker == marker && receipt.sequence == sequence)
140 {
141 return false;
142 }
143 while receipts.pending.len() >= MAX_PENDING_STREAM_RECEIPTS {
144 receipts.pending.pop_front();
145 }
146 receipts.pending.push_back(TrustedStreamReceipt {
147 marker: marker.to_string(),
148 sequence,
149 fingerprint,
150 thread_id: std::thread::current().id(),
151 });
152 true
153 }
154
155 fn consume_trusted_stream_receipt(&self, canonical: &BTreeMap<String, Value>) -> bool {
156 let Some(fingerprint) = canonical_stream_fingerprint(canonical) else {
157 return false;
158 };
159 let thread_id = std::thread::current().id();
160 let Ok(mut receipts) = self.trusted_stream_receipts.lock() else {
161 return false;
162 };
163 let Some(index) = receipts.pending.iter().position(|receipt| {
164 receipt.thread_id == thread_id && receipt.fingerprint == fingerprint
165 }) else {
166 return false;
167 };
168 receipts.pending.remove(index).is_some()
169 }
170}
171
172pub struct RunEventStream {
173 events: Arc<Mutex<Vec<RunEvent>>>,
174 next_index: usize,
175 receiver: Option<broadcast::Receiver<RunEvent>>,
176 shared_result: Option<SharedRunResult>,
177 completion: tokio::sync::watch::Receiver<bool>,
178}
179
180impl RunEventStream {
181 pub(crate) fn from_live(
182 receiver: Option<broadcast::Receiver<RunEvent>>,
183 result: Option<SharedRunResult>,
184 events: Arc<Mutex<Vec<RunEvent>>>,
185 completion: tokio::sync::watch::Receiver<bool>,
186 ) -> Self {
187 Self {
188 events,
189 next_index: 0,
190 receiver,
191 shared_result: result,
192 completion,
193 }
194 }
195
196 pub async fn next(&mut self) -> Option<Result<RunEvent, String>> {
197 loop {
198 if let Some(event) = self.next_journal_event() {
199 return Some(Ok(event));
200 }
201 if *self.completion.borrow() && self.active_sub_runs().is_empty() {
202 return None;
203 }
204 match self.receiver.as_mut() {
205 Some(receiver) => {
206 tokio::select! {
207 event = receiver.recv() => {
208 if matches!(event, Err(broadcast::error::RecvError::Closed)) {
209 self.receiver = None;
210 }
211 },
212 _ = self.completion.changed() => {},
213 }
214 }
215 None => {
216 if self.completion.changed().await.is_err() {
217 return self.next_journal_event().map(Ok);
218 }
219 }
220 }
221 }
222 }
223
224 fn next_journal_event(&mut self) -> Option<RunEvent> {
225 let event = self
226 .events
227 .lock()
228 .unwrap_or_else(std::sync::PoisonError::into_inner)
229 .get(self.next_index)
230 .cloned();
231 if event.is_some() {
232 self.next_index += 1;
233 }
234 event
235 }
236
237 fn active_sub_runs(&self) -> std::collections::HashSet<String> {
238 let events = self
239 .events
240 .lock()
241 .unwrap_or_else(std::sync::PoisonError::into_inner);
242 active_sub_run_ids(&events)
243 }
244
245 pub async fn into_result(mut self) -> Result<RunResult, String> {
246 if let Some(result) = self.shared_result.take() {
247 return result.wait().await;
248 }
249 Err("stream result already taken".to_string())
250 }
251}
252
253#[doc(hidden)]
254pub fn map_runtime_event(
255 event: &str,
256 payload: &std::collections::BTreeMap<String, Value>,
257 context: &RuntimeEventContext,
258) -> Option<RunEvent> {
259 let mapped = match event {
260 "sub_agent_assistant_delta"
261 | "sub_agent_reasoning_delta"
262 | "sub_agent_tool_call_started"
263 | "sub_agent_tool_call_progress" => {
264 let stream_event = event.strip_prefix("sub_agent_")?;
265 let canonical = canonical_sub_agent_stream_payload(stream_event, payload)?;
266 if !context.register_trusted_stream_receipt(payload, &canonical) {
267 return None;
268 }
269 map_canonical_sub_agent_stream_event(stream_event, &canonical)
270 }
271 "run_started" => Some(RunEvent::run_started(
272 &context.run_id,
273 &context.trace_id,
274 &context.agent_name,
275 &context.input,
276 )),
277 "cycle_started" => Some(RunEvent::cycle_started(
278 &context.run_id,
279 &context.trace_id,
280 &context.agent_name,
281 payload
282 .get("cycle")
283 .and_then(Value::as_u64)
284 .unwrap_or_default() as u32,
285 )),
286 "agent_started" => Some(RunEvent::new(
287 &context.run_id,
288 &context.trace_id,
289 &context.agent_name,
290 payload
291 .get("cycle")
292 .and_then(Value::as_u64)
293 .map(|cycle| cycle as u32),
294 RunEventPayload::AgentStarted,
295 )),
296 "llm_started" => Some(RunEvent::new(
297 &context.run_id,
298 &context.trace_id,
299 &context.agent_name,
300 payload
301 .get("cycle")
302 .and_then(Value::as_u64)
303 .map(|cycle| cycle as u32),
304 RunEventPayload::LlmStarted {
305 model: payload
306 .get("model")
307 .and_then(Value::as_str)
308 .unwrap_or_default()
309 .to_string(),
310 },
311 )),
312 "run_state_changed" => Some(RunEvent::new(
313 &context.run_id,
314 &context.trace_id,
315 &context.agent_name,
316 payload
317 .get("cycle")
318 .and_then(Value::as_u64)
319 .map(|cycle| cycle as u32),
320 RunEventPayload::RunStateChanged {
321 state: payload
322 .get("state")
323 .and_then(Value::as_str)
324 .unwrap_or_default()
325 .to_string(),
326 },
327 )),
328 "session_persisted" => Some(RunEvent::new(
329 &context.run_id,
330 &context.trace_id,
331 &context.agent_name,
332 payload
333 .get("cycle")
334 .and_then(Value::as_u64)
335 .map(|cycle| cycle as u32),
336 RunEventPayload::SessionPersisted,
337 )),
338 "budget_snapshot" => budget_events::map_budget_snapshot(payload, context),
339 "budget_exhausted" => budget_events::map_budget_exhausted(payload, context),
340 "assistant_delta" => Some(RunEvent::assistant_delta(
341 &context.run_id,
342 &context.trace_id,
343 &context.agent_name,
344 payload
345 .get("cycle")
346 .and_then(Value::as_u64)
347 .unwrap_or_default() as u32,
348 payload
349 .get("delta")
350 .or_else(|| payload.get("content_delta"))
351 .and_then(Value::as_str)
352 .unwrap_or_default(),
353 )),
354 "cycle_llm_response" => None,
358 "tool_call_started" => Some(RunEvent::tool_call_started(
359 &context.run_id,
360 &context.trace_id,
361 &context.agent_name,
362 payload
363 .get("cycle")
364 .and_then(Value::as_u64)
365 .unwrap_or_default() as u32,
366 payload
367 .get("tool_call_id")
368 .and_then(Value::as_str)
369 .unwrap_or_default(),
370 payload
371 .get("tool_name")
372 .and_then(Value::as_str)
373 .unwrap_or_default(),
374 payload
375 .get("arguments")
376 .or_else(|| payload.get("tool_arguments"))
377 .cloned()
378 .unwrap_or(Value::Null),
379 )),
380 "approval_requested" => {
381 let tool_name = payload_string(payload, "tool_name");
382 Some(with_selected_payload_metadata(
383 RunEvent::new(
384 &context.run_id,
385 &context.trace_id,
386 &context.agent_name,
387 payload
388 .get("cycle")
389 .and_then(Value::as_u64)
390 .map(|cycle| cycle as u32),
391 RunEventPayload::ApprovalRequested {
392 request_id: payload_string(payload, "request_id"),
393 tool_call_id: payload_string(payload, "tool_call_id"),
394 tool_name,
395 message: payload
396 .get("message")
397 .or_else(|| payload.get("preview"))
398 .and_then(Value::as_str)
399 .unwrap_or_default()
400 .to_string(),
401 },
402 ),
403 payload,
404 &["arguments", "tool_name"],
405 ))
406 }
407 "approval_resolved" => {
408 let approved = payload
409 .get("approved")
410 .and_then(Value::as_bool)
411 .unwrap_or(false);
412 let action = payload
413 .get("action")
414 .and_then(Value::as_str)
415 .and_then(ApprovalAction::parse)
416 .unwrap_or_else(|| ApprovalAction::from_approved(approved));
417 Some(with_selected_payload_metadata(
418 RunEvent::new(
419 &context.run_id,
420 &context.trace_id,
421 &context.agent_name,
422 payload
423 .get("cycle")
424 .and_then(Value::as_u64)
425 .map(|cycle| cycle as u32),
426 RunEventPayload::ApprovalResolved {
427 request_id: payload
428 .get("request_id")
429 .and_then(Value::as_str)
430 .unwrap_or_default()
431 .to_string(),
432 tool_name: payload
433 .get("tool_name")
434 .and_then(Value::as_str)
435 .unwrap_or_default()
436 .to_string(),
437 tool_call_id: payload
438 .get("tool_call_id")
439 .and_then(Value::as_str)
440 .unwrap_or_default()
441 .to_string(),
442 approved: action.is_approved(),
443 },
444 )
445 .with_approval_action(action),
446 payload,
447 &["action", "reason", "decision_metadata"],
448 ))
449 }
450 "sub_run_started" => {
451 let child_session_id = payload.get("child_session_id").and_then(Value::as_str);
452 let mut event = RunEvent::new(
453 child_run_id(payload).unwrap_or(&context.run_id),
454 payload
455 .get("trace_id")
456 .and_then(Value::as_str)
457 .unwrap_or(&context.trace_id),
458 payload
459 .get("agent_name")
460 .and_then(Value::as_str)
461 .unwrap_or(&context.agent_name),
462 payload
463 .get("cycle")
464 .and_then(Value::as_u64)
465 .map(|cycle| cycle as u32),
466 RunEventPayload::SubRunStarted {
467 parent_tool_call_id: payload_string(payload, "parent_tool_call_id"),
468 child_session_id: child_session_id.map(str::to_string),
469 task_id: payload
470 .get("task_id_hint")
471 .or_else(|| payload.get("task_id"))
472 .and_then(Value::as_str)
473 .map(str::to_string),
474 },
475 )
476 .with_parent_run_id(
477 payload
478 .get("parent_run_id")
479 .and_then(Value::as_str)
480 .unwrap_or(&context.run_id),
481 );
482 if let Some(session_id) = child_session_id {
483 event = event.with_session_id(session_id);
484 }
485 Some(with_nested_payload_metadata(event, payload))
486 }
487 "sub_run_completed" => {
488 let child_session_id = payload.get("child_session_id").and_then(Value::as_str);
489 let task_id = (payload.contains_key("child_run_id")
490 || payload.contains_key("child_session_id"))
491 .then(|| payload.get("task_id").and_then(Value::as_str))
492 .flatten()
493 .or_else(|| payload.get("task_id_hint").and_then(Value::as_str));
494 let mut event = RunEvent::new(
495 child_run_id(payload).unwrap_or(&context.run_id),
496 payload
497 .get("trace_id")
498 .and_then(Value::as_str)
499 .unwrap_or(&context.trace_id),
500 payload
501 .get("agent_name")
502 .and_then(Value::as_str)
503 .unwrap_or(&context.agent_name),
504 payload
505 .get("cycle")
506 .and_then(Value::as_u64)
507 .map(|cycle| cycle as u32),
508 RunEventPayload::SubRunCompleted {
509 parent_tool_call_id: payload_string(payload, "parent_tool_call_id"),
510 status: agent_status(payload),
511 final_output: payload
512 .get("final_output")
513 .and_then(Value::as_str)
514 .map(str::to_string),
515 },
516 )
517 .with_parent_run_id(
518 payload
519 .get("parent_run_id")
520 .and_then(Value::as_str)
521 .unwrap_or(&context.run_id),
522 )
523 .with_sub_run_details(
524 child_session_id,
525 task_id,
526 payload.get("wait_reason").and_then(Value::as_str),
527 payload.get("error").and_then(Value::as_str),
528 payload.get("token_usage").cloned(),
529 )
530 .with_budget_details(
531 payload
532 .get("budget_usage")
533 .and_then(|value| serde_json::from_value(value.clone()).ok())
534 .as_ref(),
535 payload
536 .get("budget_exhaustion")
537 .and_then(|value| serde_json::from_value(value.clone()).ok())
538 .as_ref(),
539 )
540 .with_completion_details(
541 completion_reason_from_payload(payload, None),
542 payload.get("completion_tool_name").and_then(Value::as_str),
543 payload.get("partial_output").and_then(Value::as_str),
544 );
545 if let Some(session_id) = child_session_id {
546 event = event.with_session_id(session_id);
547 }
548 Some(with_nested_payload_metadata(event, payload))
549 }
550 "tool_result" => {
551 let metadata = payload.get("metadata").and_then(Value::as_object);
552 if let Some(interruption_id) = metadata
553 .and_then(|metadata| metadata.get("approval_interruption_id"))
554 .and_then(Value::as_str)
555 {
556 let tool_name = payload_string(payload, "tool_name");
557 Some(RunEvent::new(
558 &context.run_id,
559 &context.trace_id,
560 &context.agent_name,
561 payload
562 .get("cycle")
563 .and_then(Value::as_u64)
564 .map(|cycle| cycle as u32),
565 RunEventPayload::ApprovalRequested {
566 request_id: interruption_id.to_string(),
567 tool_call_id: payload_string(payload, "tool_call_id"),
568 tool_name: tool_name.clone(),
569 message: metadata
570 .and_then(|metadata| metadata.get("message"))
571 .and_then(Value::as_str)
572 .map(str::to_string)
573 .unwrap_or_else(|| format!("Approval required for tool {tool_name}.")),
574 },
575 ))
576 } else if metadata
577 .and_then(|metadata| metadata.get("mode"))
578 .and_then(Value::as_str)
579 .is_some_and(|mode| mode == "handoff")
580 {
581 let metadata = metadata.expect("handoff metadata");
582 let mut event = RunEvent::new(
583 &context.run_id,
584 &context.trace_id,
585 &context.agent_name,
586 payload
587 .get("cycle")
588 .and_then(Value::as_u64)
589 .map(|cycle| cycle as u32),
590 RunEventPayload::Handoff {
591 source_agent: metadata
592 .get("handoff_from")
593 .and_then(Value::as_str)
594 .unwrap_or(&context.agent_name)
595 .to_string(),
596 target_agent: metadata
597 .get("handoff_to")
598 .and_then(Value::as_str)
599 .unwrap_or_default()
600 .to_string(),
601 tool_call_id: payload
602 .get("tool_call_id")
603 .and_then(Value::as_str)
604 .unwrap_or_default()
605 .to_string(),
606 },
607 );
608 for (key, value) in metadata {
609 event = event.with_metadata(key, value.clone());
610 }
611 Some(event)
612 } else {
613 let status = match payload
614 .get("status")
615 .and_then(Value::as_str)
616 .unwrap_or_default()
617 .to_ascii_lowercase()
618 .as_str()
619 {
620 "error" => ToolStatus::Error,
621 "wait_response" => ToolStatus::WaitResponse,
622 _ => ToolStatus::Success,
623 };
624 Some(RunEvent::tool_call_completed(
625 &context.run_id,
626 &context.trace_id,
627 &context.agent_name,
628 payload
629 .get("cycle")
630 .and_then(Value::as_u64)
631 .map(|cycle| cycle as u32),
632 payload
633 .get("tool_call_id")
634 .and_then(Value::as_str)
635 .unwrap_or_default(),
636 payload
637 .get("tool_name")
638 .and_then(Value::as_str)
639 .unwrap_or_default(),
640 status,
641 ))
642 }
643 }
644 "run_completed" => Some(
645 RunEvent::new(
646 &context.run_id,
647 &context.trace_id,
648 &context.agent_name,
649 payload
650 .get("cycle")
651 .and_then(Value::as_u64)
652 .map(|cycle| cycle as u32),
653 RunEventPayload::RunCompleted {
654 status: agent_status(payload),
655 },
656 )
657 .with_final_output(
658 payload
659 .get("final_output")
660 .or_else(|| payload.get("final_answer"))
661 .and_then(Value::as_str)
662 .map(str::to_string),
663 )
664 .with_completion_details(
665 completion_reason_from_payload(payload, None),
666 payload.get("completion_tool_name").and_then(Value::as_str),
667 payload.get("partial_output").and_then(Value::as_str),
668 ),
669 ),
670 "run_wait_user" => Some(
671 RunEvent::new(
672 &context.run_id,
673 &context.trace_id,
674 &context.agent_name,
675 payload
676 .get("cycle")
677 .and_then(Value::as_u64)
678 .map(|cycle| cycle as u32),
679 RunEventPayload::RunCompleted {
680 status: crate::types::AgentStatus::WaitUser,
681 },
682 )
683 .with_final_output(
684 payload
685 .get("wait_reason")
686 .and_then(Value::as_str)
687 .map(str::to_string),
688 )
689 .with_completion_details(
690 completion_reason_from_payload(
691 payload,
692 Some(crate::types::CompletionReason::WaitUser),
693 ),
694 payload.get("completion_tool_name").and_then(Value::as_str),
695 payload.get("partial_output").and_then(Value::as_str),
696 ),
697 ),
698 "run_cancelled" => Some(
699 RunEvent::new(
700 &context.run_id,
701 &context.trace_id,
702 &context.agent_name,
703 None,
704 RunEventPayload::RunCancelled {
705 reason: payload
706 .get("reason")
707 .and_then(Value::as_str)
708 .or_else(|| payload.get("error").and_then(Value::as_str))
709 .unwrap_or("run cancelled")
710 .to_string(),
711 },
712 )
713 .with_completion_details(
714 completion_reason_from_payload(
715 payload,
716 Some(crate::types::CompletionReason::Cancelled),
717 ),
718 None,
719 payload.get("partial_output").and_then(Value::as_str),
720 ),
721 ),
722 "run_max_cycles" => Some(
723 RunEvent::run_failed(
724 &context.run_id,
725 &context.trace_id,
726 &context.agent_name,
727 AgentErrorPayload::new(
728 payload
729 .get("error")
730 .and_then(Value::as_str)
731 .unwrap_or("run_max_cycles"),
732 ),
733 )
734 .with_completion_details(
735 completion_reason_from_payload(
736 payload,
737 Some(crate::types::CompletionReason::MaxCycles),
738 ),
739 None,
740 payload.get("partial_output").and_then(Value::as_str),
741 ),
742 ),
743 "run_failed" | "cycle_failed" => Some(
744 RunEvent::run_failed(
745 &context.run_id,
746 &context.trace_id,
747 &context.agent_name,
748 AgentErrorPayload::new(
749 payload
750 .get("error")
751 .and_then(Value::as_str)
752 .unwrap_or("cycle failed"),
753 ),
754 )
755 .with_completion_details(
756 completion_reason_from_payload(
757 payload,
758 Some(crate::types::CompletionReason::Failed),
759 ),
760 None,
761 payload.get("partial_output").and_then(Value::as_str),
762 ),
763 ),
764 _ => None,
765 };
766 mapped.map(|mapped_event| {
767 let handoff_payload = event == "tool_result"
768 && payload
769 .get("metadata")
770 .and_then(Value::as_object)
771 .and_then(|metadata| metadata.get("mode"))
772 .and_then(Value::as_str)
773 == Some("handoff");
774 let typed_metadata_payload = matches!(
775 event,
776 "approval_requested"
777 | "approval_resolved"
778 | "sub_agent_assistant_delta"
779 | "sub_agent_reasoning_delta"
780 | "sub_agent_tool_call_started"
781 | "sub_agent_tool_call_progress"
782 | "sub_run_started"
783 | "sub_run_completed"
784 | "budget_snapshot"
785 | "budget_exhausted"
786 );
787 let event = if handoff_payload || typed_metadata_payload {
788 mapped_event
789 } else {
790 with_payload_metadata(mapped_event, payload)
791 };
792 context.attach(event)
793 })
794}
795
796fn child_run_id(payload: &std::collections::BTreeMap<String, Value>) -> Option<&str> {
797 payload
798 .get("child_run_id")
799 .or_else(|| payload.get("task_id_hint"))
800 .and_then(Value::as_str)
801}
802
803fn payload_string(payload: &std::collections::BTreeMap<String, Value>, key: &str) -> String {
804 payload
805 .get(key)
806 .and_then(Value::as_str)
807 .unwrap_or_default()
808 .to_string()
809}
810
811fn with_payload_metadata(
812 mut event: RunEvent,
813 payload: &std::collections::BTreeMap<String, Value>,
814) -> RunEvent {
815 for (key, value) in payload {
816 event = event.with_metadata(key, value.clone());
817 }
818 event
819}
820
821fn with_nested_payload_metadata(
822 mut event: RunEvent,
823 payload: &std::collections::BTreeMap<String, Value>,
824) -> RunEvent {
825 if let Some(metadata) = payload.get("metadata").and_then(Value::as_object) {
826 for (key, value) in metadata {
827 event = event.with_metadata(key, value.clone());
828 }
829 }
830 event
831}
832
833fn with_selected_payload_metadata(
834 mut event: RunEvent,
835 payload: &std::collections::BTreeMap<String, Value>,
836 fields: &[&str],
837) -> RunEvent {
838 for field in fields {
839 if let Some(value) = payload.get(*field) {
840 event = event.with_metadata(*field, value.clone());
841 }
842 }
843 event
844}
845
846pub(super) fn map_stream_event(
847 payload: &std::collections::BTreeMap<String, Value>,
848 context: &RuntimeEventContext,
849) -> Option<RunEvent> {
850 let event = payload
851 .get("event")
852 .or_else(|| payload.get("type"))
853 .and_then(Value::as_str)?;
854 if let Some(canonical) = canonical_sub_agent_stream_payload(event, payload) {
855 if context.consume_trusted_stream_receipt(&canonical) {
856 return None;
857 }
858 }
859 match event {
860 "assistant_delta" => Some(
861 context.attach(RunEvent::assistant_delta(
862 &context.run_id,
863 &context.trace_id,
864 &context.agent_name,
865 payload
866 .get("cycle")
867 .and_then(Value::as_u64)
868 .unwrap_or_default() as u32,
869 payload
870 .get("delta")
871 .or_else(|| payload.get("content_delta"))
872 .and_then(Value::as_str)
873 .unwrap_or_default(),
874 )),
875 ),
876 _ => None,
877 }
878}
879
880fn map_canonical_sub_agent_stream_event(
881 stream_event: &str,
882 canonical: &BTreeMap<String, Value>,
883) -> Option<RunEvent> {
884 let payload = match stream_event {
885 "assistant_delta" => RunEventPayload::AssistantDelta {
886 delta: canonical
887 .get("delta")
888 .or_else(|| canonical.get("content_delta"))
889 .and_then(Value::as_str)
890 .unwrap_or_default()
891 .to_string(),
892 },
893 "reasoning_delta" => RunEventPayload::AssistantDelta {
896 delta: canonical
897 .get("reasoning_delta")
898 .and_then(Value::as_str)
899 .unwrap_or_default()
900 .to_string(),
901 },
902 "tool_call_started" | "tool_call_progress" => RunEventPayload::ToolCallStarted {
905 tool_call_id: canonical
906 .get("tool_call_id")
907 .and_then(Value::as_str)
908 .unwrap_or_default()
909 .to_string(),
910 tool_name: canonical
911 .get("function_name")
912 .and_then(Value::as_str)
913 .unwrap_or_default()
914 .to_string(),
915 arguments: Value::Null,
916 },
917 _ => return None,
918 };
919 let mut event = RunEvent::new(
920 canonical.get("run_id")?.as_str()?,
921 canonical.get("trace_id")?.as_str()?,
922 canonical.get("agent_name")?.as_str()?,
923 None,
924 payload,
925 )
926 .with_session_id(canonical.get("session_id")?.as_str()?)
927 .with_parent_run_id(canonical.get("parent_run_id")?.as_str()?);
928 for (key, value) in canonical {
929 event = event.with_metadata(key, value.clone());
930 }
931 Some(event)
932}
933
934fn canonical_sub_agent_stream_payload(
935 event: &str,
936 payload: &BTreeMap<String, Value>,
937) -> Option<BTreeMap<String, Value>> {
938 let producer_fields = match event {
939 "assistant_delta" => ASSISTANT_DELTA_FIELDS,
940 "reasoning_delta" => REASONING_DELTA_FIELDS,
941 "tool_call_started" | "tool_call_progress" => TOOL_STREAM_FIELDS,
942 _ => return None,
943 };
944 let mut canonical = payload
945 .iter()
946 .filter(|(key, _)| {
947 producer_fields.contains(&key.as_str())
948 || CANONICAL_STREAM_IDENTITY_FIELDS.contains(&key.as_str())
949 })
950 .map(|(key, value)| (key.clone(), value.clone()))
951 .collect::<BTreeMap<_, _>>();
952 canonical.insert("event".to_string(), Value::String(event.to_string()));
953 if !CANONICAL_STREAM_IDENTITY_FIELDS
954 .iter()
955 .all(|key| canonical.get(*key).is_some_and(Value::is_string))
956 {
957 return None;
958 }
959 for (left, right) in [
960 ("run_id", "child_run_id"),
961 ("session_id", "child_session_id"),
962 ("agent_name", "sub_agent_name"),
963 ] {
964 if canonical.get(left) != canonical.get(right) {
965 return None;
966 }
967 }
968 Some(canonical)
969}
970
971fn canonical_stream_fingerprint(payload: &BTreeMap<String, Value>) -> Option<String> {
972 serde_json::to_string(payload).ok()
973}
974
975fn valid_stream_receipt(marker: &str) -> bool {
976 marker.strip_prefix("stream_").is_some_and(|value| {
977 value.len() == 32 && value.bytes().all(|byte| byte.is_ascii_hexdigit())
978 })
979}
980
981#[cfg(test)]
982mod tests;