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#[derive(Clone, Debug, Serialize, Deserialize)]
19#[serde(tag = "userEventType", rename_all = "camelCase")]
20pub enum UserEvent {
21 Progress { text: String },
23 Structured { event_type: String, data: Value },
25 ToolPartialResult {
28 tool_call_id: String,
29 content: String,
30 is_partial: bool,
31 },
32}
33
34#[derive(Clone, Debug, Serialize, Deserialize)]
41#[serde(tag = "runtimeEventType", rename_all = "camelCase")]
42pub enum RuntimeEvent {
43 TextDelta {
45 session_id: SessionId,
46 text: String,
47 #[serde(default, skip_serializing_if = "Option::is_none")]
49 agent_id: Option<String>,
50 #[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 #[serde(default)]
82 denied: bool,
83 #[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 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 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 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 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 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 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 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}