Skip to main content

vv_agent/runner/
event_stream.rs

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        // This is a complete cycle record, not a streaming token delta. The v1
355        // typed payload has no cycle-response variant, so keep it out of the
356        // assistant_delta channel instead of duplicating the full answer.
357        "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        // RunEvent v1 has one typed text-delta envelope. Preserve the exact
894        // producer event and fields in metadata for reasoning stream consumers.
895        "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        // The v1 tool envelope has no progress variant. The canonical metadata
903        // remains the authoritative started/progress discriminator.
904        "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;