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 "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;