Skip to main content

agent_base/types/
events.rs

1use serde::{Deserialize, Serialize};
2use serde_json::Value;
3
4use agent_types::NoticeKind;
5
6use super::approval::ApprovalRequest;
7use super::checkpoint::CheckpointData;
8use super::plan_update::PlanItem;
9use super::session::SessionId;
10
11// ---------------------------------------------------------------------------
12// UserEvent — user-space events produced by tools
13// ---------------------------------------------------------------------------
14
15/// User-space events produced by tools during execution.
16///
17/// Tools send these through `ToolContext::emit_user_event()`. The framework
18/// wraps them in [`RuntimeEvent::UserEvent`] before delivering to external
19/// consumers.
20#[derive(Clone, Debug, Serialize, Deserialize)]
21#[serde(tag = "userEventType", rename_all = "camelCase")]
22pub enum UserEvent {
23    /// Tool progress notification.
24    Progress { text: String },
25    /// User-defined structured event for custom business semantics.
26    Structured { event_type: String, data: Value },
27    /// Tool partial result — emitted during long-running tool execution.
28    /// `is_partial: true` means more output is coming; `false` means this is the final chunk.
29    ToolPartialResult {
30        tool_call_id: String,
31        content: String,
32        is_partial: bool,
33    },
34    /// Display-only notice from an engine mechanism (guard degradation,
35    /// reaper, ...). `text` is a mechanism-level English fact; product layers
36    /// may rewrite it per `source`. `source` is a `String` (not `&'static str`)
37    /// because `UserEvent` must deserialize across process boundaries.
38    Notice {
39        kind: NoticeKind,
40        source: String,
41        text: String,
42    },
43}
44
45// ---------------------------------------------------------------------------
46// RuntimeEvent — unified event stream for all consumers
47// ---------------------------------------------------------------------------
48
49/// Unified runtime event — the single event type for both internal and external
50/// consumers (frontends, CLIs, tests).
51#[derive(Clone, Debug, Serialize, Deserialize)]
52#[serde(tag = "runtimeEventType", rename_all = "camelCase")]
53pub enum RuntimeEvent {
54    // --- Framework system events ---
55    TextDelta {
56        session_id: SessionId,
57        text: String,
58        /// Identifies the agent that produced this event (root / sub-agent path).
59        #[serde(default, skip_serializing_if = "Option::is_none")]
60        agent_id: Option<String>,
61        /// Distributed tracing context carried across systems (e.g. MCP caller → phi-agent).
62        #[serde(default, skip_serializing_if = "Option::is_none")]
63        trace_id: Option<String>,
64    },
65    ThoughtDelta {
66        session_id: SessionId,
67        text: String,
68        #[serde(default, skip_serializing_if = "Option::is_none")]
69        agent_id: Option<String>,
70        #[serde(default, skip_serializing_if = "Option::is_none")]
71        trace_id: Option<String>,
72    },
73    /// A tool call that is still being streamed into existence.
74    ///
75    /// Tool-call arguments arrive as fragments and are only materialized into
76    /// a [`RuntimeEvent::ToolCallStarted`] once the whole turn's stream is
77    /// drained — multi-KB arguments (a `spawn_agent` task brief) are invisible
78    /// to every consumer until then. This event makes that window observable:
79    /// coarse progress on the call being drafted (`index`), enough for a UI to
80    /// say "writing tool calls" instead of looking hung. It carries no
81    /// argument text, only a size, and is throttled at the source.
82    ToolCallDraft {
83        session_id: SessionId,
84        /// The call's stream index within this turn (its eventual tool-call
85        /// id is not known yet).
86        index: usize,
87        /// Function name, empty until the fragment carrying it arrives.
88        name: String,
89        /// Argument characters accumulated for this call so far.
90        args_len: usize,
91        /// How many calls this turn has produced so far (drafts included).
92        count: usize,
93        #[serde(default, skip_serializing_if = "Option::is_none")]
94        agent_id: Option<String>,
95        #[serde(default, skip_serializing_if = "Option::is_none")]
96        trace_id: Option<String>,
97    },
98    ToolCallStarted {
99        session_id: SessionId,
100        tool_name: String,
101        args_json: String,
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    ToolCallFinished {
108        session_id: SessionId,
109        tool_name: String,
110        summary: String,
111        #[serde(default, skip_serializing_if = "Option::is_none")]
112        agent_id: Option<String>,
113        #[serde(default, skip_serializing_if = "Option::is_none")]
114        trace_id: Option<String>,
115        /// `true` when this finish is the result of an approval denial (not a
116        /// normal tool result or execution error).
117        #[serde(default)]
118        denied: bool,
119        /// Structured metadata from the tool result (e.g. edit line numbers).
120        /// Carried through to UI consumers; never sent to the LLM.
121        #[serde(default, skip_serializing_if = "Option::is_none")]
122        details: Option<Value>,
123    },
124    AwaitingApproval {
125        session_id: SessionId,
126        request: ApprovalRequest,
127        #[serde(default, skip_serializing_if = "Option::is_none")]
128        agent_id: Option<String>,
129        #[serde(default, skip_serializing_if = "Option::is_none")]
130        trace_id: Option<String>,
131    },
132    Checkpoint {
133        session_id: SessionId,
134        checkpoint: CheckpointData,
135        #[serde(default, skip_serializing_if = "Option::is_none")]
136        agent_id: Option<String>,
137        #[serde(default, skip_serializing_if = "Option::is_none")]
138        trace_id: Option<String>,
139    },
140    RunFinished {
141        session_id: SessionId,
142        #[serde(default, skip_serializing_if = "Option::is_none")]
143        agent_id: Option<String>,
144        #[serde(default, skip_serializing_if = "Option::is_none")]
145        trace_id: Option<String>,
146    },
147    RunCancelled {
148        session_id: SessionId,
149        #[serde(default, skip_serializing_if = "Option::is_none")]
150        agent_id: Option<String>,
151        #[serde(default, skip_serializing_if = "Option::is_none")]
152        trace_id: Option<String>,
153    },
154    // --- Lightweight plan update (display-only, no execution semantics) ---
155    PlanUpdated {
156        session_id: SessionId,
157        objective: String,
158        explanation: Option<String>,
159        plan: Vec<PlanItem>,
160        #[serde(default, skip_serializing_if = "Option::is_none")]
161        agent_id: Option<String>,
162        #[serde(default, skip_serializing_if = "Option::is_none")]
163        trace_id: Option<String>,
164    },
165    // --- User-space events ---
166    /// A user-space event produced by a tool.
167    UserEvent {
168        session_id: SessionId,
169        event: UserEvent,
170        #[serde(default, skip_serializing_if = "Option::is_none")]
171        agent_id: Option<String>,
172        #[serde(default, skip_serializing_if = "Option::is_none")]
173        trace_id: Option<String>,
174    },
175}
176
177impl RuntimeEvent {
178    /// Get the session ID associated with this event.
179    pub fn session_id(&self) -> &SessionId {
180        match self {
181            RuntimeEvent::TextDelta { session_id, .. } => session_id,
182            RuntimeEvent::ThoughtDelta { session_id, .. } => session_id,
183            RuntimeEvent::ToolCallDraft { session_id, .. } => session_id,
184            RuntimeEvent::ToolCallStarted { session_id, .. } => session_id,
185            RuntimeEvent::ToolCallFinished { session_id, .. } => session_id,
186            RuntimeEvent::AwaitingApproval { session_id, .. } => session_id,
187            RuntimeEvent::Checkpoint { session_id, .. } => session_id,
188            RuntimeEvent::RunFinished { session_id, .. } => session_id,
189            RuntimeEvent::RunCancelled { session_id, .. } => session_id,
190            RuntimeEvent::PlanUpdated { session_id, .. } => session_id,
191            RuntimeEvent::UserEvent { session_id, .. } => session_id,
192        }
193    }
194
195    /// Get the agent ID associated with this event, if any.
196    pub fn agent_id(&self) -> Option<&str> {
197        match self {
198            RuntimeEvent::TextDelta { agent_id, .. } => agent_id.as_deref(),
199            RuntimeEvent::ThoughtDelta { agent_id, .. } => agent_id.as_deref(),
200            RuntimeEvent::ToolCallDraft { agent_id, .. } => agent_id.as_deref(),
201            RuntimeEvent::ToolCallStarted { agent_id, .. } => agent_id.as_deref(),
202            RuntimeEvent::ToolCallFinished { agent_id, .. } => agent_id.as_deref(),
203            RuntimeEvent::AwaitingApproval { agent_id, .. } => agent_id.as_deref(),
204            RuntimeEvent::Checkpoint { agent_id, .. } => agent_id.as_deref(),
205            RuntimeEvent::RunFinished { agent_id, .. } => agent_id.as_deref(),
206            RuntimeEvent::RunCancelled { agent_id, .. } => agent_id.as_deref(),
207            RuntimeEvent::PlanUpdated { agent_id, .. } => agent_id.as_deref(),
208            RuntimeEvent::UserEvent { agent_id, .. } => agent_id.as_deref(),
209        }
210    }
211
212    /// Get the trace ID associated with this event, if any.
213    pub fn trace_id(&self) -> Option<&str> {
214        match self {
215            RuntimeEvent::TextDelta { trace_id, .. } => trace_id.as_deref(),
216            RuntimeEvent::ThoughtDelta { trace_id, .. } => trace_id.as_deref(),
217            RuntimeEvent::ToolCallDraft { trace_id, .. } => trace_id.as_deref(),
218            RuntimeEvent::ToolCallStarted { trace_id, .. } => trace_id.as_deref(),
219            RuntimeEvent::ToolCallFinished { trace_id, .. } => trace_id.as_deref(),
220            RuntimeEvent::AwaitingApproval { trace_id, .. } => trace_id.as_deref(),
221            RuntimeEvent::Checkpoint { trace_id, .. } => trace_id.as_deref(),
222            RuntimeEvent::RunFinished { trace_id, .. } => trace_id.as_deref(),
223            RuntimeEvent::RunCancelled { trace_id, .. } => trace_id.as_deref(),
224            RuntimeEvent::PlanUpdated { trace_id, .. } => trace_id.as_deref(),
225            RuntimeEvent::UserEvent { trace_id, .. } => trace_id.as_deref(),
226        }
227    }
228
229    /// Set the agent ID (sub-agent path) on this event.
230    pub fn with_agent_id(mut self, id: impl Into<String>) -> Self {
231        let id = id.into();
232        match &mut self {
233            RuntimeEvent::TextDelta { agent_id, .. }
234            | RuntimeEvent::ThoughtDelta { agent_id, .. }
235            | RuntimeEvent::ToolCallDraft { agent_id, .. }
236            | RuntimeEvent::ToolCallStarted { agent_id, .. }
237            | RuntimeEvent::ToolCallFinished { agent_id, .. }
238            | RuntimeEvent::AwaitingApproval { agent_id, .. }
239            | RuntimeEvent::Checkpoint { agent_id, .. }
240            | RuntimeEvent::RunFinished { agent_id, .. }
241            | RuntimeEvent::RunCancelled { agent_id, .. }
242            | RuntimeEvent::PlanUpdated { agent_id, .. }
243            | RuntimeEvent::UserEvent { agent_id, .. } => *agent_id = Some(id),
244        }
245        self
246    }
247}
248
249#[cfg(test)]
250mod tests {
251    use super::*;
252    use crate::types::{CheckpointStep, PlanStepStatus, RiskLevel};
253
254    fn sid(id: u64) -> SessionId {
255        SessionId::new(id)
256    }
257
258    fn approval_request() -> ApprovalRequest {
259        ApprovalRequest {
260            title: "title".to_string(),
261            message: "message".to_string(),
262            action_key: None,
263            risk_level: RiskLevel::Safe,
264            raw: None,
265            source: None,
266        }
267    }
268
269    fn checkpoint() -> CheckpointData {
270        CheckpointData {
271            session_id: sid(42),
272            user_input: "input".to_string(),
273            step: CheckpointStep::AfterUserInput,
274            turn_count: 0,
275        }
276    }
277
278    #[test]
279    fn accessors_return_embedded_ids_for_all_variants() {
280        let events: Vec<RuntimeEvent> = vec![
281            RuntimeEvent::TextDelta {
282                session_id: sid(1),
283                text: "hi".into(),
284                agent_id: Some("a".into()),
285                trace_id: Some("t".into()),
286            },
287            RuntimeEvent::ThoughtDelta {
288                session_id: sid(2),
289                text: "hmm".into(),
290                agent_id: Some("a".into()),
291                trace_id: Some("t".into()),
292            },
293            RuntimeEvent::ToolCallStarted {
294                session_id: sid(3),
295                tool_name: "read".into(),
296                args_json: "{}".into(),
297                agent_id: Some("a".into()),
298                trace_id: Some("t".into()),
299            },
300            RuntimeEvent::ToolCallFinished {
301                session_id: sid(4),
302                tool_name: "read".into(),
303                summary: "ok".into(),
304                agent_id: Some("a".into()),
305                trace_id: Some("t".into()),
306                denied: false,
307                details: None,
308            },
309            RuntimeEvent::AwaitingApproval {
310                session_id: sid(5),
311                request: approval_request(),
312                agent_id: Some("a".into()),
313                trace_id: Some("t".into()),
314            },
315            RuntimeEvent::Checkpoint {
316                session_id: sid(6),
317                checkpoint: checkpoint(),
318                agent_id: Some("a".into()),
319                trace_id: Some("t".into()),
320            },
321            RuntimeEvent::RunFinished {
322                session_id: sid(7),
323                agent_id: Some("a".into()),
324                trace_id: Some("t".into()),
325            },
326            RuntimeEvent::RunCancelled {
327                session_id: sid(8),
328                agent_id: Some("a".into()),
329                trace_id: Some("t".into()),
330            },
331            RuntimeEvent::PlanUpdated {
332                session_id: sid(9),
333                objective: "goal".into(),
334                explanation: None,
335                plan: vec![PlanItem {
336                    step: "s".into(),
337                    status: PlanStepStatus::Pending,
338                }],
339                agent_id: Some("a".into()),
340                trace_id: Some("t".into()),
341            },
342            RuntimeEvent::UserEvent {
343                session_id: sid(10),
344                event: UserEvent::Progress { text: "p".into() },
345                agent_id: Some("a".into()),
346                trace_id: Some("t".into()),
347            },
348            RuntimeEvent::ToolCallDraft {
349                session_id: sid(11),
350                index: 0,
351                name: "spawn_agent".into(),
352                args_len: 512,
353                count: 1,
354                agent_id: Some("a".into()),
355                trace_id: Some("t".into()),
356            },
357        ];
358
359        for (i, ev) in events.iter().enumerate() {
360            let expected = i as u64 + 1;
361            assert_eq!(ev.session_id(), &sid(expected), "variant {i}");
362            assert_eq!(ev.agent_id(), Some("a"), "variant {i}");
363            assert_eq!(ev.trace_id(), Some("t"), "variant {i}");
364        }
365    }
366
367    #[test]
368    fn with_agent_id_sets_id_on_all_variants() {
369        let events: Vec<RuntimeEvent> = vec![
370            RuntimeEvent::TextDelta {
371                session_id: sid(1),
372                text: "hi".into(),
373                agent_id: None,
374                trace_id: None,
375            },
376            RuntimeEvent::ThoughtDelta {
377                session_id: sid(2),
378                text: "hmm".into(),
379                agent_id: None,
380                trace_id: None,
381            },
382            RuntimeEvent::ToolCallStarted {
383                session_id: sid(3),
384                tool_name: "read".into(),
385                args_json: "{}".into(),
386                agent_id: None,
387                trace_id: None,
388            },
389            RuntimeEvent::ToolCallFinished {
390                session_id: sid(4),
391                tool_name: "read".into(),
392                summary: "ok".into(),
393                agent_id: None,
394                trace_id: None,
395                denied: false,
396                details: None,
397            },
398            RuntimeEvent::AwaitingApproval {
399                session_id: sid(5),
400                request: approval_request(),
401                agent_id: None,
402                trace_id: None,
403            },
404            RuntimeEvent::Checkpoint {
405                session_id: sid(6),
406                checkpoint: checkpoint(),
407                agent_id: None,
408                trace_id: None,
409            },
410            RuntimeEvent::RunFinished {
411                session_id: sid(7),
412                agent_id: None,
413                trace_id: None,
414            },
415            RuntimeEvent::RunCancelled {
416                session_id: sid(8),
417                agent_id: None,
418                trace_id: None,
419            },
420            RuntimeEvent::PlanUpdated {
421                session_id: sid(9),
422                objective: "goal".into(),
423                explanation: None,
424                plan: vec![PlanItem {
425                    step: "s".into(),
426                    status: PlanStepStatus::Pending,
427                }],
428                agent_id: None,
429                trace_id: None,
430            },
431            RuntimeEvent::UserEvent {
432                session_id: sid(10),
433                event: UserEvent::Progress { text: "p".into() },
434                agent_id: None,
435                trace_id: None,
436            },
437            RuntimeEvent::ToolCallDraft {
438                session_id: sid(11),
439                index: 0,
440                name: "spawn_agent".into(),
441                args_len: 512,
442                count: 1,
443                agent_id: None,
444                trace_id: None,
445            },
446        ];
447
448        for (i, ev) in events.into_iter().enumerate() {
449            let tagged = ev.with_agent_id("sub/1");
450            assert_eq!(tagged.agent_id(), Some("sub/1"), "variant {i}");
451            assert_eq!(tagged.session_id(), &sid(i as u64 + 1), "variant {i}");
452        }
453    }
454
455    #[test]
456    fn accessors_return_none_when_ids_absent() {
457        let ev = RuntimeEvent::RunFinished {
458            session_id: sid(1),
459            agent_id: None,
460            trace_id: None,
461        };
462        assert_eq!(ev.session_id(), &sid(1));
463        assert_eq!(ev.agent_id(), None);
464        assert_eq!(ev.trace_id(), None);
465    }
466
467    #[test]
468    fn runtime_event_serde_uses_camel_case_tag() {
469        let ev = RuntimeEvent::ToolCallStarted {
470            session_id: SessionId::with_external_id(7, "ext"),
471            tool_name: "read".into(),
472            args_json: "{}".into(),
473            agent_id: Some("a".into()),
474            trace_id: None,
475        };
476        let v = serde_json::to_value(&ev).unwrap();
477        assert_eq!(v["runtimeEventType"], "toolCallStarted");
478        assert_eq!(v["session_id"]["id"], serde_json::json!(7));
479        assert_eq!(v["session_id"]["external_id"], "ext");
480        // trace_id is None → field skipped
481        assert!(v.get("trace_id").is_none());
482
483        let back: RuntimeEvent = serde_json::from_value(v).unwrap();
484        assert_eq!(back.session_id(), &SessionId::with_external_id(7, "ext"));
485        assert_eq!(back.agent_id(), Some("a"));
486        assert_eq!(back.trace_id(), None);
487    }
488
489    #[test]
490    fn tool_call_draft_serde_round_trip() {
491        let ev = RuntimeEvent::ToolCallDraft {
492            session_id: sid(3),
493            index: 2,
494            name: "spawn_agent".into(),
495            args_len: 640,
496            count: 4,
497            agent_id: None,
498            trace_id: None,
499        };
500        let v = serde_json::to_value(&ev).unwrap();
501        assert_eq!(v["runtimeEventType"], "toolCallDraft");
502        assert_eq!(v["args_len"], 640);
503        assert!(v.get("agent_id").is_none(), "None ids stay skipped");
504
505        let back: RuntimeEvent = serde_json::from_value(v).unwrap();
506        match back {
507            RuntimeEvent::ToolCallDraft {
508                index,
509                name,
510                args_len,
511                count,
512                ..
513            } => {
514                assert_eq!(index, 2);
515                assert_eq!(name, "spawn_agent");
516                assert_eq!(args_len, 640);
517                assert_eq!(count, 4);
518            }
519            other => panic!("expected ToolCallDraft, got {other:?}"),
520        }
521    }
522
523    #[test]
524    fn user_event_serde_tag() {
525        let v = serde_json::to_value(UserEvent::Progress {
526            text: "working".into(),
527        })
528        .unwrap();
529        assert_eq!(v["userEventType"], "progress");
530        assert_eq!(v["text"], "working");
531
532        let v = serde_json::to_value(UserEvent::Structured {
533            event_type: "custom".into(),
534            data: serde_json::json!({"k": 1}),
535        })
536        .unwrap();
537        assert_eq!(v["userEventType"], "structured");
538        assert_eq!(v["data"]["k"], serde_json::json!(1));
539    }
540
541    // Design doc §9: UserEvent::Notice serialization tag + serde round-trip
542    // (source must be String — &'static str cannot deserialize).
543    #[test]
544    fn user_event_notice_serde_round_trip() {
545        let ev = UserEvent::Notice {
546            kind: NoticeKind::Warning,
547            source: "guard".into(),
548            text: "guard judge unparsed — treating as complete".into(),
549        };
550        let v = serde_json::to_value(&ev).unwrap();
551        assert_eq!(v["userEventType"], "notice");
552        assert_eq!(v["kind"], serde_json::json!("warning"));
553        assert_eq!(v["source"], serde_json::json!("guard"));
554        assert_eq!(
555            v["text"],
556            serde_json::json!("guard judge unparsed — treating as complete")
557        );
558
559        let back: UserEvent = serde_json::from_value(v).unwrap();
560        match back {
561            UserEvent::Notice { kind, source, text } => {
562                assert_eq!(kind, NoticeKind::Warning);
563                assert_eq!(source, "guard");
564                assert_eq!(text, "guard judge unparsed — treating as complete");
565            }
566            other => panic!("expected Notice, got {other:?}"),
567        }
568    }
569
570    // Design doc §11: old JSONL replay — a Notice serialized elsewhere must
571    // deserialize cleanly alongside existing variants (serde is add-only).
572    #[test]
573    fn notice_coexists_with_legacy_variants_on_replay() {
574        let json = r#"{"userEventType":"progress","text":"legacy line"}"#;
575        let back: UserEvent = serde_json::from_str(json).unwrap();
576        assert!(matches!(back, UserEvent::Progress { .. }));
577    }
578}