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
57impl serde::Serialize for EventEnvelope {
58    fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
59        let mut value = serde_json::to_value(&self.event).map_err(serde::ser::Error::custom)?;
60        if let serde_json::Value::Object(ref mut map) = value {
61            map.insert("seq".into(), serde_json::Value::Number(self.seq.into()));
62            map.insert("ts".into(), serde_json::Value::String(self.ts.to_rfc3339()));
63        }
64        value.serialize(serializer)
65    }
66}
67
68impl<'de> serde::Deserialize<'de> for EventEnvelope {
69    fn deserialize<D: serde::Deserializer<'de>>(deserializer: D) -> Result<Self, D::Error> {
70        let value = serde_json::Value::deserialize(deserializer)?;
71        // seq: required for new events, default 0 for legacy backward compat
72        let seq = value.get("seq").and_then(|v| v.as_u64()).unwrap_or(0);
73        // ts: parse when present; fall back to now for legacy events
74        let ts = value
75            .get("ts")
76            .and_then(|v| v.as_str())
77            .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
78            .map(|dt| dt.with_timezone(&chrono::Utc))
79            .unwrap_or_else(chrono::Utc::now);
80        let event = serde_json::from_value(value).map_err(serde::de::Error::custom)?;
81        Ok(EventEnvelope { seq, ts, event })
82    }
83}
84
85#[derive(Debug, Clone, Serialize, Deserialize)]
86#[serde(tag = "type", rename_all = "snake_case")]
87pub enum Event {
88    FlowStart {
89        run_id: FlowRunId,
90        flow_name: String,
91        parent_run_id: Option<FlowRunId>,
92        parent_node_id: Option<String>,
93    },
94    FlowEnd {
95        run_id: FlowRunId,
96        flow_name: String,
97        status: FlowStatus,
98    },
99    LlmCall {
100        model: String,
101        provider: String,
102        usage: crate::provider::TokenUsage,
103        wallclock_ms: u64,
104        ttft_ms: Option<u64>,
105        tokens_per_second: Option<f64>,
106        status: LlmCallStatus,
107        #[serde(default, skip_serializing_if = "Option::is_none")]
108        run_id: Option<crate::event::FlowRunId>,
109        #[serde(default, skip_serializing_if = "Option::is_none")]
110        node_id: Option<String>,
111    },
112    TurnStart {
113        turn_id: TurnId,
114    },
115    TurnEnd {
116        turn_id: TurnId,
117    },
118    UserMsg {
119        turn_id: TurnId,
120        message: crate::message::Message,
121    },
122    AssistantMsg {
123        turn_id: TurnId,
124        flow_run_id: Option<FlowRunId>,
125        message: crate::message::Message,
126    },
127    ToolResultMsg {
128        turn_id: TurnId,
129        flow_run_id: Option<FlowRunId>,
130        message: crate::message::Message,
131    },
132    DiffPreview {
133        turn_id: Option<TurnId>,
134        flow_run_id: Option<FlowRunId>,
135        title: String,
136        old_content: Option<String>,
137        new_content: Option<String>,
138        unified_diff: Option<String>,
139    },
140    CompactionSummary {
141        session_id: String,
142        range_start: u64,
143        range_end: u64,
144        compacted_count: usize,
145        before_tokens: u64,
146        after_tokens: u64,
147        summary: String,
148    },
149    SystemMsg {
150        turn_id: TurnId,
151        message: crate::message::Message,
152    },
153    UserInject {
154        turn_id: TurnId,
155        injection: crate::injection::Injection,
156    },
157    ContentFilterHit {
158        turn_id: Option<TurnId>,
159        flow_run_id: Option<FlowRunId>,
160        provider: String,
161        model: String,
162        category: String,
163        action: String,
164    },
165    ContextCompact {
166        session_id: String,
167        before_tokens: u64,
168        after_tokens: u64,
169        compacted_range_start: u64,
170        compacted_range_end: u64,
171        #[serde(default, skip_serializing_if = "Option::is_none")]
172        summary_text: Option<String>,
173        #[serde(default, skip_serializing_if = "Option::is_none")]
174        replacement_msg_seq: Option<u64>,
175    },
176    Checkpoint {
177        session_id: String,
178        messages: Vec<crate::message::Message>,
179        window_tokens: u64,
180    },
181    ContextTruncated {
182        turn_id: Option<TurnId>,
183        flow_run_id: Option<FlowRunId>,
184        original_chars: u64,
185        result_chars: u64,
186        dropped_chars: u64,
187        budget_tokens: u64,
188    },
189    WatchWarn {
190        turn_id: Option<TurnId>,
191        flow_run_id: Option<FlowRunId>,
192        target: String,
193        trigger: String,
194        message: String,
195    },
196    LlmPartialCall {
197        turn_id: Option<TurnId>,
198        flow_run_id: Option<FlowRunId>,
199        model: String,
200        provider: String,
201        tokens_before_abort: u64,
202        restart_reason: String,
203    },
204    PendingPrompt {
205        prompt_id: uuid::Uuid,
206        kind: String,
207        payload: serde_json::Value,
208    },
209    PromptResolved {
210        prompt_id: uuid::Uuid,
211        answer: serde_json::Value,
212    },
213    FlowGraph {
214        run_id: FlowRunId,
215        graph: crate::nodegraph::FlowGraph,
216    },
217    FlowNodeStart {
218        run_id: FlowRunId,
219        node_id: String,
220        kind: crate::nodegraph::NodeKind,
221        label: String,
222        parent_node_id: Option<String>,
223    },
224    FlowNodeEnd {
225        run_id: FlowRunId,
226        node_id: String,
227        status: FlowNodeStatus,
228        output_preview: Option<String>,
229    },
230    ToolNode {
231        run_id: FlowRunId,
232        parent_node_id: String,
233        tool_use_id: String,
234        tool_name: String,
235        args_preview: String,
236    },
237    AttachmentDegraded {
238        turn_id: Option<TurnId>,
239        flow_run_id: Option<FlowRunId>,
240        message_seq: u64,
241        part_index: usize,
242        file_basename: String,
243        reason: String,
244    },
245    ToolPendingApproval {
246        run_id: FlowRunId,
247        tool_use_id: String,
248        tool_name: String,
249        args_preview: String,
250        level: String,
251        #[serde(default, skip_serializing_if = "Option::is_none")]
252        preview: Option<String>,
253    },
254    ToolApproved {
255        run_id: FlowRunId,
256        tool_use_id: String,
257        decided_by: String,
258    },
259    ToolDenied {
260        run_id: FlowRunId,
261        tool_use_id: String,
262        reason: String,
263    },
264    /// Persisted when a terminal's reader loop exits (normal exit or kill).
265    /// Carries the last screen state so TUI restore can show it instead of
266    /// the empty placeholder from the spawn-time ToolResultMsg.
267    TerminalFinalState {
268        handle: String,
269        screen: crate::tools::term::TerminalScreen,
270        state: crate::tools::term::TermStateSnapshot,
271    },
272    MermaidDiagram {
273        source: String,
274    },
275}
276
277#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
278#[serde(rename_all = "snake_case")]
279pub enum FlowNodeStatus {
280    Ok,
281    Err,
282    Cancelled,
283}
284
285#[derive(Debug, Clone, Serialize, Deserialize)]
286#[serde(tag = "kind", rename_all = "snake_case")]
287pub enum FlowStatus {
288    Ok,
289    Errored { message: String },
290    Cancelled,
291}
292
293impl FlowStatus {
294    pub fn errored(msg: impl Into<String>) -> Self {
295        Self::Errored {
296            message: msg.into(),
297        }
298    }
299}
300
301#[derive(Debug, Clone, Serialize, Deserialize)]
302#[serde(tag = "kind", rename_all = "snake_case")]
303pub enum LlmCallStatus {
304    Ok,
305    Errored { message: String },
306}
307
308impl LlmCallStatus {
309    pub fn errored(msg: impl Into<String>) -> Self {
310        Self::Errored {
311            message: msg.into(),
312        }
313    }
314}
315
316#[derive(Debug, Clone)]
317pub enum NodeEvent {
318    LlmChunk {
319        text: String,
320        cumulative_tokens: u64,
321    },
322    ThinkingChunk {
323        text: String,
324    },
325    LlmDone {
326        total_tokens: u64,
327    },
328    ToolStdoutLine {
329        line: String,
330    },
331    ToolStderrLine {
332        line: String,
333    },
334    ToolDone {
335        exit: i32,
336    },
337}
338
339pub struct Observable<T> {
340    pub output: crate::tool::BoxFut<'static, Result<T, crate::error::RuntimeError>>,
341    pub events: broadcast::Receiver<NodeEvent>,
342    pub cancel: CancellationToken,
343}
344
345#[derive(Default, Clone)]
346pub struct EventSink {
347    events: Arc<Mutex<Vec<EventEnvelope>>>,
348    forwarder: Option<mpsc::UnboundedSender<EventEnvelope>>,
349    seq_counter: Arc<std::sync::atomic::AtomicU64>,
350    redactor: Option<Arc<crate::redact::Redactor>>,
351    last_compact_at: Arc<Mutex<Option<chrono::DateTime<chrono::Utc>>>>,
352}
353
354impl EventSink {
355    pub fn new() -> Self {
356        Self::default()
357    }
358
359    pub fn with_forwarder(mut self, tx: mpsc::UnboundedSender<EventEnvelope>) -> Self {
360        self.forwarder = Some(tx);
361        self
362    }
363
364    pub fn with_redactor(mut self, redactor: Arc<crate::redact::Redactor>) -> Self {
365        self.redactor = Some(redactor);
366        self
367    }
368
369    // Best-effort peek for anchor labels. NOT reserved: two peekers see the same value.
370    // Safe today because eval is single-threaded per flow (see tool.rs comment on !Send).
371    // If parallel-tool dispatch is introduced, switch call sites to reserve_seq.
372    pub fn next_seq_peek(&self) -> u64 {
373        self.seq_counter.load(std::sync::atomic::Ordering::SeqCst) + 1
374    }
375
376    pub fn restore_seq(&self, last_seq: u64) {
377        self.seq_counter
378            .store(last_seq, std::sync::atomic::Ordering::SeqCst);
379    }
380
381    // Atomic reserve for the future parallel-dispatch case: returns a seq value that
382    // no other reservation can obtain, at the cost of advancing the counter even if
383    // the caller never emits (a hole in seq numbering). Not used yet; kept ready.
384    pub fn reserve_seq(&self) -> u64 {
385        self.seq_counter
386            .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
387            + 1
388    }
389
390    pub fn emit_returning_seq(&self, event: Event) -> u64 {
391        let next = self
392            .seq_counter
393            .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
394            + 1;
395        let envelope = EventEnvelope::new(next, event);
396        if let Some(tx) = &self.forwarder {
397            let _ = tx.send(envelope.clone());
398        }
399        self.events
400            .lock()
401            .expect("event sink poisoned")
402            .push(envelope);
403        next
404    }
405
406    pub fn emit(&self, event: Event) {
407        let next = self
408            .seq_counter
409            .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
410            + 1;
411        let envelope = EventEnvelope::new(next, event);
412        if let Some(tx) = &self.forwarder {
413            let _ = tx.send(envelope.clone());
414        }
415        self.events
416            .lock()
417            .expect("event sink poisoned")
418            .push(envelope);
419    }
420
421    pub fn events_handle(&self) -> Arc<Mutex<Vec<EventEnvelope>>> {
422        self.events.clone()
423    }
424
425    pub fn redactor(&self) -> Option<Arc<crate::redact::Redactor>> {
426        self.redactor.clone()
427    }
428
429    pub fn mark_compacted(&self) {
430        *self.last_compact_at.lock().expect("last_compact poisoned") = Some(chrono::Utc::now());
431    }
432
433    pub fn last_compact_ago_seconds(&self) -> Option<i64> {
434        self.last_compact_at
435            .lock()
436            .expect("last_compact poisoned")
437            .map(|t| (chrono::Utc::now() - t).num_seconds())
438    }
439
440    pub fn drain(&self) -> Vec<Event> {
441        std::mem::take(&mut *self.events.lock().expect("event sink poisoned"))
442            .into_iter()
443            .map(|envelope| envelope.event)
444            .collect()
445    }
446
447    pub fn snapshot(&self) -> Vec<Event> {
448        self.events
449            .lock()
450            .expect("event sink poisoned")
451            .iter()
452            .map(|envelope| envelope.event.clone())
453            .collect()
454    }
455
456    pub fn snapshot_envelopes(&self) -> Vec<EventEnvelope> {
457        self.events
458            .lock()
459            .expect("event sink poisoned")
460            .iter()
461            .cloned()
462            .collect()
463    }
464}
465
466#[cfg(test)]
467mod tests {
468    use super::*;
469
470    #[test]
471    fn flow_start_serializes_parent_linkage() {
472        let parent = FlowRunId::now();
473        let ev = Event::FlowStart {
474            run_id: FlowRunId::now(),
475            flow_name: "child".into(),
476            parent_run_id: Some(parent.clone()),
477            parent_node_id: Some("stmt_3".into()),
478        };
479        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
480        assert_eq!(v["type"], "flow_start");
481        assert_eq!(v["parent_run_id"], serde_json::json!(parent.0.to_string()));
482        assert_eq!(v["parent_node_id"], "stmt_3");
483    }
484
485    #[test]
486    fn flow_node_start_carries_parent_node_id() {
487        let ev = Event::FlowNodeStart {
488            run_id: FlowRunId::now(),
489            node_id: "stmt_1.branch[0]".into(),
490            kind: crate::nodegraph::NodeKind::UserConfirm,
491            label: "fanout".into(),
492            parent_node_id: Some("stmt_1".into()),
493        };
494        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
495        assert_eq!(v["type"], "flow_node_start");
496        assert_eq!(v["parent_node_id"], "stmt_1");
497    }
498
499    #[test]
500    fn tool_node_serializes_all_fields() {
501        let run_id = FlowRunId::now();
502        let ev = Event::ToolNode {
503            run_id: run_id.clone(),
504            parent_node_id: "stmt_2".into(),
505            tool_use_id: "tu_abc".into(),
506            tool_name: "fs.read".into(),
507            args_preview: "{\"path\":\"a.rs\"}".into(),
508        };
509        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
510        assert_eq!(v["type"], "tool_node");
511        assert_eq!(v["run_id"], run_id.0.to_string());
512        assert_eq!(v["parent_node_id"], "stmt_2");
513        assert_eq!(v["tool_use_id"], "tu_abc");
514        assert_eq!(v["tool_name"], "fs.read");
515        assert_eq!(v["args_preview"], "{\"path\":\"a.rs\"}");
516    }
517
518    #[test]
519    fn seq_and_set_seq_cover_tool_node() {
520        let _ev = Event::ToolNode {
521            run_id: FlowRunId::now(),
522            parent_node_id: "s".into(),
523            tool_use_id: "t".into(),
524            tool_name: "n".into(),
525            args_preview: "{}".into(),
526        };
527    }
528
529    #[test]
530    fn attachment_degraded_serializes_all_fields() {
531        let turn = TurnId::now();
532        let flow = FlowRunId::now();
533        let ev = Event::AttachmentDegraded {
534            turn_id: Some(turn.clone()),
535            flow_run_id: Some(flow.clone()),
536            message_seq: 42,
537            part_index: 1,
538            file_basename: "photo.png".into(),
539            reason: "image_too_large".into(),
540        };
541        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
542        assert_eq!(v["type"], "attachment_degraded");
543        assert_eq!(v["message_seq"], 42);
544        assert_eq!(v["part_index"], 1);
545        assert_eq!(v["file_basename"], "photo.png");
546        assert_eq!(v["reason"], "image_too_large");
547        assert_eq!(v["turn_id"], serde_json::json!(turn.0.to_string()));
548        assert_eq!(v["flow_run_id"], serde_json::json!(flow.0.to_string()));
549    }
550
551    #[test]
552    fn seq_and_set_seq_cover_attachment_degraded() {
553        let _ev = Event::AttachmentDegraded {
554            turn_id: None,
555            flow_run_id: None,
556            message_seq: 10,
557            part_index: 0,
558            file_basename: "x".into(),
559            reason: "y".into(),
560        };
561    }
562
563    #[test]
564    fn tool_pending_approval_round_trip() {
565        let ev = Event::ToolPendingApproval {
566            run_id: FlowRunId::now(),
567            tool_use_id: "tu1".into(),
568            tool_name: "fs.write".into(),
569            args_preview: "{}".into(),
570            level: "approve".into(),
571            preview: None,
572        };
573        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
574        assert_eq!(v["type"], "tool_pending_approval");
575        assert_eq!(v["tool_use_id"], "tu1");
576        assert_eq!(v["level"], "approve");
577    }
578
579    #[test]
580    fn seq_and_set_seq_cover_approval_variants() {
581        let rid = FlowRunId::now();
582        for _ev in [
583            Event::ToolPendingApproval {
584                run_id: rid.clone(),
585                tool_use_id: "t".into(),
586                tool_name: "n".into(),
587                args_preview: "{}".into(),
588                level: "approve".into(),
589                preview: None,
590            },
591            Event::ToolApproved {
592                run_id: rid.clone(),
593                tool_use_id: "t".into(),
594                decided_by: "user".into(),
595            },
596            Event::ToolDenied {
597                run_id: rid.clone(),
598                tool_use_id: "t".into(),
599                reason: "no".into(),
600            },
601        ] {}
602    }
603
604    #[test]
605    fn compaction_summary_serializes_all_fields() {
606        let ev = Event::CompactionSummary {
607            session_id: "sess".into(),
608            range_start: 2,
609            range_end: 8,
610            compacted_count: 7,
611            before_tokens: 1000,
612            after_tokens: 250,
613            summary: "gist".into(),
614        };
615        let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
616        assert_eq!(v["type"], "compaction_summary");
617        assert_eq!(v["session_id"], "sess");
618        assert_eq!(v["range_start"], 2);
619        assert_eq!(v["range_end"], 8);
620        assert_eq!(v["compacted_count"], 7);
621        assert_eq!(v["before_tokens"], 1000);
622        assert_eq!(v["after_tokens"], 250);
623        assert_eq!(v["summary"], "gist");
624    }
625
626    #[test]
627    fn seq_and_set_seq_cover_compaction_summary() {
628        let _ev = Event::CompactionSummary {
629            session_id: "sess".into(),
630            range_start: 0,
631            range_end: 1,
632            compacted_count: 2,
633            before_tokens: 10,
634            after_tokens: 3,
635            summary: String::new(),
636        };
637    }
638
639    #[test]
640    fn envelope_round_trips_through_json() {
641        let env = EventEnvelope::new(
642            42,
643            Event::UserMsg {
644                turn_id: TurnId::now(),
645                message: crate::message::Message::user_text(TurnId::now(), "hello"),
646            },
647        );
648        let json = serde_json::to_string(&env).unwrap();
649        let back: EventEnvelope = serde_json::from_str(&json).unwrap();
650        assert_eq!(back.seq, 42);
651        assert!(matches!(back.event, Event::UserMsg { .. }));
652    }
653}