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        #[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    /// Persisted when a terminal's reader loop exits (normal exit or kill).
395    /// Carries the last screen state so TUI restore can show it instead of
396    /// the empty placeholder from the spawn-time ToolResultMsg.
397    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    // Best-effort peek for anchor labels. NOT reserved: two peekers see the same value.
514    // Safe while evaluator dispatch remains sequential per flow. Parallel dispatch must
515    // reserve sequence numbers instead.
516    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    // Atomic reserve for the future parallel-dispatch case: returns a seq value that
526    // no other reservation can obtain, at the cost of advancing the counter even if
527    // the caller never emits (a hole in seq numbering). Not used yet; kept ready.
528    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}