Skip to main content

agent_base/types/
events.rs

1use serde::{Deserialize, Serialize};
2use serde_json::Value;
3
4use super::approval::ApprovalRequest;
5use super::checkpoint::CheckpointData;
6use super::plan_update::PlanItem;
7use super::session::SessionId;
8
9// ---------------------------------------------------------------------------
10// UserEvent — user-space events produced by tools
11// ---------------------------------------------------------------------------
12
13/// User-space events produced by tools during execution.
14///
15/// Tools send these through `ToolContext::emit_user_event()`. The framework
16/// wraps them in [`RuntimeEvent::UserEvent`] before delivering to external
17/// consumers.
18#[derive(Clone, Debug, Serialize, Deserialize)]
19#[serde(tag = "userEventType", rename_all = "camelCase")]
20pub enum UserEvent {
21    /// Tool progress notification.
22    Progress { text: String },
23    /// User-defined structured event for custom business semantics.
24    Structured { event_type: String, data: Value },
25    /// Tool partial result — emitted during long-running tool execution.
26    /// `is_partial: true` means more output is coming; `false` means this is the final chunk.
27    ToolPartialResult {
28        tool_call_id: String,
29        content: String,
30        is_partial: bool,
31    },
32}
33
34// ---------------------------------------------------------------------------
35// RuntimeEvent — unified event stream for all consumers
36// ---------------------------------------------------------------------------
37
38/// Unified runtime event — the single event type for both internal and external
39/// consumers (frontends, CLIs, tests).
40#[derive(Clone, Debug, Serialize, Deserialize)]
41#[serde(tag = "runtimeEventType", rename_all = "camelCase")]
42pub enum RuntimeEvent {
43    // --- Framework system events ---
44    TextDelta {
45        session_id: SessionId,
46        text: String,
47        /// Identifies the agent that produced this event (root / sub-agent path).
48        #[serde(default, skip_serializing_if = "Option::is_none")]
49        agent_id: Option<String>,
50        /// Distributed tracing context carried across systems (e.g. MCP caller → phi-agent).
51        #[serde(default, skip_serializing_if = "Option::is_none")]
52        trace_id: Option<String>,
53    },
54    ThoughtDelta {
55        session_id: SessionId,
56        text: String,
57        #[serde(default, skip_serializing_if = "Option::is_none")]
58        agent_id: Option<String>,
59        #[serde(default, skip_serializing_if = "Option::is_none")]
60        trace_id: Option<String>,
61    },
62    ToolCallStarted {
63        session_id: SessionId,
64        tool_name: String,
65        args_json: String,
66        #[serde(default, skip_serializing_if = "Option::is_none")]
67        agent_id: Option<String>,
68        #[serde(default, skip_serializing_if = "Option::is_none")]
69        trace_id: Option<String>,
70    },
71    ToolCallFinished {
72        session_id: SessionId,
73        tool_name: String,
74        summary: String,
75        #[serde(default, skip_serializing_if = "Option::is_none")]
76        agent_id: Option<String>,
77        #[serde(default, skip_serializing_if = "Option::is_none")]
78        trace_id: Option<String>,
79        /// `true` when this finish is the result of an approval denial (not a
80        /// normal tool result or execution error).
81        #[serde(default)]
82        denied: bool,
83    },
84    AwaitingApproval {
85        session_id: SessionId,
86        request: ApprovalRequest,
87        #[serde(default, skip_serializing_if = "Option::is_none")]
88        agent_id: Option<String>,
89        #[serde(default, skip_serializing_if = "Option::is_none")]
90        trace_id: Option<String>,
91    },
92    Checkpoint {
93        session_id: SessionId,
94        checkpoint: CheckpointData,
95        #[serde(default, skip_serializing_if = "Option::is_none")]
96        agent_id: Option<String>,
97        #[serde(default, skip_serializing_if = "Option::is_none")]
98        trace_id: Option<String>,
99    },
100    RunFinished {
101        session_id: SessionId,
102        #[serde(default, skip_serializing_if = "Option::is_none")]
103        agent_id: Option<String>,
104        #[serde(default, skip_serializing_if = "Option::is_none")]
105        trace_id: Option<String>,
106    },
107    RunCancelled {
108        session_id: SessionId,
109        #[serde(default, skip_serializing_if = "Option::is_none")]
110        agent_id: Option<String>,
111        #[serde(default, skip_serializing_if = "Option::is_none")]
112        trace_id: Option<String>,
113    },
114    // --- Lightweight plan update (display-only, no execution semantics) ---
115    PlanUpdated {
116        session_id: SessionId,
117        objective: String,
118        explanation: Option<String>,
119        plan: Vec<PlanItem>,
120        #[serde(default, skip_serializing_if = "Option::is_none")]
121        agent_id: Option<String>,
122        #[serde(default, skip_serializing_if = "Option::is_none")]
123        trace_id: Option<String>,
124    },
125    // --- User-space events ---
126    /// A user-space event produced by a tool.
127    UserEvent {
128        session_id: SessionId,
129        event: UserEvent,
130        #[serde(default, skip_serializing_if = "Option::is_none")]
131        agent_id: Option<String>,
132        #[serde(default, skip_serializing_if = "Option::is_none")]
133        trace_id: Option<String>,
134    },
135}
136
137impl RuntimeEvent {
138    /// Get the session ID associated with this event.
139    pub fn session_id(&self) -> &SessionId {
140        match self {
141            RuntimeEvent::TextDelta { session_id, .. } => session_id,
142            RuntimeEvent::ThoughtDelta { session_id, .. } => session_id,
143            RuntimeEvent::ToolCallStarted { session_id, .. } => session_id,
144            RuntimeEvent::ToolCallFinished { session_id, .. } => session_id,
145            RuntimeEvent::AwaitingApproval { session_id, .. } => session_id,
146            RuntimeEvent::Checkpoint { session_id, .. } => session_id,
147            RuntimeEvent::RunFinished { session_id, .. } => session_id,
148            RuntimeEvent::RunCancelled { session_id, .. } => session_id,
149            RuntimeEvent::PlanUpdated { session_id, .. } => session_id,
150            RuntimeEvent::UserEvent { session_id, .. } => session_id,
151        }
152    }
153
154    /// Get the agent ID associated with this event, if any.
155    pub fn agent_id(&self) -> Option<&str> {
156        match self {
157            RuntimeEvent::TextDelta { agent_id, .. } => agent_id.as_deref(),
158            RuntimeEvent::ThoughtDelta { agent_id, .. } => agent_id.as_deref(),
159            RuntimeEvent::ToolCallStarted { agent_id, .. } => agent_id.as_deref(),
160            RuntimeEvent::ToolCallFinished { agent_id, .. } => agent_id.as_deref(),
161            RuntimeEvent::AwaitingApproval { agent_id, .. } => agent_id.as_deref(),
162            RuntimeEvent::Checkpoint { agent_id, .. } => agent_id.as_deref(),
163            RuntimeEvent::RunFinished { agent_id, .. } => agent_id.as_deref(),
164            RuntimeEvent::RunCancelled { agent_id, .. } => agent_id.as_deref(),
165            RuntimeEvent::PlanUpdated { agent_id, .. } => agent_id.as_deref(),
166            RuntimeEvent::UserEvent { agent_id, .. } => agent_id.as_deref(),
167        }
168    }
169
170    /// Get the trace ID associated with this event, if any.
171    pub fn trace_id(&self) -> Option<&str> {
172        match self {
173            RuntimeEvent::TextDelta { trace_id, .. } => trace_id.as_deref(),
174            RuntimeEvent::ThoughtDelta { trace_id, .. } => trace_id.as_deref(),
175            RuntimeEvent::ToolCallStarted { trace_id, .. } => trace_id.as_deref(),
176            RuntimeEvent::ToolCallFinished { trace_id, .. } => trace_id.as_deref(),
177            RuntimeEvent::AwaitingApproval { trace_id, .. } => trace_id.as_deref(),
178            RuntimeEvent::Checkpoint { trace_id, .. } => trace_id.as_deref(),
179            RuntimeEvent::RunFinished { trace_id, .. } => trace_id.as_deref(),
180            RuntimeEvent::RunCancelled { trace_id, .. } => trace_id.as_deref(),
181            RuntimeEvent::PlanUpdated { trace_id, .. } => trace_id.as_deref(),
182            RuntimeEvent::UserEvent { trace_id, .. } => trace_id.as_deref(),
183        }
184    }
185
186    /// Set the agent ID (sub-agent path) on this event.
187    pub fn with_agent_id(mut self, id: impl Into<String>) -> Self {
188        let id = id.into();
189        match &mut self {
190            RuntimeEvent::TextDelta { agent_id, .. }
191            | RuntimeEvent::ThoughtDelta { agent_id, .. }
192            | RuntimeEvent::ToolCallStarted { agent_id, .. }
193            | RuntimeEvent::ToolCallFinished { agent_id, .. }
194            | RuntimeEvent::AwaitingApproval { agent_id, .. }
195            | RuntimeEvent::Checkpoint { agent_id, .. }
196            | RuntimeEvent::RunFinished { agent_id, .. }
197            | RuntimeEvent::RunCancelled { agent_id, .. }
198            | RuntimeEvent::PlanUpdated { agent_id, .. }
199            | RuntimeEvent::UserEvent { agent_id, .. } => *agent_id = Some(id),
200        }
201        self
202    }
203}
204
205#[cfg(test)]
206mod tests {
207    use super::*;
208    use crate::types::{CheckpointStep, PlanStepStatus, RiskLevel};
209
210    fn sid(id: u64) -> SessionId {
211        SessionId::new(id)
212    }
213
214    fn approval_request() -> ApprovalRequest {
215        ApprovalRequest {
216            title: "title".to_string(),
217            message: "message".to_string(),
218            action_key: None,
219            risk_level: RiskLevel::Safe,
220            raw: None,
221        }
222    }
223
224    fn checkpoint() -> CheckpointData {
225        CheckpointData {
226            session_id: sid(42),
227            user_input: "input".to_string(),
228            step: CheckpointStep::AfterUserInput,
229            turn_count: 0,
230        }
231    }
232
233    #[test]
234    fn accessors_return_embedded_ids_for_all_variants() {
235        let events: Vec<RuntimeEvent> = vec![
236            RuntimeEvent::TextDelta {
237                session_id: sid(1),
238                text: "hi".into(),
239                agent_id: Some("a".into()),
240                trace_id: Some("t".into()),
241            },
242            RuntimeEvent::ThoughtDelta {
243                session_id: sid(2),
244                text: "hmm".into(),
245                agent_id: Some("a".into()),
246                trace_id: Some("t".into()),
247            },
248            RuntimeEvent::ToolCallStarted {
249                session_id: sid(3),
250                tool_name: "read".into(),
251                args_json: "{}".into(),
252                agent_id: Some("a".into()),
253                trace_id: Some("t".into()),
254            },
255            RuntimeEvent::ToolCallFinished {
256                session_id: sid(4),
257                tool_name: "read".into(),
258                summary: "ok".into(),
259                agent_id: Some("a".into()),
260                trace_id: Some("t".into()),
261                denied: false,
262            },
263            RuntimeEvent::AwaitingApproval {
264                session_id: sid(5),
265                request: approval_request(),
266                agent_id: Some("a".into()),
267                trace_id: Some("t".into()),
268            },
269            RuntimeEvent::Checkpoint {
270                session_id: sid(6),
271                checkpoint: checkpoint(),
272                agent_id: Some("a".into()),
273                trace_id: Some("t".into()),
274            },
275            RuntimeEvent::RunFinished {
276                session_id: sid(7),
277                agent_id: Some("a".into()),
278                trace_id: Some("t".into()),
279            },
280            RuntimeEvent::RunCancelled {
281                session_id: sid(8),
282                agent_id: Some("a".into()),
283                trace_id: Some("t".into()),
284            },
285            RuntimeEvent::PlanUpdated {
286                session_id: sid(9),
287                objective: "goal".into(),
288                explanation: None,
289                plan: vec![PlanItem {
290                    step: "s".into(),
291                    status: PlanStepStatus::Pending,
292                }],
293                agent_id: Some("a".into()),
294                trace_id: Some("t".into()),
295            },
296            RuntimeEvent::UserEvent {
297                session_id: sid(10),
298                event: UserEvent::Progress { text: "p".into() },
299                agent_id: Some("a".into()),
300                trace_id: Some("t".into()),
301            },
302        ];
303
304        for (i, ev) in events.iter().enumerate() {
305            let expected = i as u64 + 1;
306            assert_eq!(ev.session_id(), &sid(expected), "variant {i}");
307            assert_eq!(ev.agent_id(), Some("a"), "variant {i}");
308            assert_eq!(ev.trace_id(), Some("t"), "variant {i}");
309        }
310    }
311
312    #[test]
313    fn with_agent_id_sets_id_on_all_variants() {
314        let events: Vec<RuntimeEvent> = vec![
315            RuntimeEvent::TextDelta {
316                session_id: sid(1),
317                text: "hi".into(),
318                agent_id: None,
319                trace_id: None,
320            },
321            RuntimeEvent::ThoughtDelta {
322                session_id: sid(2),
323                text: "hmm".into(),
324                agent_id: None,
325                trace_id: None,
326            },
327            RuntimeEvent::ToolCallStarted {
328                session_id: sid(3),
329                tool_name: "read".into(),
330                args_json: "{}".into(),
331                agent_id: None,
332                trace_id: None,
333            },
334            RuntimeEvent::ToolCallFinished {
335                session_id: sid(4),
336                tool_name: "read".into(),
337                summary: "ok".into(),
338                agent_id: None,
339                trace_id: None,
340                denied: false,
341            },
342            RuntimeEvent::AwaitingApproval {
343                session_id: sid(5),
344                request: approval_request(),
345                agent_id: None,
346                trace_id: None,
347            },
348            RuntimeEvent::Checkpoint {
349                session_id: sid(6),
350                checkpoint: checkpoint(),
351                agent_id: None,
352                trace_id: None,
353            },
354            RuntimeEvent::RunFinished {
355                session_id: sid(7),
356                agent_id: None,
357                trace_id: None,
358            },
359            RuntimeEvent::RunCancelled {
360                session_id: sid(8),
361                agent_id: None,
362                trace_id: None,
363            },
364            RuntimeEvent::PlanUpdated {
365                session_id: sid(9),
366                objective: "goal".into(),
367                explanation: None,
368                plan: vec![PlanItem {
369                    step: "s".into(),
370                    status: PlanStepStatus::Pending,
371                }],
372                agent_id: None,
373                trace_id: None,
374            },
375            RuntimeEvent::UserEvent {
376                session_id: sid(10),
377                event: UserEvent::Progress { text: "p".into() },
378                agent_id: None,
379                trace_id: None,
380            },
381        ];
382
383        for (i, ev) in events.into_iter().enumerate() {
384            let tagged = ev.with_agent_id("sub/1");
385            assert_eq!(tagged.agent_id(), Some("sub/1"), "variant {i}");
386            assert_eq!(tagged.session_id(), &sid(i as u64 + 1), "variant {i}");
387        }
388    }
389
390    #[test]
391    fn accessors_return_none_when_ids_absent() {
392        let ev = RuntimeEvent::RunFinished {
393            session_id: sid(1),
394            agent_id: None,
395            trace_id: None,
396        };
397        assert_eq!(ev.session_id(), &sid(1));
398        assert_eq!(ev.agent_id(), None);
399        assert_eq!(ev.trace_id(), None);
400    }
401
402    #[test]
403    fn runtime_event_serde_uses_camel_case_tag() {
404        let ev = RuntimeEvent::ToolCallStarted {
405            session_id: SessionId::with_external_id(7, "ext"),
406            tool_name: "read".into(),
407            args_json: "{}".into(),
408            agent_id: Some("a".into()),
409            trace_id: None,
410        };
411        let v = serde_json::to_value(&ev).unwrap();
412        assert_eq!(v["runtimeEventType"], "toolCallStarted");
413        assert_eq!(v["session_id"]["id"], serde_json::json!(7));
414        assert_eq!(v["session_id"]["external_id"], "ext");
415        // trace_id is None → field skipped
416        assert!(v.get("trace_id").is_none());
417
418        let back: RuntimeEvent = serde_json::from_value(v).unwrap();
419        assert_eq!(back.session_id(), &SessionId::with_external_id(7, "ext"));
420        assert_eq!(back.agent_id(), Some("a"));
421        assert_eq!(back.trace_id(), None);
422    }
423
424    #[test]
425    fn user_event_serde_tag() {
426        let v = serde_json::to_value(UserEvent::Progress {
427            text: "working".into(),
428        })
429        .unwrap();
430        assert_eq!(v["userEventType"], "progress");
431        assert_eq!(v["text"], "working");
432
433        let v = serde_json::to_value(UserEvent::Structured {
434            event_type: "custom".into(),
435            data: serde_json::json!({"k": 1}),
436        })
437        .unwrap();
438        assert_eq!(v["userEventType"], "structured");
439        assert_eq!(v["data"]["k"], serde_json::json!(1));
440    }
441}