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#[derive(Clone, Debug, Serialize, Deserialize)]
21#[serde(tag = "userEventType", rename_all = "camelCase")]
22pub enum UserEvent {
23 Progress { text: String },
25 Structured { event_type: String, data: Value },
27 ToolPartialResult {
30 tool_call_id: String,
31 content: String,
32 is_partial: bool,
33 },
34 Notice {
39 kind: NoticeKind,
40 source: String,
41 text: String,
42 },
43}
44
45#[derive(Clone, Debug, Serialize, Deserialize)]
52#[serde(tag = "runtimeEventType", rename_all = "camelCase")]
53pub enum RuntimeEvent {
54 TextDelta {
56 session_id: SessionId,
57 text: String,
58 #[serde(default, skip_serializing_if = "Option::is_none")]
60 agent_id: Option<String>,
61 #[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 ToolCallDraft {
83 session_id: SessionId,
84 index: usize,
87 name: String,
89 args_len: usize,
91 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 #[serde(default)]
118 denied: bool,
119 #[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 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 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 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 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 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 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 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 #[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 #[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}