Skip to main content

vv_agent/runner/
event_stream.rs

1mod budget_events;
2mod memory_events;
3mod payload;
4mod stream_projection;
5
6use std::collections::{BTreeMap, HashSet};
7use std::sync::{Arc, Mutex};
8
9use serde_json::Value;
10use tokio::sync::broadcast;
11
12use crate::events::{AgentErrorPayload, ApprovalAction, RunEvent, RunEventPayload, ToolStatus};
13use crate::result::RunResult;
14use crate::run_handle::{active_sub_run_ids, SharedRunResult};
15use crate::tools::ToolMetadata;
16
17use payload::{agent_status, completion_reason_from_payload};
18pub(crate) use stream_projection::map_stream_event;
19
20#[doc(hidden)]
21#[derive(Clone, Debug)]
22pub struct RuntimeEventContext {
23    run_id: String,
24    trace_id: String,
25    agent_name: String,
26    session_id: Option<String>,
27    input: String,
28    observed_tool_completions: Arc<Mutex<HashSet<String>>>,
29}
30
31impl RuntimeEventContext {
32    pub fn new(
33        run_id: impl Into<String>,
34        trace_id: impl Into<String>,
35        agent_name: impl Into<String>,
36        session_id: Option<String>,
37        input: impl Into<String>,
38    ) -> Self {
39        Self {
40            run_id: run_id.into(),
41            trace_id: trace_id.into(),
42            agent_name: agent_name.into(),
43            session_id,
44            input: input.into(),
45            observed_tool_completions: Arc::new(Mutex::new(HashSet::new())),
46        }
47    }
48
49    #[doc(hidden)]
50    pub fn map_stream_payload(&self, payload: &BTreeMap<String, Value>) -> Option<RunEvent> {
51        map_stream_event(payload, self)
52    }
53
54    fn attach(&self, event: RunEvent) -> RunEvent {
55        if event.session_id().is_some() {
56            return event;
57        }
58        match &self.session_id {
59            Some(session_id) => event.with_session_id(session_id),
60            None => event,
61        }
62    }
63}
64
65pub struct RunEventStream {
66    events: Arc<Mutex<Vec<RunEvent>>>,
67    next_index: usize,
68    receiver: Option<broadcast::Receiver<RunEvent>>,
69    shared_result: Option<SharedRunResult>,
70    completion: tokio::sync::watch::Receiver<bool>,
71}
72
73impl RunEventStream {
74    pub(crate) fn from_live(
75        receiver: Option<broadcast::Receiver<RunEvent>>,
76        result: Option<SharedRunResult>,
77        events: Arc<Mutex<Vec<RunEvent>>>,
78        completion: tokio::sync::watch::Receiver<bool>,
79    ) -> Self {
80        Self {
81            events,
82            next_index: 0,
83            receiver,
84            shared_result: result,
85            completion,
86        }
87    }
88
89    pub async fn next(&mut self) -> Option<Result<RunEvent, String>> {
90        loop {
91            if let Some(event) = self.next_journal_event() {
92                return Some(Ok(event));
93            }
94            if *self.completion.borrow() && self.active_sub_runs().is_empty() {
95                return None;
96            }
97            match self.receiver.as_mut() {
98                Some(receiver) => {
99                    tokio::select! {
100                        event = receiver.recv() => {
101                            if matches!(event, Err(broadcast::error::RecvError::Closed)) {
102                                self.receiver = None;
103                            }
104                        },
105                        _ = self.completion.changed() => {},
106                    }
107                }
108                None => {
109                    if self.completion.changed().await.is_err() {
110                        return self.next_journal_event().map(Ok);
111                    }
112                }
113            }
114        }
115    }
116
117    fn next_journal_event(&mut self) -> Option<RunEvent> {
118        let event = self
119            .events
120            .lock()
121            .unwrap_or_else(std::sync::PoisonError::into_inner)
122            .get(self.next_index)
123            .cloned();
124        if event.is_some() {
125            self.next_index += 1;
126        }
127        event
128    }
129
130    fn active_sub_runs(&self) -> std::collections::HashSet<String> {
131        let events = self
132            .events
133            .lock()
134            .unwrap_or_else(std::sync::PoisonError::into_inner);
135        active_sub_run_ids(&events)
136    }
137
138    pub async fn into_result(mut self) -> Result<RunResult, String> {
139        if let Some(result) = self.shared_result.take() {
140            return result.wait().await;
141        }
142        Err("stream result already taken".to_string())
143    }
144}
145
146#[doc(hidden)]
147pub fn map_runtime_event(
148    event: &str,
149    payload: &std::collections::BTreeMap<String, Value>,
150    context: &RuntimeEventContext,
151) -> Option<RunEvent> {
152    let mapped = match event {
153        "run_started" => Some(RunEvent::run_started(
154            &context.run_id,
155            &context.trace_id,
156            &context.agent_name,
157            &context.input,
158        )),
159        "cycle_started" => Some(RunEvent::cycle_started(
160            &context.run_id,
161            &context.trace_id,
162            &context.agent_name,
163            payload
164                .get("cycle")
165                .and_then(Value::as_u64)
166                .unwrap_or_default() as u32,
167        )),
168        "agent_started" => Some(RunEvent::new(
169            &context.run_id,
170            &context.trace_id,
171            &context.agent_name,
172            payload
173                .get("cycle")
174                .and_then(Value::as_u64)
175                .map(|cycle| cycle as u32),
176            RunEventPayload::AgentStarted,
177        )),
178        "run_state_changed" => Some(RunEvent::new(
179            &context.run_id,
180            &context.trace_id,
181            &context.agent_name,
182            payload
183                .get("cycle")
184                .and_then(Value::as_u64)
185                .map(|cycle| cycle as u32),
186            RunEventPayload::RunStateChanged {
187                state: payload
188                    .get("state")
189                    .and_then(Value::as_str)
190                    .unwrap_or_default()
191                    .to_string(),
192            },
193        )),
194        "session_persisted" => Some(RunEvent::new(
195            &context.run_id,
196            &context.trace_id,
197            &context.agent_name,
198            payload
199                .get("cycle")
200                .and_then(Value::as_u64)
201                .map(|cycle| cycle as u32),
202            RunEventPayload::SessionPersisted,
203        )),
204        "budget_snapshot" => budget_events::map_budget_snapshot(payload, context),
205        "budget_exhausted" => budget_events::map_budget_exhausted(payload, context),
206        "memory_compact_started" => memory_events::map_memory_compact_started(payload, context),
207        "memory_compact_completed" => memory_events::map_memory_compact_completed(payload, context),
208        "assistant_delta" => Some(RunEvent::assistant_delta(
209            &context.run_id,
210            &context.trace_id,
211            &context.agent_name,
212            payload
213                .get("cycle")
214                .and_then(Value::as_u64)
215                .unwrap_or_default() as u32,
216            payload
217                .get("delta")
218                .or_else(|| payload.get("content_delta"))
219                .and_then(Value::as_str)
220                .unwrap_or_default(),
221        )),
222        // This is a complete cycle record, not a streaming token delta. The v1
223        // typed payload has no cycle-response variant, so keep it out of the
224        // assistant_delta channel instead of duplicating the full answer.
225        "cycle_llm_response" => None,
226        "tool_call_planned" => map_runtime_tool_call(payload, context, true),
227        "tool_call_started" => map_runtime_tool_call(payload, context, false),
228        "tool_call_completed" => {
229            let event = map_runtime_tool_completion(payload, context)?;
230            if let Ok(mut observed) = context.observed_tool_completions.lock() {
231                observed.insert(tool_completion_key(payload));
232            }
233            Some(event)
234        }
235        "approval_requested" => {
236            let tool_name = payload_string(payload, "tool_name");
237            Some(with_selected_payload_metadata(
238                RunEvent::new(
239                    &context.run_id,
240                    &context.trace_id,
241                    &context.agent_name,
242                    payload
243                        .get("cycle")
244                        .and_then(Value::as_u64)
245                        .map(|cycle| cycle as u32),
246                    RunEventPayload::ApprovalRequested {
247                        request_id: payload_string_non_empty(payload, "request_id")?.to_string(),
248                        tool_call_id: payload_string_non_empty(payload, "tool_call_id")?
249                            .to_string(),
250                        tool_name,
251                        message: payload.get("message")?.as_str()?.to_string(),
252                    },
253                ),
254                payload,
255                &["arguments", "tool_name"],
256            ))
257        }
258        "approval_resolved" => {
259            let action = payload
260                .get("action")
261                .and_then(Value::as_str)
262                .and_then(ApprovalAction::parse)?;
263            Some(with_selected_payload_metadata(
264                RunEvent::new(
265                    &context.run_id,
266                    &context.trace_id,
267                    &context.agent_name,
268                    payload
269                        .get("cycle")
270                        .and_then(Value::as_u64)
271                        .map(|cycle| cycle as u32),
272                    RunEventPayload::ApprovalResolved {
273                        request_id: payload_string_non_empty(payload, "request_id")?.to_string(),
274                        tool_name: payload_string_non_empty(payload, "tool_name")?.to_string(),
275                        tool_call_id: payload_string_non_empty(payload, "tool_call_id")?
276                            .to_string(),
277                        action,
278                    },
279                ),
280                payload,
281                &["reason", "decision_metadata"],
282            ))
283        }
284        "sub_run_started" => {
285            let child_session_id = payload.get("child_session_id").and_then(Value::as_str);
286            let mut event = RunEvent::new(
287                child_run_id(payload).unwrap_or(&context.run_id),
288                payload
289                    .get("trace_id")
290                    .and_then(Value::as_str)
291                    .unwrap_or(&context.trace_id),
292                payload
293                    .get("agent_name")
294                    .and_then(Value::as_str)
295                    .unwrap_or(&context.agent_name),
296                payload
297                    .get("cycle")
298                    .and_then(Value::as_u64)
299                    .map(|cycle| cycle as u32),
300                RunEventPayload::SubRunStarted {
301                    parent_tool_call_id: payload_string(payload, "parent_tool_call_id"),
302                    child_session_id: child_session_id.map(str::to_string),
303                    task_id: payload
304                        .get("task_id_hint")
305                        .or_else(|| payload.get("task_id"))
306                        .and_then(Value::as_str)
307                        .map(str::to_string),
308                },
309            )
310            .with_parent_run_id(
311                payload
312                    .get("parent_run_id")
313                    .and_then(Value::as_str)
314                    .unwrap_or(&context.run_id),
315            );
316            if let Some(session_id) = child_session_id {
317                event = event.with_session_id(session_id);
318            }
319            Some(with_nested_payload_metadata(event, payload))
320        }
321        "sub_run_completed" => {
322            let child_session_id = payload.get("child_session_id").and_then(Value::as_str);
323            let task_id = (payload.contains_key("child_run_id")
324                || payload.contains_key("child_session_id"))
325            .then(|| payload.get("task_id").and_then(Value::as_str))
326            .flatten()
327            .or_else(|| payload.get("task_id_hint").and_then(Value::as_str));
328            let mut event = RunEvent::new(
329                child_run_id(payload).unwrap_or(&context.run_id),
330                payload
331                    .get("trace_id")
332                    .and_then(Value::as_str)
333                    .unwrap_or(&context.trace_id),
334                payload
335                    .get("agent_name")
336                    .and_then(Value::as_str)
337                    .unwrap_or(&context.agent_name),
338                payload
339                    .get("cycle")
340                    .and_then(Value::as_u64)
341                    .map(|cycle| cycle as u32),
342                RunEventPayload::SubRunCompleted {
343                    parent_tool_call_id: payload_string(payload, "parent_tool_call_id"),
344                    status: agent_status(payload),
345                    final_output: payload
346                        .get("final_output")
347                        .and_then(Value::as_str)
348                        .map(str::to_string),
349                },
350            )
351            .with_parent_run_id(
352                payload
353                    .get("parent_run_id")
354                    .and_then(Value::as_str)
355                    .unwrap_or(&context.run_id),
356            )
357            .with_sub_run_details(
358                child_session_id,
359                task_id,
360                payload.get("wait_reason").and_then(Value::as_str),
361                payload.get("error").and_then(Value::as_str),
362                payload.get("token_usage").cloned(),
363            )
364            .with_budget_details(
365                payload
366                    .get("budget_usage")
367                    .and_then(|value| serde_json::from_value(value.clone()).ok())
368                    .as_ref(),
369                payload
370                    .get("budget_exhaustion")
371                    .and_then(|value| serde_json::from_value(value.clone()).ok())
372                    .as_ref(),
373            )
374            .with_completion_details(
375                completion_reason_from_payload(payload, None),
376                payload.get("completion_tool_name").and_then(Value::as_str),
377                payload.get("partial_output").and_then(Value::as_str),
378            );
379            if let Some(session_id) = child_session_id {
380                event = event.with_session_id(session_id);
381            }
382            Some(with_nested_payload_metadata(event, payload))
383        }
384        "tool_result" => {
385            if payload
386                .get("lifecycle_suppressed")
387                .and_then(Value::as_bool)
388                .unwrap_or(false)
389            {
390                return None;
391            }
392            let metadata = payload.get("metadata").and_then(Value::as_object);
393            if let Some(interruption_id) = metadata
394                .and_then(|metadata| metadata.get("approval_interruption_id"))
395                .and_then(Value::as_str)
396            {
397                let tool_name = payload_string(payload, "tool_name");
398                Some(RunEvent::new(
399                    &context.run_id,
400                    &context.trace_id,
401                    &context.agent_name,
402                    payload
403                        .get("cycle")
404                        .and_then(Value::as_u64)
405                        .map(|cycle| cycle as u32),
406                    RunEventPayload::ApprovalRequested {
407                        request_id: interruption_id.to_string(),
408                        tool_call_id: payload_string(payload, "tool_call_id"),
409                        tool_name: tool_name.clone(),
410                        message: metadata
411                            .and_then(|metadata| metadata.get("message"))
412                            .and_then(Value::as_str)
413                            .map(str::to_string)
414                            .unwrap_or_else(|| format!("Approval required for tool {tool_name}.")),
415                    },
416                ))
417            } else if metadata
418                .and_then(|metadata| metadata.get("mode"))
419                .and_then(Value::as_str)
420                .is_some_and(|mode| mode == "handoff")
421            {
422                None
423            } else {
424                let already_observed = context
425                    .observed_tool_completions
426                    .lock()
427                    .map(|observed| observed.contains(&tool_completion_key(payload)))
428                    .unwrap_or(false);
429                (!already_observed)
430                    .then(|| map_runtime_tool_completion(payload, context))
431                    .flatten()
432            }
433        }
434        "run_completed" => Some(
435            RunEvent::new(
436                &context.run_id,
437                &context.trace_id,
438                &context.agent_name,
439                payload
440                    .get("cycle")
441                    .and_then(Value::as_u64)
442                    .map(|cycle| cycle as u32),
443                RunEventPayload::RunCompleted {
444                    status: agent_status(payload),
445                },
446            )
447            .with_final_output(
448                payload
449                    .get("final_output")
450                    .or_else(|| payload.get("final_answer"))
451                    .and_then(Value::as_str)
452                    .map(str::to_string),
453            )
454            .with_completion_details(
455                completion_reason_from_payload(payload, None),
456                payload.get("completion_tool_name").and_then(Value::as_str),
457                payload.get("partial_output").and_then(Value::as_str),
458            ),
459        ),
460        "run_wait_user" => Some(
461            RunEvent::new(
462                &context.run_id,
463                &context.trace_id,
464                &context.agent_name,
465                payload
466                    .get("cycle")
467                    .and_then(Value::as_u64)
468                    .map(|cycle| cycle as u32),
469                RunEventPayload::RunCompleted {
470                    status: crate::types::AgentStatus::WaitUser,
471                },
472            )
473            .with_final_output(
474                payload
475                    .get("wait_reason")
476                    .and_then(Value::as_str)
477                    .map(str::to_string),
478            )
479            .with_completion_details(
480                completion_reason_from_payload(
481                    payload,
482                    Some(crate::types::CompletionReason::WaitUser),
483                ),
484                payload.get("completion_tool_name").and_then(Value::as_str),
485                payload.get("partial_output").and_then(Value::as_str),
486            ),
487        ),
488        "run_cancelled" => Some(
489            RunEvent::new(
490                &context.run_id,
491                &context.trace_id,
492                &context.agent_name,
493                None,
494                RunEventPayload::RunCancelled {
495                    reason: payload
496                        .get("reason")
497                        .and_then(Value::as_str)
498                        .or_else(|| payload.get("error").and_then(Value::as_str))
499                        .unwrap_or("run cancelled")
500                        .to_string(),
501                },
502            )
503            .with_completion_details(
504                completion_reason_from_payload(
505                    payload,
506                    Some(crate::types::CompletionReason::Cancelled),
507                ),
508                None,
509                payload.get("partial_output").and_then(Value::as_str),
510            ),
511        ),
512        "run_max_cycles" => Some(
513            RunEvent::run_failed(
514                &context.run_id,
515                &context.trace_id,
516                &context.agent_name,
517                AgentErrorPayload::new(
518                    payload
519                        .get("error")
520                        .and_then(Value::as_str)
521                        .unwrap_or("run_max_cycles"),
522                ),
523            )
524            .with_completion_details(
525                completion_reason_from_payload(
526                    payload,
527                    Some(crate::types::CompletionReason::MaxCycles),
528                ),
529                None,
530                payload.get("partial_output").and_then(Value::as_str),
531            ),
532        ),
533        "run_failed" | "cycle_failed" => Some(
534            RunEvent::run_failed(
535                &context.run_id,
536                &context.trace_id,
537                &context.agent_name,
538                AgentErrorPayload::new(
539                    payload
540                        .get("error")
541                        .and_then(Value::as_str)
542                        .unwrap_or("cycle failed"),
543                ),
544            )
545            .with_completion_details(
546                completion_reason_from_payload(
547                    payload,
548                    Some(crate::types::CompletionReason::Failed),
549                ),
550                None,
551                payload.get("partial_output").and_then(Value::as_str),
552            ),
553        ),
554        _ => None,
555    };
556    mapped.map(|mapped_event| {
557        let handoff_payload = event == "tool_result"
558            && payload
559                .get("metadata")
560                .and_then(Value::as_object)
561                .and_then(|metadata| metadata.get("mode"))
562                .and_then(Value::as_str)
563                == Some("handoff");
564        let typed_metadata_payload = matches!(
565            event,
566            "approval_requested"
567                | "approval_resolved"
568                | "sub_agent_assistant_delta"
569                | "sub_agent_reasoning_delta"
570                | "sub_agent_tool_call_started"
571                | "sub_agent_tool_call_progress"
572                | "sub_run_started"
573                | "sub_run_completed"
574                | "memory_compact_started"
575                | "memory_compact_completed"
576                | "budget_snapshot"
577                | "budget_exhausted"
578        );
579        let event = if handoff_payload || typed_metadata_payload {
580            mapped_event
581        } else {
582            with_payload_metadata(mapped_event, payload)
583        };
584        context.attach(event)
585    })
586}
587
588fn map_runtime_tool_call(
589    payload: &BTreeMap<String, Value>,
590    context: &RuntimeEventContext,
591    planned: bool,
592) -> Option<RunEvent> {
593    let tool_call_id = payload_string_non_empty(payload, "tool_call_id")?;
594    let tool_name = payload_string_non_empty(payload, "tool_name")?;
595    let arguments = payload
596        .get("arguments")
597        .or_else(|| payload.get("tool_arguments"))?
598        .clone();
599    if !arguments.is_object() {
600        return None;
601    }
602    let tool_metadata = runtime_tool_metadata(payload)?;
603    let cycle_index = payload
604        .get("cycle")
605        .and_then(Value::as_u64)
606        .unwrap_or_default() as u32;
607    let event = if planned {
608        RunEvent::tool_call_planned(
609            &context.run_id,
610            &context.trace_id,
611            &context.agent_name,
612            cycle_index,
613            tool_call_id,
614            tool_name,
615            arguments,
616        )
617    } else {
618        RunEvent::tool_call_started(
619            &context.run_id,
620            &context.trace_id,
621            &context.agent_name,
622            cycle_index,
623            tool_call_id,
624            tool_name,
625            arguments,
626        )
627    };
628    Some(event.with_tool_metadata(tool_metadata.as_ref()))
629}
630
631fn map_runtime_tool_completion(
632    payload: &BTreeMap<String, Value>,
633    context: &RuntimeEventContext,
634) -> Option<RunEvent> {
635    const JSON_SAFE_INTEGER_MAX: u64 = (1_u64 << 53) - 1;
636
637    let status = runtime_tool_status(payload.get("status")?.as_str()?)?;
638    let tool_call_id = payload_string_non_empty(payload, "tool_call_id")?;
639    let tool_name = payload_string_non_empty(payload, "tool_name")?;
640    let tool_metadata = runtime_tool_metadata(payload)?;
641    let directive =
642        serde_json::from_value::<crate::types::ToolDirective>(payload.get("directive")?.clone())
643            .ok()?;
644    let error_code = match payload.get("error_code")? {
645        Value::Null => None,
646        Value::String(value) => Some(value.clone()),
647        _ => return None,
648    };
649    let execution_started = payload.get("execution_started")?.as_bool()?;
650    let duration_ms = match payload.get("duration_ms")? {
651        Value::Null => None,
652        value => Some(
653            value
654                .as_u64()
655                .filter(|value| *value <= JSON_SAFE_INTEGER_MAX)?,
656        ),
657    };
658    if !execution_started && duration_ms.is_some() {
659        return None;
660    }
661    let event = RunEvent::new(
662        &context.run_id,
663        &context.trace_id,
664        &context.agent_name,
665        payload
666            .get("cycle")
667            .and_then(Value::as_u64)
668            .map(|cycle| cycle as u32),
669        RunEventPayload::ToolCallCompleted {
670            tool_call_id: tool_call_id.to_string(),
671            tool_name: tool_name.to_string(),
672            status,
673            directive,
674            error_code,
675            execution_started,
676            duration_ms,
677        },
678    )
679    .with_tool_metadata(tool_metadata.as_ref());
680    Some(event)
681}
682
683fn runtime_tool_status(status: &str) -> Option<ToolStatus> {
684    match status {
685        "success" => Some(ToolStatus::Success),
686        "error" => Some(ToolStatus::Error),
687        "wait_response" => Some(ToolStatus::WaitResponse),
688        "running" => Some(ToolStatus::Running),
689        "pending_compress" => Some(ToolStatus::PendingCompress),
690        _ => None,
691    }
692}
693
694fn runtime_tool_metadata(payload: &BTreeMap<String, Value>) -> Option<Option<ToolMetadata>> {
695    payload
696        .get("tool_metadata")
697        .map(|value| serde_json::from_value(value.clone()).ok().map(Some))
698        .unwrap_or(Some(None))
699}
700
701fn payload_string_non_empty<'a>(
702    payload: &'a BTreeMap<String, Value>,
703    field: &str,
704) -> Option<&'a str> {
705    payload
706        .get(field)
707        .and_then(Value::as_str)
708        .map(str::trim)
709        .filter(|value| !value.is_empty())
710}
711
712fn tool_completion_key(payload: &BTreeMap<String, Value>) -> String {
713    format!(
714        "{}\0{}",
715        payload
716            .get("cycle")
717            .and_then(Value::as_u64)
718            .unwrap_or_default(),
719        payload
720            .get("tool_call_id")
721            .and_then(Value::as_str)
722            .unwrap_or_default()
723    )
724}
725
726fn child_run_id(payload: &std::collections::BTreeMap<String, Value>) -> Option<&str> {
727    payload
728        .get("child_run_id")
729        .or_else(|| payload.get("task_id_hint"))
730        .and_then(Value::as_str)
731}
732
733fn payload_string(payload: &std::collections::BTreeMap<String, Value>, key: &str) -> String {
734    payload
735        .get(key)
736        .and_then(Value::as_str)
737        .unwrap_or_default()
738        .to_string()
739}
740
741fn with_payload_metadata(
742    mut event: RunEvent,
743    payload: &std::collections::BTreeMap<String, Value>,
744) -> RunEvent {
745    for (key, value) in payload {
746        event = event.with_metadata(key, value.clone());
747    }
748    event
749}
750
751fn with_nested_payload_metadata(
752    mut event: RunEvent,
753    payload: &std::collections::BTreeMap<String, Value>,
754) -> RunEvent {
755    if let Some(metadata) = payload.get("metadata").and_then(Value::as_object) {
756        for (key, value) in metadata {
757            event = event.with_metadata(key, value.clone());
758        }
759    }
760    event
761}
762
763fn with_selected_payload_metadata(
764    mut event: RunEvent,
765    payload: &std::collections::BTreeMap<String, Value>,
766    fields: &[&str],
767) -> RunEvent {
768    for field in fields {
769        if let Some(value) = payload.get(*field) {
770            event = event.with_metadata(*field, value.clone());
771        }
772    }
773    event
774}
775
776#[cfg(test)]
777mod tests;