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