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        /// Structured metadata from the tool result (e.g. edit line numbers).
84        /// Carried through to UI consumers; never sent to the LLM.
85        #[serde(default, skip_serializing_if = "Option::is_none")]
86        details: Option<Value>,
87    },
88    AwaitingApproval {
89        session_id: SessionId,
90        request: ApprovalRequest,
91        #[serde(default, skip_serializing_if = "Option::is_none")]
92        agent_id: Option<String>,
93        #[serde(default, skip_serializing_if = "Option::is_none")]
94        trace_id: Option<String>,
95    },
96    Checkpoint {
97        session_id: SessionId,
98        checkpoint: CheckpointData,
99        #[serde(default, skip_serializing_if = "Option::is_none")]
100        agent_id: Option<String>,
101        #[serde(default, skip_serializing_if = "Option::is_none")]
102        trace_id: Option<String>,
103    },
104    RunFinished {
105        session_id: SessionId,
106        #[serde(default, skip_serializing_if = "Option::is_none")]
107        agent_id: Option<String>,
108        #[serde(default, skip_serializing_if = "Option::is_none")]
109        trace_id: Option<String>,
110    },
111    RunCancelled {
112        session_id: SessionId,
113        #[serde(default, skip_serializing_if = "Option::is_none")]
114        agent_id: Option<String>,
115        #[serde(default, skip_serializing_if = "Option::is_none")]
116        trace_id: Option<String>,
117    },
118    // --- Lightweight plan update (display-only, no execution semantics) ---
119    PlanUpdated {
120        session_id: SessionId,
121        objective: String,
122        explanation: Option<String>,
123        plan: Vec<PlanItem>,
124        #[serde(default, skip_serializing_if = "Option::is_none")]
125        agent_id: Option<String>,
126        #[serde(default, skip_serializing_if = "Option::is_none")]
127        trace_id: Option<String>,
128    },
129    // --- User-space events ---
130    /// A user-space event produced by a tool.
131    UserEvent {
132        session_id: SessionId,
133        event: UserEvent,
134        #[serde(default, skip_serializing_if = "Option::is_none")]
135        agent_id: Option<String>,
136        #[serde(default, skip_serializing_if = "Option::is_none")]
137        trace_id: Option<String>,
138    },
139}
140
141impl RuntimeEvent {
142    /// Get the session ID associated with this event.
143    pub fn session_id(&self) -> &SessionId {
144        match self {
145            RuntimeEvent::TextDelta { session_id, .. } => session_id,
146            RuntimeEvent::ThoughtDelta { session_id, .. } => session_id,
147            RuntimeEvent::ToolCallStarted { session_id, .. } => session_id,
148            RuntimeEvent::ToolCallFinished { session_id, .. } => session_id,
149            RuntimeEvent::AwaitingApproval { session_id, .. } => session_id,
150            RuntimeEvent::Checkpoint { session_id, .. } => session_id,
151            RuntimeEvent::RunFinished { session_id, .. } => session_id,
152            RuntimeEvent::RunCancelled { session_id, .. } => session_id,
153            RuntimeEvent::PlanUpdated { session_id, .. } => session_id,
154            RuntimeEvent::UserEvent { session_id, .. } => session_id,
155        }
156    }
157
158    /// Get the agent ID associated with this event, if any.
159    pub fn agent_id(&self) -> Option<&str> {
160        match self {
161            RuntimeEvent::TextDelta { agent_id, .. } => agent_id.as_deref(),
162            RuntimeEvent::ThoughtDelta { agent_id, .. } => agent_id.as_deref(),
163            RuntimeEvent::ToolCallStarted { agent_id, .. } => agent_id.as_deref(),
164            RuntimeEvent::ToolCallFinished { agent_id, .. } => agent_id.as_deref(),
165            RuntimeEvent::AwaitingApproval { agent_id, .. } => agent_id.as_deref(),
166            RuntimeEvent::Checkpoint { agent_id, .. } => agent_id.as_deref(),
167            RuntimeEvent::RunFinished { agent_id, .. } => agent_id.as_deref(),
168            RuntimeEvent::RunCancelled { agent_id, .. } => agent_id.as_deref(),
169            RuntimeEvent::PlanUpdated { agent_id, .. } => agent_id.as_deref(),
170            RuntimeEvent::UserEvent { agent_id, .. } => agent_id.as_deref(),
171        }
172    }
173
174    /// Get the trace ID associated with this event, if any.
175    pub fn trace_id(&self) -> Option<&str> {
176        match self {
177            RuntimeEvent::TextDelta { trace_id, .. } => trace_id.as_deref(),
178            RuntimeEvent::ThoughtDelta { trace_id, .. } => trace_id.as_deref(),
179            RuntimeEvent::ToolCallStarted { trace_id, .. } => trace_id.as_deref(),
180            RuntimeEvent::ToolCallFinished { trace_id, .. } => trace_id.as_deref(),
181            RuntimeEvent::AwaitingApproval { trace_id, .. } => trace_id.as_deref(),
182            RuntimeEvent::Checkpoint { trace_id, .. } => trace_id.as_deref(),
183            RuntimeEvent::RunFinished { trace_id, .. } => trace_id.as_deref(),
184            RuntimeEvent::RunCancelled { trace_id, .. } => trace_id.as_deref(),
185            RuntimeEvent::PlanUpdated { trace_id, .. } => trace_id.as_deref(),
186            RuntimeEvent::UserEvent { trace_id, .. } => trace_id.as_deref(),
187        }
188    }
189
190    /// Set the agent ID (sub-agent path) on this event.
191    pub fn with_agent_id(mut self, id: impl Into<String>) -> Self {
192        let id = id.into();
193        match &mut self {
194            RuntimeEvent::TextDelta { agent_id, .. }
195            | RuntimeEvent::ThoughtDelta { agent_id, .. }
196            | RuntimeEvent::ToolCallStarted { agent_id, .. }
197            | RuntimeEvent::ToolCallFinished { agent_id, .. }
198            | RuntimeEvent::AwaitingApproval { agent_id, .. }
199            | RuntimeEvent::Checkpoint { agent_id, .. }
200            | RuntimeEvent::RunFinished { agent_id, .. }
201            | RuntimeEvent::RunCancelled { agent_id, .. }
202            | RuntimeEvent::PlanUpdated { agent_id, .. }
203            | RuntimeEvent::UserEvent { agent_id, .. } => *agent_id = Some(id),
204        }
205        self
206    }
207}
208
209#[cfg(test)]
210mod tests {
211    use super::*;
212    use crate::types::{CheckpointStep, PlanStepStatus, RiskLevel};
213
214    fn sid(id: u64) -> SessionId {
215        SessionId::new(id)
216    }
217
218    fn approval_request() -> ApprovalRequest {
219        ApprovalRequest {
220            title: "title".to_string(),
221            message: "message".to_string(),
222            action_key: None,
223            risk_level: RiskLevel::Safe,
224            raw: None,
225        }
226    }
227
228    fn checkpoint() -> CheckpointData {
229        CheckpointData {
230            session_id: sid(42),
231            user_input: "input".to_string(),
232            step: CheckpointStep::AfterUserInput,
233            turn_count: 0,
234        }
235    }
236
237    #[test]
238    fn accessors_return_embedded_ids_for_all_variants() {
239        let events: Vec<RuntimeEvent> = vec![
240            RuntimeEvent::TextDelta {
241                session_id: sid(1),
242                text: "hi".into(),
243                agent_id: Some("a".into()),
244                trace_id: Some("t".into()),
245            },
246            RuntimeEvent::ThoughtDelta {
247                session_id: sid(2),
248                text: "hmm".into(),
249                agent_id: Some("a".into()),
250                trace_id: Some("t".into()),
251            },
252            RuntimeEvent::ToolCallStarted {
253                session_id: sid(3),
254                tool_name: "read".into(),
255                args_json: "{}".into(),
256                agent_id: Some("a".into()),
257                trace_id: Some("t".into()),
258            },
259            RuntimeEvent::ToolCallFinished {
260                session_id: sid(4),
261                tool_name: "read".into(),
262                summary: "ok".into(),
263                agent_id: Some("a".into()),
264                trace_id: Some("t".into()),
265                denied: false,
266                details: None,
267            },
268            RuntimeEvent::AwaitingApproval {
269                session_id: sid(5),
270                request: approval_request(),
271                agent_id: Some("a".into()),
272                trace_id: Some("t".into()),
273            },
274            RuntimeEvent::Checkpoint {
275                session_id: sid(6),
276                checkpoint: checkpoint(),
277                agent_id: Some("a".into()),
278                trace_id: Some("t".into()),
279            },
280            RuntimeEvent::RunFinished {
281                session_id: sid(7),
282                agent_id: Some("a".into()),
283                trace_id: Some("t".into()),
284            },
285            RuntimeEvent::RunCancelled {
286                session_id: sid(8),
287                agent_id: Some("a".into()),
288                trace_id: Some("t".into()),
289            },
290            RuntimeEvent::PlanUpdated {
291                session_id: sid(9),
292                objective: "goal".into(),
293                explanation: None,
294                plan: vec![PlanItem {
295                    step: "s".into(),
296                    status: PlanStepStatus::Pending,
297                }],
298                agent_id: Some("a".into()),
299                trace_id: Some("t".into()),
300            },
301            RuntimeEvent::UserEvent {
302                session_id: sid(10),
303                event: UserEvent::Progress { text: "p".into() },
304                agent_id: Some("a".into()),
305                trace_id: Some("t".into()),
306            },
307        ];
308
309        for (i, ev) in events.iter().enumerate() {
310            let expected = i as u64 + 1;
311            assert_eq!(ev.session_id(), &sid(expected), "variant {i}");
312            assert_eq!(ev.agent_id(), Some("a"), "variant {i}");
313            assert_eq!(ev.trace_id(), Some("t"), "variant {i}");
314        }
315    }
316
317    #[test]
318    fn with_agent_id_sets_id_on_all_variants() {
319        let events: Vec<RuntimeEvent> = vec![
320            RuntimeEvent::TextDelta {
321                session_id: sid(1),
322                text: "hi".into(),
323                agent_id: None,
324                trace_id: None,
325            },
326            RuntimeEvent::ThoughtDelta {
327                session_id: sid(2),
328                text: "hmm".into(),
329                agent_id: None,
330                trace_id: None,
331            },
332            RuntimeEvent::ToolCallStarted {
333                session_id: sid(3),
334                tool_name: "read".into(),
335                args_json: "{}".into(),
336                agent_id: None,
337                trace_id: None,
338            },
339            RuntimeEvent::ToolCallFinished {
340                session_id: sid(4),
341                tool_name: "read".into(),
342                summary: "ok".into(),
343                agent_id: None,
344                trace_id: None,
345                denied: false,
346                details: None,
347            },
348            RuntimeEvent::AwaitingApproval {
349                session_id: sid(5),
350                request: approval_request(),
351                agent_id: None,
352                trace_id: None,
353            },
354            RuntimeEvent::Checkpoint {
355                session_id: sid(6),
356                checkpoint: checkpoint(),
357                agent_id: None,
358                trace_id: None,
359            },
360            RuntimeEvent::RunFinished {
361                session_id: sid(7),
362                agent_id: None,
363                trace_id: None,
364            },
365            RuntimeEvent::RunCancelled {
366                session_id: sid(8),
367                agent_id: None,
368                trace_id: None,
369            },
370            RuntimeEvent::PlanUpdated {
371                session_id: sid(9),
372                objective: "goal".into(),
373                explanation: None,
374                plan: vec![PlanItem {
375                    step: "s".into(),
376                    status: PlanStepStatus::Pending,
377                }],
378                agent_id: None,
379                trace_id: None,
380            },
381            RuntimeEvent::UserEvent {
382                session_id: sid(10),
383                event: UserEvent::Progress { text: "p".into() },
384                agent_id: None,
385                trace_id: None,
386            },
387        ];
388
389        for (i, ev) in events.into_iter().enumerate() {
390            let tagged = ev.with_agent_id("sub/1");
391            assert_eq!(tagged.agent_id(), Some("sub/1"), "variant {i}");
392            assert_eq!(tagged.session_id(), &sid(i as u64 + 1), "variant {i}");
393        }
394    }
395
396    #[test]
397    fn accessors_return_none_when_ids_absent() {
398        let ev = RuntimeEvent::RunFinished {
399            session_id: sid(1),
400            agent_id: None,
401            trace_id: None,
402        };
403        assert_eq!(ev.session_id(), &sid(1));
404        assert_eq!(ev.agent_id(), None);
405        assert_eq!(ev.trace_id(), None);
406    }
407
408    #[test]
409    fn runtime_event_serde_uses_camel_case_tag() {
410        let ev = RuntimeEvent::ToolCallStarted {
411            session_id: SessionId::with_external_id(7, "ext"),
412            tool_name: "read".into(),
413            args_json: "{}".into(),
414            agent_id: Some("a".into()),
415            trace_id: None,
416        };
417        let v = serde_json::to_value(&ev).unwrap();
418        assert_eq!(v["runtimeEventType"], "toolCallStarted");
419        assert_eq!(v["session_id"]["id"], serde_json::json!(7));
420        assert_eq!(v["session_id"]["external_id"], "ext");
421        // trace_id is None → field skipped
422        assert!(v.get("trace_id").is_none());
423
424        let back: RuntimeEvent = serde_json::from_value(v).unwrap();
425        assert_eq!(back.session_id(), &SessionId::with_external_id(7, "ext"));
426        assert_eq!(back.agent_id(), Some("a"));
427        assert_eq!(back.trace_id(), None);
428    }
429
430    #[test]
431    fn user_event_serde_tag() {
432        let v = serde_json::to_value(UserEvent::Progress {
433            text: "working".into(),
434        })
435        .unwrap();
436        assert_eq!(v["userEventType"], "progress");
437        assert_eq!(v["text"], "working");
438
439        let v = serde_json::to_value(UserEvent::Structured {
440            event_type: "custom".into(),
441            data: serde_json::json!({"k": 1}),
442        })
443        .unwrap();
444        assert_eq!(v["userEventType"], "structured");
445        assert_eq!(v["data"]["k"], serde_json::json!(1));
446    }
447}