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 title: String,
187 #[serde(default)]
188 old_content: Option<String>,
189 #[serde(default)]
190 new_content: Option<String>,
191 #[serde(default)]
192 unified_diff: Option<String>,
193 },
194 CompactionSummary {
195 session_id: String,
196 #[serde(default, skip_serializing_if = "Option::is_none")]
197 flow_run_id: Option<FlowRunId>,
198 range_start: u64,
199 range_end: u64,
200 compacted_count: usize,
201 before_tokens: u64,
202 after_tokens: u64,
203 summary: String,
204 },
205 SystemMsg {
206 turn_id: TurnId,
207 #[serde(default, skip_serializing_if = "Option::is_none")]
208 flow_run_id: Option<FlowRunId>,
209 message: crate::message::Message,
210 },
211 UserInject {
212 turn_id: TurnId,
213 injection: crate::injection::Injection,
214 },
215 ContentFilterHit {
216 turn_id: Option<TurnId>,
217 flow_run_id: Option<FlowRunId>,
218 provider: String,
219 model: String,
220 category: String,
221 action: String,
222 },
223 ContextCompact {
224 session_id: String,
225 #[serde(default, skip_serializing_if = "Option::is_none")]
226 flow_run_id: Option<FlowRunId>,
227 before_tokens: u64,
228 after_tokens: u64,
229 compacted_range_start: u64,
230 compacted_range_end: u64,
231 #[serde(default, skip_serializing_if = "Option::is_none")]
232 summary_text: Option<String>,
233 #[serde(default, skip_serializing_if = "Option::is_none")]
234 replacement_msg_seq: Option<u64>,
235 },
236 Checkpoint {
237 session_id: String,
238 #[serde(default, skip_serializing_if = "Option::is_none")]
239 flow_run_id: Option<FlowRunId>,
240 messages: Vec<crate::message::Message>,
241 window_tokens: u64,
242 },
243 ContextTruncated {
244 turn_id: Option<TurnId>,
245 flow_run_id: Option<FlowRunId>,
246 original_chars: u64,
247 result_chars: u64,
248 dropped_chars: u64,
249 budget_tokens: u64,
250 },
251 WatchWarn {
252 turn_id: Option<TurnId>,
253 flow_run_id: Option<FlowRunId>,
254 target: String,
255 trigger: String,
256 message: String,
257 },
258 LlmPartialCall {
259 turn_id: Option<TurnId>,
260 flow_run_id: Option<FlowRunId>,
261 model: String,
262 provider: String,
263 tokens_before_abort: u64,
264 restart_reason: String,
265 },
266 PendingPrompt {
267 prompt_id: uuid::Uuid,
268 kind: String,
269 payload: serde_json::Value,
270 },
271 PromptResolved {
272 prompt_id: uuid::Uuid,
273 answer: serde_json::Value,
274 },
275 FlowGraph {
276 run_id: FlowRunId,
277 graph: crate::nodegraph::FlowGraph,
278 },
279 FlowNodeStart {
280 run_id: FlowRunId,
281 node_id: String,
282 #[serde(default = "default_replay_node_kind")]
283 kind: crate::nodegraph::NodeKind,
284 #[serde(default)]
285 label: String,
286 #[serde(default)]
287 parent_node_id: Option<String>,
288 },
289 FlowNodeEnd {
290 run_id: FlowRunId,
291 node_id: String,
292 #[serde(default)]
293 status: FlowNodeStatus,
294 #[serde(default)]
295 output_preview: Option<String>,
296 },
297 ToolNode {
298 run_id: FlowRunId,
299 parent_node_id: String,
300 tool_use_id: String,
301 tool_name: String,
302 args_preview: String,
303 #[serde(default, skip_serializing_if = "Option::is_none")]
304 call_intent: Option<crate::message::ToolCallIntent>,
305 },
306 AttachmentDegraded {
307 turn_id: Option<TurnId>,
308 flow_run_id: Option<FlowRunId>,
309 message_seq: u64,
310 part_index: usize,
311 file_basename: String,
312 reason: String,
313 },
314 ToolPendingApproval {
315 run_id: FlowRunId,
316 tool_use_id: String,
317 tool_name: String,
318 args_preview: String,
319 level: String,
320 #[serde(default, skip_serializing_if = "Option::is_none")]
321 preview: Option<String>,
322 },
323 ToolApproved {
324 run_id: FlowRunId,
325 tool_use_id: String,
326 decided_by: String,
327 },
328 ToolDenied {
329 run_id: FlowRunId,
330 tool_use_id: String,
331 reason: String,
332 },
333 PermissionRequestCreated {
334 payload: crate::permission_audit::PermissionRequestAudit,
335 },
336 PermissionRequestTargeted {
337 payload: crate::permission_audit::PermissionRequestAudit,
338 },
339 PermissionRequestDeferred {
340 payload: crate::permission_audit::PermissionRequestAudit,
341 },
342 PermissionRequestApproved {
343 payload: crate::permission_audit::PermissionRequestAudit,
344 },
345 PermissionRequestDenied {
346 payload: crate::permission_audit::PermissionRequestAudit,
347 },
348 PermissionRequestCancelled {
349 payload: crate::permission_audit::PermissionRequestAudit,
350 },
351 PermissionGroupCreated {
352 payload: crate::permission_audit::PermissionGroupAudit,
353 },
354 PermissionGroupUpdated {
355 payload: crate::permission_audit::PermissionGroupAudit,
356 },
357 PermissionGroupResolved {
358 payload: crate::permission_audit::PermissionGroupAudit,
359 },
360 PermissionGrantCreated {
361 payload: crate::permission_audit::PermissionGrantAudit,
362 },
363 PermissionGrantExpired {
364 payload: crate::permission_audit::PermissionGrantAudit,
365 },
366 UnrestrictedExecution {
367 payload: crate::permission_audit::PermissionRequestAudit,
368 },
369 TerminalFinalState {
373 handle: String,
374 screen: crate::tools::term::TerminalScreen,
375 state: crate::tools::term::TermStateSnapshot,
376 },
377 MermaidDiagram {
378 source: String,
379 },
380}
381
382#[derive(Debug, Clone, Default, Serialize, Deserialize, PartialEq)]
383#[serde(rename_all = "snake_case")]
384pub enum FlowNodeStatus {
385 #[default]
386 Ok,
387 Err,
388 Cancelled,
389}
390
391#[derive(Debug, Clone, Serialize, Deserialize)]
392#[serde(tag = "kind", rename_all = "snake_case")]
393pub enum FlowStatus {
394 Ok,
395 Errored { message: String },
396 Cancelled,
397}
398
399impl FlowStatus {
400 pub fn errored(msg: impl Into<String>) -> Self {
401 Self::Errored {
402 message: msg.into(),
403 }
404 }
405}
406
407#[derive(Debug, Clone, Default, Serialize, Deserialize)]
408#[serde(tag = "kind", rename_all = "snake_case")]
409pub enum LlmCallStatus {
410 #[default]
411 Ok,
412 Errored {
413 message: String,
414 },
415}
416
417fn default_replay_node_kind() -> crate::nodegraph::NodeKind {
418 crate::nodegraph::NodeKind::UserConfirm
419}
420
421impl LlmCallStatus {
422 pub fn errored(msg: impl Into<String>) -> Self {
423 Self::Errored {
424 message: msg.into(),
425 }
426 }
427}
428
429#[derive(Debug, Clone)]
430pub enum NodeEvent {
431 LlmChunk {
432 text: String,
433 cumulative_tokens: u64,
434 },
435 ThinkingChunk {
436 text: String,
437 },
438 LlmDone {
439 total_tokens: u64,
440 },
441 ToolStdoutLine {
442 line: String,
443 },
444 ToolStderrLine {
445 line: String,
446 },
447 ToolDone {
448 exit: i32,
449 },
450}
451
452pub struct Observable<T> {
453 pub output: crate::tool::BoxFut<'static, Result<T, crate::error::RuntimeError>>,
454 pub events: broadcast::Receiver<NodeEvent>,
455 pub cancel: CancellationToken,
456}
457
458#[derive(Default, Clone)]
459pub struct EventSink {
460 events: Arc<Mutex<Vec<EventEnvelope>>>,
461 forwarder: Option<mpsc::UnboundedSender<EventEnvelope>>,
462 seq_counter: Arc<std::sync::atomic::AtomicU64>,
463 redactor: Option<Arc<crate::redact::Redactor>>,
464 last_compact_at: Arc<Mutex<Option<chrono::DateTime<chrono::Utc>>>>,
465}
466
467impl EventSink {
468 pub fn new() -> Self {
469 Self::default()
470 }
471
472 pub fn with_forwarder(mut self, tx: mpsc::UnboundedSender<EventEnvelope>) -> Self {
473 self.forwarder = Some(tx);
474 self
475 }
476
477 pub fn with_redactor(mut self, redactor: Arc<crate::redact::Redactor>) -> Self {
478 self.redactor = Some(redactor);
479 self
480 }
481
482 pub fn next_seq_peek(&self) -> u64 {
486 self.seq_counter.load(std::sync::atomic::Ordering::SeqCst) + 1
487 }
488
489 pub fn restore_seq(&self, last_seq: u64) {
490 self.seq_counter
491 .store(last_seq, std::sync::atomic::Ordering::SeqCst);
492 }
493
494 pub fn reserve_seq(&self) -> u64 {
498 self.seq_counter
499 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
500 + 1
501 }
502
503 pub fn emit_returning_seq(&self, event: Event) -> u64 {
504 let next = self
505 .seq_counter
506 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
507 + 1;
508 let envelope = EventEnvelope::new(next, event);
509 if let Some(tx) = &self.forwarder {
510 let _ = tx.send(envelope.clone());
511 }
512 self.events
513 .lock()
514 .expect("event sink poisoned")
515 .push(envelope);
516 next
517 }
518
519 pub fn emit(&self, event: Event) {
520 let next = self
521 .seq_counter
522 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
523 + 1;
524 let envelope = EventEnvelope::new(next, event);
525 if let Some(tx) = &self.forwarder {
526 let _ = tx.send(envelope.clone());
527 }
528 self.events
529 .lock()
530 .expect("event sink poisoned")
531 .push(envelope);
532 }
533
534 pub fn events_handle(&self) -> Arc<Mutex<Vec<EventEnvelope>>> {
535 self.events.clone()
536 }
537
538 pub fn redactor(&self) -> Option<Arc<crate::redact::Redactor>> {
539 self.redactor.clone()
540 }
541
542 pub fn mark_compacted(&self) {
543 *self.last_compact_at.lock().expect("last_compact poisoned") = Some(chrono::Utc::now());
544 }
545
546 pub fn last_compact_ago_seconds(&self) -> Option<i64> {
547 self.last_compact_at
548 .lock()
549 .expect("last_compact poisoned")
550 .map(|t| (chrono::Utc::now() - t).num_seconds())
551 }
552
553 pub fn drain(&self) -> Vec<Event> {
554 std::mem::take(&mut *self.events.lock().expect("event sink poisoned"))
555 .into_iter()
556 .map(|envelope| envelope.event)
557 .collect()
558 }
559
560 pub fn snapshot(&self) -> Vec<Event> {
561 self.events
562 .lock()
563 .expect("event sink poisoned")
564 .iter()
565 .map(|envelope| envelope.event.clone())
566 .collect()
567 }
568
569 pub fn snapshot_envelopes(&self) -> Vec<EventEnvelope> {
570 self.events
571 .lock()
572 .expect("event sink poisoned")
573 .iter()
574 .cloned()
575 .collect()
576 }
577}
578
579#[cfg(test)]
580mod tests {
581 use super::*;
582
583 #[test]
584 fn flow_start_serializes_parent_linkage() {
585 let parent = FlowRunId::now();
586 let ev = Event::FlowStart {
587 run_id: FlowRunId::now(),
588 flow_name: "child".into(),
589 parent_run_id: Some(parent.clone()),
590 parent_node_id: Some("stmt_3".into()),
591 spawned: false,
592 };
593 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
594 assert_eq!(v["type"], "flow_start");
595 assert_eq!(v["parent_run_id"], serde_json::json!(parent.0.to_string()));
596 assert_eq!(v["parent_node_id"], "stmt_3");
597 }
598
599 #[test]
600 fn flow_node_start_carries_parent_node_id() {
601 let ev = Event::FlowNodeStart {
602 run_id: FlowRunId::now(),
603 node_id: "stmt_1.branch[0]".into(),
604 kind: crate::nodegraph::NodeKind::UserConfirm,
605 label: "fanout".into(),
606 parent_node_id: Some("stmt_1".into()),
607 };
608 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
609 assert_eq!(v["type"], "flow_node_start");
610 assert_eq!(v["parent_node_id"], "stmt_1");
611 }
612
613 #[test]
614 fn tool_node_serializes_all_fields() {
615 let run_id = FlowRunId::now();
616 let ev = Event::ToolNode {
617 run_id: run_id.clone(),
618 parent_node_id: "stmt_2".into(),
619 tool_use_id: "tu_abc".into(),
620 tool_name: "fs.read".into(),
621 args_preview: "{\"path\":\"a.rs\"}".into(),
622 call_intent: crate::message::ToolCallIntent::new("Inspect source"),
623 };
624 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
625 assert_eq!(v["type"], "tool_node");
626 assert_eq!(v["run_id"], run_id.0.to_string());
627 assert_eq!(v["parent_node_id"], "stmt_2");
628 assert_eq!(v["tool_use_id"], "tu_abc");
629 assert_eq!(v["tool_name"], "fs.read");
630 assert_eq!(v["args_preview"], "{\"path\":\"a.rs\"}");
631 assert_eq!(v["call_intent"], "Inspect source");
632 }
633
634 #[test]
635 fn tool_result_metrics_serialize_raw_and_excerpt_sizes() {
636 let ev = Event::ToolResultMetrics {
637 turn_id: TurnId::now(),
638 flow_run_id: Some(FlowRunId::now()),
639 tool_use_id: "call-1".into(),
640 raw_bytes: 10_000,
641 excerpt_bytes: 1_000,
642 truncated: true,
643 };
644 let value = serde_json::to_value(ev).unwrap();
645
646 assert_eq!(value["type"], "tool_result_metrics");
647 assert_eq!(value["tool_use_id"], "call-1");
648 assert_eq!(value["raw_bytes"], 10_000);
649 assert_eq!(value["excerpt_bytes"], 1_000);
650 assert_eq!(value["truncated"], true);
651 }
652
653 #[test]
654 fn seq_and_set_seq_cover_tool_node() {
655 let _ev = Event::ToolNode {
656 run_id: FlowRunId::now(),
657 parent_node_id: "s".into(),
658 tool_use_id: "t".into(),
659 tool_name: "n".into(),
660 args_preview: "{}".into(),
661 call_intent: None,
662 };
663 }
664
665 #[test]
666 fn attachment_degraded_serializes_all_fields() {
667 let turn = TurnId::now();
668 let flow = FlowRunId::now();
669 let ev = Event::AttachmentDegraded {
670 turn_id: Some(turn.clone()),
671 flow_run_id: Some(flow.clone()),
672 message_seq: 42,
673 part_index: 1,
674 file_basename: "photo.png".into(),
675 reason: "image_too_large".into(),
676 };
677 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
678 assert_eq!(v["type"], "attachment_degraded");
679 assert_eq!(v["message_seq"], 42);
680 assert_eq!(v["part_index"], 1);
681 assert_eq!(v["file_basename"], "photo.png");
682 assert_eq!(v["reason"], "image_too_large");
683 assert_eq!(v["turn_id"], serde_json::json!(turn.0.to_string()));
684 assert_eq!(v["flow_run_id"], serde_json::json!(flow.0.to_string()));
685 }
686
687 #[test]
688 fn seq_and_set_seq_cover_attachment_degraded() {
689 let _ev = Event::AttachmentDegraded {
690 turn_id: None,
691 flow_run_id: None,
692 message_seq: 10,
693 part_index: 0,
694 file_basename: "x".into(),
695 reason: "y".into(),
696 };
697 }
698
699 #[test]
700 fn tool_pending_approval_round_trip() {
701 let ev = Event::ToolPendingApproval {
702 run_id: FlowRunId::now(),
703 tool_use_id: "tu1".into(),
704 tool_name: "fs.write".into(),
705 args_preview: "{}".into(),
706 level: "approve".into(),
707 preview: None,
708 };
709 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
710 assert_eq!(v["type"], "tool_pending_approval");
711 assert_eq!(v["tool_use_id"], "tu1");
712 assert_eq!(v["level"], "approve");
713 }
714
715 #[test]
716 fn seq_and_set_seq_cover_approval_variants() {
717 let rid = FlowRunId::now();
718 for _ev in [
719 Event::ToolPendingApproval {
720 run_id: rid.clone(),
721 tool_use_id: "t".into(),
722 tool_name: "n".into(),
723 args_preview: "{}".into(),
724 level: "approve".into(),
725 preview: None,
726 },
727 Event::ToolApproved {
728 run_id: rid.clone(),
729 tool_use_id: "t".into(),
730 decided_by: "user".into(),
731 },
732 Event::ToolDenied {
733 run_id: rid.clone(),
734 tool_use_id: "t".into(),
735 reason: "no".into(),
736 },
737 ] {}
738 }
739
740 #[test]
741 fn compaction_summary_serializes_all_fields() {
742 let ev = Event::CompactionSummary {
743 session_id: "sess".into(),
744 flow_run_id: None,
745 range_start: 2,
746 range_end: 8,
747 compacted_count: 7,
748 before_tokens: 1000,
749 after_tokens: 250,
750 summary: "gist".into(),
751 };
752 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
753 assert_eq!(v["type"], "compaction_summary");
754 assert_eq!(v["session_id"], "sess");
755 assert_eq!(v["range_start"], 2);
756 assert_eq!(v["range_end"], 8);
757 assert_eq!(v["compacted_count"], 7);
758 assert_eq!(v["before_tokens"], 1000);
759 assert_eq!(v["after_tokens"], 250);
760 assert_eq!(v["summary"], "gist");
761 }
762
763 #[test]
764 fn seq_and_set_seq_cover_compaction_summary() {
765 let _ev = Event::CompactionSummary {
766 session_id: "sess".into(),
767 flow_run_id: None,
768 range_start: 0,
769 range_end: 1,
770 compacted_count: 2,
771 before_tokens: 10,
772 after_tokens: 3,
773 summary: String::new(),
774 };
775 }
776
777 #[test]
778 fn envelope_round_trips_through_json() {
779 let env = EventEnvelope::new(
780 42,
781 Event::UserMsg {
782 turn_id: TurnId::now(),
783 flow_run_id: None,
784 message: crate::message::Message::user_text(TurnId::now(), "hello"),
785 },
786 );
787 let json = serde_json::to_string(&env).unwrap();
788 let back: EventEnvelope = serde_json::from_str(&json).unwrap();
789 assert_eq!(back.seq, 42);
790 assert!(matches!(back.event, Event::UserMsg { .. }));
791 }
792
793 #[test]
794 fn legacy_llm_call_without_context_plan_id_still_deserializes() {
795 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}"#;
796 let event: Event = serde_json::from_str(json).unwrap();
797
798 assert!(matches!(
799 event,
800 Event::LlmCall {
801 context_plan_id: None,
802 context_epoch: None,
803 context_tokens: None,
804 usage_source: None,
805 context_call_purpose: None,
806 context_call_identity: None,
807 context_cache: None,
808 assistant_tool_batch_width: None,
809 ..
810 }
811 ));
812 }
813
814 #[test]
815 fn legacy_system_message_without_flow_owner_still_deserializes() {
816 let turn_id = TurnId::now();
817 let message = crate::message::Message::context_record(
818 turn_id.clone(),
819 crate::context_plan::ContextRecord::new(
820 "session.goal",
821 1,
822 crate::context_plan::ContextRecordAuthority::User,
823 crate::context_plan::ContextRecordRetention::Latest,
824 crate::context_plan::ContextRecordBody::text("ship"),
825 ),
826 );
827 let event: Event = serde_json::from_value(serde_json::json!({
828 "type": "system_msg",
829 "turn_id": turn_id,
830 "message": message,
831 }))
832 .unwrap();
833
834 assert!(matches!(
835 event,
836 Event::SystemMsg {
837 flow_run_id: None,
838 ..
839 }
840 ));
841 }
842
843 #[test]
844 fn legacy_compaction_events_without_flow_owner_still_deserialize() {
845 for value in [
846 serde_json::json!({
847 "type": "compaction_summary",
848 "session_id": "session",
849 "range_start": 0,
850 "range_end": 1,
851 "compacted_count": 2,
852 "before_tokens": 100,
853 "after_tokens": 10,
854 "summary": "summary",
855 }),
856 serde_json::json!({
857 "type": "context_compact",
858 "session_id": "session",
859 "before_tokens": 100,
860 "after_tokens": 10,
861 "compacted_range_start": 0,
862 "compacted_range_end": 1,
863 }),
864 serde_json::json!({
865 "type": "checkpoint",
866 "session_id": "session",
867 "messages": [],
868 "window_tokens": 10,
869 }),
870 ] {
871 let event: Event = serde_json::from_value(value).unwrap();
872 assert!(match event {
873 Event::CompactionSummary { flow_run_id, .. }
874 | Event::ContextCompact { flow_run_id, .. }
875 | Event::Checkpoint { flow_run_id, .. } => flow_run_id.is_none(),
876 _ => false,
877 });
878 }
879 }
880}