1use std::sync::{Arc, Mutex};
2
3use serde::{Deserialize, Serialize};
4use tokio::sync::broadcast;
5use tokio::sync::mpsc;
6use tokio_util::sync::CancellationToken;
7use uuid::Uuid;
8
9#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
10#[serde(transparent)]
11pub struct FlowRunId(pub Uuid);
12
13impl FlowRunId {
14 pub fn now() -> Self {
15 Self(Uuid::now_v7())
16 }
17}
18
19impl std::fmt::Display for FlowRunId {
20 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
21 self.0.fmt(f)
22 }
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
26#[serde(transparent)]
27pub struct TurnId(pub Uuid);
28
29impl TurnId {
30 pub fn now() -> Self {
31 Self(Uuid::now_v7())
32 }
33}
34
35impl std::fmt::Display for TurnId {
36 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
37 self.0.fmt(f)
38 }
39}
40
41#[derive(Debug, Clone)]
44pub struct EventEnvelope {
45 pub seq: u64,
46 pub ts: chrono::DateTime<chrono::Utc>,
47 pub event: Event,
48}
49
50impl EventEnvelope {
51 pub fn new(seq: u64, event: Event) -> Self {
52 let ts = chrono::Utc::now();
53 Self { seq, ts, event }
54 }
55
56 pub(crate) fn from_json_value(value: serde_json::Value) -> serde_json::Result<Self> {
57 let seq = value.get("seq").and_then(|v| v.as_u64()).unwrap_or(0);
58 let ts = value
59 .get("ts")
60 .and_then(|v| v.as_str())
61 .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
62 .map(|dt| dt.with_timezone(&chrono::Utc))
63 .unwrap_or_else(chrono::Utc::now);
64 let event = serde_json::from_value(value)?;
65 Ok(Self { seq, ts, event })
66 }
67}
68
69impl serde::Serialize for EventEnvelope {
70 fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
71 let mut value = serde_json::to_value(&self.event).map_err(serde::ser::Error::custom)?;
72 if let serde_json::Value::Object(ref mut map) = value {
73 map.insert("seq".into(), serde_json::Value::Number(self.seq.into()));
74 map.insert("ts".into(), serde_json::Value::String(self.ts.to_rfc3339()));
75 }
76 value.serialize(serializer)
77 }
78}
79
80impl<'de> serde::Deserialize<'de> for EventEnvelope {
81 fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
82 let value = serde_json::Value::deserialize(deserializer)?;
83 Self::from_json_value(value).map_err(serde::de::Error::custom)
84 }
85}
86
87#[derive(Debug, Clone, Serialize, Deserialize)]
88#[serde(tag = "type", rename_all = "snake_case")]
89pub enum Event {
90 FlowStart {
91 run_id: FlowRunId,
92 #[serde(default)]
93 flow_name: String,
94 #[serde(default)]
95 parent_run_id: Option<FlowRunId>,
96 #[serde(default)]
97 parent_node_id: Option<String>,
98 #[serde(default)]
99 spawned: bool,
100 },
101 FlowEnd {
102 run_id: FlowRunId,
103 flow_name: String,
104 status: FlowStatus,
105 },
106 WorkspaceLifecycle {
107 run_id: FlowRunId,
108 workspace_id: String,
109 path: String,
110 state: String,
111 #[serde(default, skip_serializing_if = "Option::is_none")]
112 cleanup_error: Option<String>,
113 },
114 LlmCall {
115 model: String,
116 provider: String,
117 #[serde(default, skip_serializing_if = "Option::is_none")]
118 context_plan_id: Option<crate::context_plan::ContextPlanId>,
119 #[serde(default, skip_serializing_if = "Option::is_none")]
120 context_epoch: Option<crate::context_plan::ContextEpoch>,
121 #[serde(default, skip_serializing_if = "Option::is_none")]
122 context_tokens: Option<crate::context_plan::ContextTokenLanes>,
123 #[serde(default, skip_serializing_if = "Option::is_none")]
124 usage_source: Option<crate::context_plan::TokenUsageSource>,
125 #[serde(default, skip_serializing_if = "Option::is_none")]
126 context_call_purpose: Option<crate::context_plan::ContextCallPurpose>,
127 #[serde(default, skip_serializing_if = "Option::is_none")]
128 context_call_identity: Option<crate::context_plan::ContextCallIdentity>,
129 #[serde(default, skip_serializing_if = "Option::is_none")]
130 context_cache: Option<crate::context_plan::ContextCacheObservation>,
131 #[serde(default, skip_serializing_if = "Option::is_none")]
132 assistant_tool_batch_width: Option<u64>,
133 #[serde(default)]
134 usage: crate::provider::TokenUsage,
135 #[serde(default)]
136 wallclock_ms: u64,
137 #[serde(default)]
138 ttft_ms: Option<u64>,
139 #[serde(default)]
140 tokens_per_second: Option<f64>,
141 #[serde(default)]
142 status: LlmCallStatus,
143 #[serde(default, skip_serializing_if = "Option::is_none")]
144 run_id: Option<crate::event::FlowRunId>,
145 #[serde(default, skip_serializing_if = "Option::is_none")]
146 node_id: Option<String>,
147 },
148 TurnStart {
149 turn_id: TurnId,
150 },
151 TurnEnd {
152 turn_id: TurnId,
153 },
154 UserMsg {
155 turn_id: TurnId,
156 #[serde(default)]
157 flow_run_id: Option<FlowRunId>,
158 message: crate::message::Message,
159 },
160 AssistantMsg {
161 turn_id: TurnId,
162 #[serde(default)]
163 flow_run_id: Option<FlowRunId>,
164 message: crate::message::Message,
165 },
166 ToolResultMsg {
167 turn_id: TurnId,
168 #[serde(default)]
169 flow_run_id: Option<FlowRunId>,
170 message: crate::message::Message,
171 },
172 ToolResultMetrics {
173 turn_id: TurnId,
174 #[serde(default, skip_serializing_if = "Option::is_none")]
175 flow_run_id: Option<FlowRunId>,
176 tool_use_id: String,
177 raw_bytes: u64,
178 excerpt_bytes: u64,
179 truncated: bool,
180 },
181 DiffPreview {
182 #[serde(default)]
183 turn_id: Option<TurnId>,
184 #[serde(default)]
185 flow_run_id: Option<FlowRunId>,
186 #[serde(default, skip_serializing_if = "Option::is_none")]
187 tool_use_id: Option<String>,
188 title: String,
189 #[serde(default)]
190 old_content: Option<String>,
191 #[serde(default)]
192 new_content: Option<String>,
193 #[serde(default)]
194 unified_diff: Option<String>,
195 },
196 FileEditApplied {
197 #[serde(default, skip_serializing_if = "Option::is_none")]
198 turn_id: Option<TurnId>,
199 #[serde(default, skip_serializing_if = "Option::is_none")]
200 flow_run_id: Option<FlowRunId>,
201 #[serde(default, skip_serializing_if = "Option::is_none")]
202 tool_use_id: Option<String>,
203 tool_name: String,
204 path: String,
205 metrics: crate::activity::EditMetrics,
206 },
207 CompactionSummary {
208 session_id: String,
209 #[serde(default, skip_serializing_if = "Option::is_none")]
210 flow_run_id: Option<FlowRunId>,
211 range_start: u64,
212 range_end: u64,
213 compacted_count: usize,
214 before_tokens: u64,
215 after_tokens: u64,
216 summary: String,
217 },
218 SystemMsg {
219 turn_id: TurnId,
220 #[serde(default, skip_serializing_if = "Option::is_none")]
221 flow_run_id: Option<FlowRunId>,
222 message: crate::message::Message,
223 },
224 UserInject {
225 turn_id: TurnId,
226 injection: crate::injection::Injection,
227 },
228 ContentFilterHit {
229 turn_id: Option<TurnId>,
230 flow_run_id: Option<FlowRunId>,
231 provider: String,
232 model: String,
233 category: String,
234 action: String,
235 },
236 ContextCompact {
237 session_id: String,
238 #[serde(default, skip_serializing_if = "Option::is_none")]
239 flow_run_id: Option<FlowRunId>,
240 before_tokens: u64,
241 after_tokens: u64,
242 compacted_range_start: u64,
243 compacted_range_end: u64,
244 #[serde(default, skip_serializing_if = "Option::is_none")]
245 summary_text: Option<String>,
246 #[serde(default, skip_serializing_if = "Option::is_none")]
247 replacement_msg_seq: Option<u64>,
248 },
249 Checkpoint {
250 session_id: String,
251 #[serde(default, skip_serializing_if = "Option::is_none")]
252 flow_run_id: Option<FlowRunId>,
253 messages: Vec<crate::message::Message>,
254 window_tokens: u64,
255 },
256 ContextTruncated {
257 turn_id: Option<TurnId>,
258 flow_run_id: Option<FlowRunId>,
259 original_chars: u64,
260 result_chars: u64,
261 dropped_chars: u64,
262 budget_tokens: u64,
263 },
264 WatchWarn {
265 turn_id: Option<TurnId>,
266 flow_run_id: Option<FlowRunId>,
267 target: String,
268 trigger: String,
269 message: String,
270 },
271 LlmPartialCall {
272 turn_id: Option<TurnId>,
273 flow_run_id: Option<FlowRunId>,
274 model: String,
275 provider: String,
276 tokens_before_abort: u64,
277 restart_reason: String,
278 },
279 PendingPrompt {
280 prompt_id: uuid::Uuid,
281 kind: String,
282 payload: serde_json::Value,
283 },
284 PromptExpired {
285 prompt_id: uuid::Uuid,
286 },
287 PromptResolved {
288 prompt_id: uuid::Uuid,
289 answer: serde_json::Value,
290 },
291 DeferredFormRecorded {
292 answer: crate::form::DeferredFormAnswer,
293 },
294 DeferredFormApplied {
295 prompt_id: String,
296 #[serde(default)]
297 flow_run_id: Option<FlowRunId>,
298 message: crate::message::Message,
299 },
300 FlowGraph {
301 run_id: FlowRunId,
302 graph: crate::nodegraph::FlowGraph,
303 },
304 FlowNodeStart {
305 run_id: FlowRunId,
306 node_id: String,
307 #[serde(default = "default_replay_node_kind")]
308 kind: crate::nodegraph::NodeKind,
309 #[serde(default)]
310 label: String,
311 #[serde(default)]
312 parent_node_id: Option<String>,
313 },
314 FlowNodeEnd {
315 run_id: FlowRunId,
316 node_id: String,
317 #[serde(default)]
318 status: FlowNodeStatus,
319 #[serde(default)]
320 output_preview: Option<String>,
321 },
322 ToolNode {
323 run_id: FlowRunId,
324 parent_node_id: String,
325 tool_use_id: String,
326 tool_name: String,
327 args_preview: String,
328 #[serde(default, skip_serializing_if = "Option::is_none")]
329 call_intent: Option<crate::message::ToolCallIntent>,
330 },
331 AttachmentDegraded {
332 turn_id: Option<TurnId>,
333 flow_run_id: Option<FlowRunId>,
334 message_seq: u64,
335 part_index: usize,
336 file_basename: String,
337 reason: String,
338 },
339 ToolPendingApproval {
340 run_id: FlowRunId,
341 tool_use_id: String,
342 tool_name: String,
343 args_preview: String,
344 level: String,
345 #[serde(default, skip_serializing_if = "Option::is_none")]
346 preview: Option<String>,
347 },
348 ToolApproved {
349 run_id: FlowRunId,
350 tool_use_id: String,
351 decided_by: String,
352 },
353 ToolDenied {
354 run_id: FlowRunId,
355 tool_use_id: String,
356 reason: String,
357 },
358 PermissionRequestCreated {
359 payload: crate::permission_audit::PermissionRequestAudit,
360 },
361 PermissionRequestTargeted {
362 payload: crate::permission_audit::PermissionRequestAudit,
363 },
364 PermissionRequestDeferred {
365 payload: crate::permission_audit::PermissionRequestAudit,
366 },
367 PermissionRequestApproved {
368 payload: crate::permission_audit::PermissionRequestAudit,
369 },
370 PermissionRequestDenied {
371 payload: crate::permission_audit::PermissionRequestAudit,
372 },
373 PermissionRequestCancelled {
374 payload: crate::permission_audit::PermissionRequestAudit,
375 },
376 PermissionGroupCreated {
377 payload: crate::permission_audit::PermissionGroupAudit,
378 },
379 PermissionGroupUpdated {
380 payload: crate::permission_audit::PermissionGroupAudit,
381 },
382 PermissionGroupResolved {
383 payload: crate::permission_audit::PermissionGroupAudit,
384 },
385 PermissionGrantCreated {
386 payload: crate::permission_audit::PermissionGrantAudit,
387 },
388 PermissionGrantExpired {
389 payload: crate::permission_audit::PermissionGrantAudit,
390 },
391 UnrestrictedExecution {
392 payload: crate::permission_audit::PermissionRequestAudit,
393 },
394 TerminalFinalState {
398 handle: String,
399 screen: crate::tools::term::TerminalScreen,
400 state: crate::tools::term::TermStateSnapshot,
401 },
402 MermaidDiagram {
403 source: String,
404 },
405}
406
407#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
408#[serde(rename_all = "snake_case")]
409pub enum FlowNodeStatus {
410 #[default]
411 Ok,
412 Err,
413 Cancelled,
414}
415
416#[derive(Debug, Clone, Serialize, Deserialize)]
417#[serde(tag = "kind", rename_all = "snake_case")]
418pub enum FlowStatus {
419 Ok,
420 Errored { message: String },
421 Cancelled,
422}
423
424impl FlowStatus {
425 pub fn errored(msg: impl Into<String>) -> Self {
426 Self::Errored {
427 message: msg.into(),
428 }
429 }
430}
431
432#[derive(Debug, Clone, Default, Serialize, Deserialize)]
433#[serde(tag = "kind", rename_all = "snake_case")]
434pub enum LlmCallStatus {
435 #[default]
436 Ok,
437 Errored {
438 message: String,
439 },
440}
441
442fn default_replay_node_kind() -> crate::nodegraph::NodeKind {
443 crate::nodegraph::NodeKind::UserConfirm
444}
445
446impl LlmCallStatus {
447 pub fn errored(msg: impl Into<String>) -> Self {
448 Self::Errored {
449 message: msg.into(),
450 }
451 }
452}
453
454#[derive(Debug, Clone)]
455pub enum NodeEvent {
456 LlmChunk {
457 text: String,
458 cumulative_tokens: u64,
459 },
460 ThinkingChunk {
461 text: String,
462 },
463 ToolCallDraft {
464 index: usize,
465 call_id: String,
466 name: String,
467 arguments_delta: String,
468 },
469 LlmDone {
470 total_tokens: u64,
471 },
472 ToolStdoutLine {
473 line: String,
474 },
475 ToolStderrLine {
476 line: String,
477 },
478 ToolDone {
479 exit: i32,
480 },
481}
482
483pub struct Observable<T> {
484 pub output: crate::tool::BoxFut<'static, Result<T, crate::error::RuntimeError>>,
485 pub events: broadcast::Receiver<NodeEvent>,
486 pub cancel: CancellationToken,
487}
488
489#[derive(Default, Clone)]
490pub struct EventSink {
491 events: Arc<Mutex<Vec<EventEnvelope>>>,
492 forwarder: Option<mpsc::UnboundedSender<EventEnvelope>>,
493 seq_counter: Arc<std::sync::atomic::AtomicU64>,
494 redactor: Option<Arc<crate::redact::Redactor>>,
495 last_compact_at: Arc<Mutex<Option<chrono::DateTime<chrono::Utc>>>>,
496}
497
498impl EventSink {
499 pub fn new() -> Self {
500 Self::default()
501 }
502
503 pub fn with_forwarder(mut self, tx: mpsc::UnboundedSender<EventEnvelope>) -> Self {
504 self.forwarder = Some(tx);
505 self
506 }
507
508 pub fn with_redactor(mut self, redactor: Arc<crate::redact::Redactor>) -> Self {
509 self.redactor = Some(redactor);
510 self
511 }
512
513 pub fn next_seq_peek(&self) -> u64 {
517 self.seq_counter.load(std::sync::atomic::Ordering::SeqCst) + 1
518 }
519
520 pub fn restore_seq(&self, last_seq: u64) {
521 self.seq_counter
522 .store(last_seq, std::sync::atomic::Ordering::SeqCst);
523 }
524
525 pub fn reserve_seq(&self) -> u64 {
529 self.seq_counter
530 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
531 + 1
532 }
533
534 pub fn emit_returning_seq(&self, event: Event) -> u64 {
535 let next = self
536 .seq_counter
537 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
538 + 1;
539 let envelope = EventEnvelope::new(next, event);
540 if let Some(tx) = &self.forwarder {
541 let _ = tx.send(envelope.clone());
542 }
543 self.events
544 .lock()
545 .expect("event sink poisoned")
546 .push(envelope);
547 next
548 }
549
550 pub fn emit(&self, event: Event) {
551 let next = self
552 .seq_counter
553 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
554 + 1;
555 let envelope = EventEnvelope::new(next, event);
556 if let Some(tx) = &self.forwarder {
557 let _ = tx.send(envelope.clone());
558 }
559 self.events
560 .lock()
561 .expect("event sink poisoned")
562 .push(envelope);
563 }
564
565 pub fn events_handle(&self) -> Arc<Mutex<Vec<EventEnvelope>>> {
566 self.events.clone()
567 }
568
569 pub fn redactor(&self) -> Option<Arc<crate::redact::Redactor>> {
570 self.redactor.clone()
571 }
572
573 pub fn mark_compacted(&self) {
574 *self.last_compact_at.lock().expect("last_compact poisoned") = Some(chrono::Utc::now());
575 }
576
577 pub fn last_compact_ago_seconds(&self) -> Option<i64> {
578 self.last_compact_at
579 .lock()
580 .expect("last_compact poisoned")
581 .map(|t| (chrono::Utc::now() - t).num_seconds())
582 }
583
584 pub fn drain(&self) -> Vec<Event> {
585 std::mem::take(&mut *self.events.lock().expect("event sink poisoned"))
586 .into_iter()
587 .map(|envelope| envelope.event)
588 .collect()
589 }
590
591 pub fn snapshot(&self) -> Vec<Event> {
592 self.events
593 .lock()
594 .expect("event sink poisoned")
595 .iter()
596 .map(|envelope| envelope.event.clone())
597 .collect()
598 }
599
600 pub fn snapshot_envelopes(&self) -> Vec<EventEnvelope> {
601 self.events
602 .lock()
603 .expect("event sink poisoned")
604 .iter()
605 .cloned()
606 .collect()
607 }
608}
609
610#[cfg(test)]
611mod tests {
612 use super::*;
613
614 #[test]
615 fn flow_start_serializes_parent_linkage() {
616 let parent = FlowRunId::now();
617 let ev = Event::FlowStart {
618 run_id: FlowRunId::now(),
619 flow_name: "child".into(),
620 parent_run_id: Some(parent.clone()),
621 parent_node_id: Some("stmt_3".into()),
622 spawned: false,
623 };
624 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
625 assert_eq!(v["type"], "flow_start");
626 assert_eq!(v["parent_run_id"], serde_json::json!(parent.0.to_string()));
627 assert_eq!(v["parent_node_id"], "stmt_3");
628 }
629
630 #[test]
631 fn flow_node_start_carries_parent_node_id() {
632 let ev = Event::FlowNodeStart {
633 run_id: FlowRunId::now(),
634 node_id: "stmt_1.branch[0]".into(),
635 kind: crate::nodegraph::NodeKind::UserConfirm,
636 label: "fanout".into(),
637 parent_node_id: Some("stmt_1".into()),
638 };
639 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
640 assert_eq!(v["type"], "flow_node_start");
641 assert_eq!(v["parent_node_id"], "stmt_1");
642 }
643
644 #[test]
645 fn tool_node_serializes_all_fields() {
646 let run_id = FlowRunId::now();
647 let ev = Event::ToolNode {
648 run_id: run_id.clone(),
649 parent_node_id: "stmt_2".into(),
650 tool_use_id: "tu_abc".into(),
651 tool_name: "fs.read".into(),
652 args_preview: "{\"path\":\"a.rs\"}".into(),
653 call_intent: crate::message::ToolCallIntent::new("Inspect source"),
654 };
655 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
656 assert_eq!(v["type"], "tool_node");
657 assert_eq!(v["run_id"], run_id.0.to_string());
658 assert_eq!(v["parent_node_id"], "stmt_2");
659 assert_eq!(v["tool_use_id"], "tu_abc");
660 assert_eq!(v["tool_name"], "fs.read");
661 assert_eq!(v["args_preview"], "{\"path\":\"a.rs\"}");
662 assert_eq!(v["call_intent"], "Inspect source");
663 }
664
665 #[test]
666 fn tool_result_metrics_serialize_raw_and_excerpt_sizes() {
667 let ev = Event::ToolResultMetrics {
668 turn_id: TurnId::now(),
669 flow_run_id: Some(FlowRunId::now()),
670 tool_use_id: "call-1".into(),
671 raw_bytes: 10_000,
672 excerpt_bytes: 1_000,
673 truncated: true,
674 };
675 let value = serde_json::to_value(ev).unwrap();
676
677 assert_eq!(value["type"], "tool_result_metrics");
678 assert_eq!(value["tool_use_id"], "call-1");
679 assert_eq!(value["raw_bytes"], 10_000);
680 assert_eq!(value["excerpt_bytes"], 1_000);
681 assert_eq!(value["truncated"], true);
682 }
683
684 #[test]
685 fn seq_and_set_seq_cover_tool_node() {
686 let _ev = Event::ToolNode {
687 run_id: FlowRunId::now(),
688 parent_node_id: "s".into(),
689 tool_use_id: "t".into(),
690 tool_name: "n".into(),
691 args_preview: "{}".into(),
692 call_intent: None,
693 };
694 }
695
696 #[test]
697 fn attachment_degraded_serializes_all_fields() {
698 let turn = TurnId::now();
699 let flow = FlowRunId::now();
700 let ev = Event::AttachmentDegraded {
701 turn_id: Some(turn.clone()),
702 flow_run_id: Some(flow.clone()),
703 message_seq: 42,
704 part_index: 1,
705 file_basename: "photo.png".into(),
706 reason: "image_too_large".into(),
707 };
708 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
709 assert_eq!(v["type"], "attachment_degraded");
710 assert_eq!(v["message_seq"], 42);
711 assert_eq!(v["part_index"], 1);
712 assert_eq!(v["file_basename"], "photo.png");
713 assert_eq!(v["reason"], "image_too_large");
714 assert_eq!(v["turn_id"], serde_json::json!(turn.0.to_string()));
715 assert_eq!(v["flow_run_id"], serde_json::json!(flow.0.to_string()));
716 }
717
718 #[test]
719 fn seq_and_set_seq_cover_attachment_degraded() {
720 let _ev = Event::AttachmentDegraded {
721 turn_id: None,
722 flow_run_id: None,
723 message_seq: 10,
724 part_index: 0,
725 file_basename: "x".into(),
726 reason: "y".into(),
727 };
728 }
729
730 #[test]
731 fn tool_pending_approval_round_trip() {
732 let ev = Event::ToolPendingApproval {
733 run_id: FlowRunId::now(),
734 tool_use_id: "tu1".into(),
735 tool_name: "fs.write".into(),
736 args_preview: "{}".into(),
737 level: "approve".into(),
738 preview: None,
739 };
740 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
741 assert_eq!(v["type"], "tool_pending_approval");
742 assert_eq!(v["tool_use_id"], "tu1");
743 assert_eq!(v["level"], "approve");
744 }
745
746 #[test]
747 fn seq_and_set_seq_cover_approval_variants() {
748 let rid = FlowRunId::now();
749 for _ev in [
750 Event::ToolPendingApproval {
751 run_id: rid.clone(),
752 tool_use_id: "t".into(),
753 tool_name: "n".into(),
754 args_preview: "{}".into(),
755 level: "approve".into(),
756 preview: None,
757 },
758 Event::ToolApproved {
759 run_id: rid.clone(),
760 tool_use_id: "t".into(),
761 decided_by: "user".into(),
762 },
763 Event::ToolDenied {
764 run_id: rid.clone(),
765 tool_use_id: "t".into(),
766 reason: "no".into(),
767 },
768 ] {}
769 }
770
771 #[test]
772 fn compaction_summary_serializes_all_fields() {
773 let ev = Event::CompactionSummary {
774 session_id: "sess".into(),
775 flow_run_id: None,
776 range_start: 2,
777 range_end: 8,
778 compacted_count: 7,
779 before_tokens: 1000,
780 after_tokens: 250,
781 summary: "gist".into(),
782 };
783 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
784 assert_eq!(v["type"], "compaction_summary");
785 assert_eq!(v["session_id"], "sess");
786 assert_eq!(v["range_start"], 2);
787 assert_eq!(v["range_end"], 8);
788 assert_eq!(v["compacted_count"], 7);
789 assert_eq!(v["before_tokens"], 1000);
790 assert_eq!(v["after_tokens"], 250);
791 assert_eq!(v["summary"], "gist");
792 }
793
794 #[test]
795 fn seq_and_set_seq_cover_compaction_summary() {
796 let _ev = Event::CompactionSummary {
797 session_id: "sess".into(),
798 flow_run_id: None,
799 range_start: 0,
800 range_end: 1,
801 compacted_count: 2,
802 before_tokens: 10,
803 after_tokens: 3,
804 summary: String::new(),
805 };
806 }
807
808 #[test]
809 fn envelope_round_trips_through_json() {
810 let env = EventEnvelope::new(
811 42,
812 Event::UserMsg {
813 turn_id: TurnId::now(),
814 flow_run_id: None,
815 message: crate::message::Message::user_text(TurnId::now(), "hello"),
816 },
817 );
818 let json = serde_json::to_string(&env).unwrap();
819 let back: EventEnvelope = serde_json::from_str(&json).unwrap();
820 assert_eq!(back.seq, 42);
821 assert!(matches!(back.event, Event::UserMsg { .. }));
822 }
823
824 #[test]
825 fn legacy_llm_call_without_context_plan_id_still_deserializes() {
826 let json = r#"{"type":"llm_call","model":"m","provider":"p","usage":{"input":1,"cached_input":0,"output":0,"cache_write":0,"reasoning_tokens":0},"wallclock_ms":1,"ttft_ms":null,"tokens_per_second":null,"status":{"kind":"ok"},"run_id":null,"node_id":null}"#;
827 let event: Event = serde_json::from_str(json).unwrap();
828
829 assert!(matches!(
830 event,
831 Event::LlmCall {
832 context_plan_id: None,
833 context_epoch: None,
834 context_tokens: None,
835 usage_source: None,
836 context_call_purpose: None,
837 context_call_identity: None,
838 context_cache: None,
839 assistant_tool_batch_width: None,
840 ..
841 }
842 ));
843 }
844
845 #[test]
846 fn legacy_system_message_without_flow_owner_still_deserializes() {
847 let turn_id = TurnId::now();
848 let message = crate::message::Message::context_record(
849 turn_id.clone(),
850 crate::context_plan::ContextRecord::new(
851 "session.goal",
852 1,
853 crate::context_plan::ContextRecordAuthority::User,
854 crate::context_plan::ContextRecordRetention::Latest,
855 crate::context_plan::ContextRecordBody::text("ship"),
856 ),
857 );
858 let event: Event = serde_json::from_value(serde_json::json!({
859 "type": "system_msg",
860 "turn_id": turn_id,
861 "message": message,
862 }))
863 .unwrap();
864
865 assert!(matches!(
866 event,
867 Event::SystemMsg {
868 flow_run_id: None,
869 ..
870 }
871 ));
872 }
873
874 #[test]
875 fn legacy_compaction_events_without_flow_owner_still_deserialize() {
876 for value in [
877 serde_json::json!({
878 "type": "compaction_summary",
879 "session_id": "session",
880 "range_start": 0,
881 "range_end": 1,
882 "compacted_count": 2,
883 "before_tokens": 100,
884 "after_tokens": 10,
885 "summary": "summary",
886 }),
887 serde_json::json!({
888 "type": "context_compact",
889 "session_id": "session",
890 "before_tokens": 100,
891 "after_tokens": 10,
892 "compacted_range_start": 0,
893 "compacted_range_end": 1,
894 }),
895 serde_json::json!({
896 "type": "checkpoint",
897 "session_id": "session",
898 "messages": [],
899 "window_tokens": 10,
900 }),
901 ] {
902 let event: Event = serde_json::from_value(value).unwrap();
903 assert!(match event {
904 Event::CompactionSummary { flow_run_id, .. }
905 | Event::ContextCompact { flow_run_id, .. }
906 | Event::Checkpoint { flow_run_id, .. } => flow_run_id.is_none(),
907 _ => false,
908 });
909 }
910 }
911}