Skip to main content

atman_runtime/
event.rs

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/// Wraps an Event with assigned sequence number and timestamp.
42/// Serializes to the same JSONL format as the flat Event for backward compat.
43#[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    /// Persisted when a terminal's reader loop exits (normal exit or kill).
370    /// Carries the last screen state so TUI restore can show it instead of
371    /// the empty placeholder from the spawn-time ToolResultMsg.
372    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    // Best-effort peek for anchor labels. NOT reserved: two peekers see the same value.
483    // Safe while evaluator dispatch remains sequential per flow. Parallel dispatch must
484    // reserve sequence numbers instead.
485    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    // Atomic reserve for the future parallel-dispatch case: returns a seq value that
495    // no other reservation can obtain, at the cost of advancing the counter even if
496    // the caller never emits (a hole in seq numbering). Not used yet; kept ready.
497    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}