1use std::sync::{Arc, Mutex};
2
3use serde::{Deserialize, Serialize};
4use tokio::sync::broadcast;
5use tokio::sync::mpsc;
6use tokio_util::sync::CancellationToken;
7use uuid::Uuid;
8
9#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
10#[serde(transparent)]
11pub struct FlowRunId(pub Uuid);
12
13impl FlowRunId {
14 pub fn now() -> Self {
15 Self(Uuid::now_v7())
16 }
17}
18
19impl std::fmt::Display for FlowRunId {
20 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
21 self.0.fmt(f)
22 }
23}
24
25#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
26#[serde(transparent)]
27pub struct TurnId(pub Uuid);
28
29impl TurnId {
30 pub fn now() -> Self {
31 Self(Uuid::now_v7())
32 }
33}
34
35impl std::fmt::Display for TurnId {
36 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
37 self.0.fmt(f)
38 }
39}
40
41#[derive(Debug, Clone)]
44pub struct EventEnvelope {
45 pub seq: u64,
46 pub ts: chrono::DateTime<chrono::Utc>,
47 pub event: Event,
48}
49
50impl EventEnvelope {
51 pub fn new(seq: u64, event: Event) -> Self {
52 let ts = chrono::Utc::now();
53 Self { seq, ts, event }
54 }
55}
56
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 let seq = value.get("seq").and_then(|v| v.as_u64()).unwrap_or(0);
73 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 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 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 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}