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, Serialize)]
42#[serde(tag = "type", rename_all = "snake_case")]
43pub enum Event {
44 FlowStart {
45 #[serde(default)]
46 seq: u64,
47 run_id: FlowRunId,
48 flow_name: String,
49 parent_run_id: Option<FlowRunId>,
50 parent_node_id: Option<String>,
51 ts: chrono::DateTime<chrono::Utc>,
52 },
53 FlowEnd {
54 #[serde(default)]
55 seq: u64,
56 run_id: FlowRunId,
57 flow_name: String,
58 status: FlowStatus,
59 ts: chrono::DateTime<chrono::Utc>,
60 },
61 LlmCall {
62 #[serde(default)]
63 seq: u64,
64 model: String,
65 provider: String,
66 usage: crate::provider::TokenUsage,
67 wallclock_ms: u64,
68 ttft_ms: Option<u64>,
69 tokens_per_second: Option<f64>,
70 status: LlmCallStatus,
71 #[serde(default, skip_serializing_if = "Option::is_none")]
72 run_id: Option<crate::event::FlowRunId>,
73 #[serde(default, skip_serializing_if = "Option::is_none")]
74 node_id: Option<String>,
75 ts: chrono::DateTime<chrono::Utc>,
76 },
77 TurnStart {
78 #[serde(default)]
79 seq: u64,
80 turn_id: TurnId,
81 ts: chrono::DateTime<chrono::Utc>,
82 },
83 TurnEnd {
84 #[serde(default)]
85 seq: u64,
86 turn_id: TurnId,
87 ts: chrono::DateTime<chrono::Utc>,
88 },
89 UserMsg {
90 #[serde(default)]
91 seq: u64,
92 turn_id: TurnId,
93 message: crate::message::Message,
94 ts: chrono::DateTime<chrono::Utc>,
95 },
96 AssistantMsg {
97 #[serde(default)]
98 seq: u64,
99 turn_id: TurnId,
100 flow_run_id: Option<FlowRunId>,
101 message: crate::message::Message,
102 ts: chrono::DateTime<chrono::Utc>,
103 },
104 ToolResultMsg {
105 #[serde(default)]
106 seq: u64,
107 turn_id: TurnId,
108 flow_run_id: Option<FlowRunId>,
109 message: crate::message::Message,
110 ts: chrono::DateTime<chrono::Utc>,
111 },
112 DiffPreview {
113 #[serde(default)]
114 seq: u64,
115 turn_id: Option<TurnId>,
116 flow_run_id: Option<FlowRunId>,
117 title: String,
118 old_content: Option<String>,
119 new_content: Option<String>,
120 unified_diff: Option<String>,
121 ts: chrono::DateTime<chrono::Utc>,
122 },
123 CompactionSummary {
124 #[serde(default)]
125 seq: u64,
126 session_id: String,
127 range_start: u64,
128 range_end: u64,
129 compacted_count: usize,
130 before_tokens: u64,
131 after_tokens: u64,
132 summary: String,
133 ts: chrono::DateTime<chrono::Utc>,
134 },
135 SystemMsg {
136 #[serde(default)]
137 seq: u64,
138 turn_id: TurnId,
139 message: crate::message::Message,
140 ts: chrono::DateTime<chrono::Utc>,
141 },
142 UserInject {
143 #[serde(default)]
144 seq: u64,
145 turn_id: TurnId,
146 injection: crate::injection::Injection,
147 ts: chrono::DateTime<chrono::Utc>,
148 },
149 ContentFilterHit {
150 #[serde(default)]
151 seq: u64,
152 turn_id: Option<TurnId>,
153 flow_run_id: Option<FlowRunId>,
154 provider: String,
155 model: String,
156 category: String,
157 action: String,
158 ts: chrono::DateTime<chrono::Utc>,
159 },
160 ContextCompact {
161 #[serde(default)]
162 seq: u64,
163 session_id: String,
164 before_tokens: u64,
165 after_tokens: u64,
166 compacted_range_start: u64,
167 compacted_range_end: u64,
168 #[serde(default, skip_serializing_if = "Option::is_none")]
169 summary_text: Option<String>,
170 #[serde(default, skip_serializing_if = "Option::is_none")]
171 replacement_msg_seq: Option<u64>,
172 ts: chrono::DateTime<chrono::Utc>,
173 },
174 Checkpoint {
175 #[serde(default)]
176 seq: u64,
177 session_id: String,
178 messages: Vec<crate::message::Message>,
179 window_tokens: u64,
180 ts: chrono::DateTime<chrono::Utc>,
181 },
182 ContextTruncated {
183 #[serde(default)]
184 seq: u64,
185 turn_id: Option<TurnId>,
186 flow_run_id: Option<FlowRunId>,
187 original_chars: u64,
188 result_chars: u64,
189 dropped_chars: u64,
190 budget_tokens: u64,
191 ts: chrono::DateTime<chrono::Utc>,
192 },
193 WatchWarn {
194 #[serde(default)]
195 seq: u64,
196 turn_id: Option<TurnId>,
197 flow_run_id: Option<FlowRunId>,
198 target: String,
199 trigger: String,
200 message: String,
201 ts: chrono::DateTime<chrono::Utc>,
202 },
203 LlmPartialCall {
204 #[serde(default)]
205 seq: u64,
206 turn_id: Option<TurnId>,
207 flow_run_id: Option<FlowRunId>,
208 model: String,
209 provider: String,
210 tokens_before_abort: u64,
211 restart_reason: String,
212 ts: chrono::DateTime<chrono::Utc>,
213 },
214 PendingPrompt {
215 #[serde(default)]
216 seq: u64,
217 prompt_id: uuid::Uuid,
218 kind: String,
219 payload: serde_json::Value,
220 ts: chrono::DateTime<chrono::Utc>,
221 },
222 PromptResolved {
223 #[serde(default)]
224 seq: u64,
225 prompt_id: uuid::Uuid,
226 answer: serde_json::Value,
227 ts: chrono::DateTime<chrono::Utc>,
228 },
229 FlowGraph {
230 #[serde(default)]
231 seq: u64,
232 run_id: FlowRunId,
233 graph: crate::nodegraph::FlowGraph,
234 ts: chrono::DateTime<chrono::Utc>,
235 },
236 FlowNodeStart {
237 #[serde(default)]
238 seq: u64,
239 run_id: FlowRunId,
240 node_id: String,
241 kind: crate::nodegraph::NodeKind,
242 label: String,
243 parent_node_id: Option<String>,
244 ts: chrono::DateTime<chrono::Utc>,
245 },
246 FlowNodeEnd {
247 #[serde(default)]
248 seq: u64,
249 run_id: FlowRunId,
250 node_id: String,
251 status: FlowNodeStatus,
252 output_preview: Option<String>,
253 ts: chrono::DateTime<chrono::Utc>,
254 },
255 ToolNode {
256 #[serde(default)]
257 seq: u64,
258 run_id: FlowRunId,
259 parent_node_id: String,
260 tool_use_id: String,
261 tool_name: String,
262 args_preview: String,
263 ts: chrono::DateTime<chrono::Utc>,
264 },
265 AttachmentDegraded {
266 #[serde(default)]
267 seq: u64,
268 turn_id: Option<TurnId>,
269 flow_run_id: Option<FlowRunId>,
270 message_seq: u64,
271 part_index: usize,
272 file_basename: String,
273 reason: String,
274 ts: chrono::DateTime<chrono::Utc>,
275 },
276 ToolPendingApproval {
277 #[serde(default)]
278 seq: u64,
279 run_id: FlowRunId,
280 tool_use_id: String,
281 tool_name: String,
282 args_preview: String,
283 level: String,
284 #[serde(default, skip_serializing_if = "Option::is_none")]
285 preview: Option<String>,
286 ts: chrono::DateTime<chrono::Utc>,
287 },
288 ToolApproved {
289 #[serde(default)]
290 seq: u64,
291 run_id: FlowRunId,
292 tool_use_id: String,
293 decided_by: String,
294 ts: chrono::DateTime<chrono::Utc>,
295 },
296 ToolDenied {
297 #[serde(default)]
298 seq: u64,
299 run_id: FlowRunId,
300 tool_use_id: String,
301 reason: String,
302 ts: chrono::DateTime<chrono::Utc>,
303 },
304}
305
306#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
307#[serde(rename_all = "snake_case")]
308pub enum FlowNodeStatus {
309 Ok,
310 Err,
311 Cancelled,
312}
313
314impl Event {
315 pub fn set_seq(&mut self, new_seq: u64) {
316 match self {
317 Event::FlowStart { seq, .. }
318 | Event::FlowEnd { seq, .. }
319 | Event::LlmCall { seq, .. }
320 | Event::TurnStart { seq, .. }
321 | Event::TurnEnd { seq, .. }
322 | Event::UserMsg { seq, .. }
323 | Event::AssistantMsg { seq, .. }
324 | Event::ToolResultMsg { seq, .. }
325 | Event::DiffPreview { seq, .. }
326 | Event::CompactionSummary { seq, .. }
327 | Event::SystemMsg { seq, .. }
328 | Event::UserInject { seq, .. }
329 | Event::ContentFilterHit { seq, .. }
330 | Event::ContextCompact { seq, .. }
331 | Event::Checkpoint { seq, .. }
332 | Event::ContextTruncated { seq, .. }
333 | Event::WatchWarn { seq, .. }
334 | Event::PendingPrompt { seq, .. }
335 | Event::PromptResolved { seq, .. }
336 | Event::LlmPartialCall { seq, .. }
337 | Event::FlowGraph { seq, .. }
338 | Event::FlowNodeStart { seq, .. }
339 | Event::FlowNodeEnd { seq, .. }
340 | Event::ToolNode { seq, .. }
341 | Event::AttachmentDegraded { seq, .. }
342 | Event::ToolPendingApproval { seq, .. }
343 | Event::ToolApproved { seq, .. }
344 | Event::ToolDenied { seq, .. } => *seq = new_seq,
345 }
346 }
347
348 pub fn seq(&self) -> u64 {
349 match self {
350 Event::FlowStart { seq, .. }
351 | Event::FlowEnd { seq, .. }
352 | Event::LlmCall { seq, .. }
353 | Event::TurnStart { seq, .. }
354 | Event::TurnEnd { seq, .. }
355 | Event::UserMsg { seq, .. }
356 | Event::AssistantMsg { seq, .. }
357 | Event::ToolResultMsg { seq, .. }
358 | Event::DiffPreview { seq, .. }
359 | Event::CompactionSummary { seq, .. }
360 | Event::SystemMsg { seq, .. }
361 | Event::UserInject { seq, .. }
362 | Event::ContentFilterHit { seq, .. }
363 | Event::ContextCompact { seq, .. }
364 | Event::Checkpoint { seq, .. }
365 | Event::ContextTruncated { seq, .. }
366 | Event::WatchWarn { seq, .. }
367 | Event::PendingPrompt { seq, .. }
368 | Event::PromptResolved { seq, .. }
369 | Event::LlmPartialCall { seq, .. }
370 | Event::FlowGraph { seq, .. }
371 | Event::FlowNodeStart { seq, .. }
372 | Event::FlowNodeEnd { seq, .. }
373 | Event::ToolNode { seq, .. }
374 | Event::AttachmentDegraded { seq, .. }
375 | Event::ToolPendingApproval { seq, .. }
376 | Event::ToolApproved { seq, .. }
377 | Event::ToolDenied { seq, .. } => *seq,
378 }
379 }
380}
381
382#[derive(Debug, Clone, Serialize)]
383#[serde(tag = "kind", rename_all = "snake_case")]
384pub enum FlowStatus {
385 Ok,
386 Errored { message: String },
387 Cancelled,
388}
389
390impl FlowStatus {
391 pub fn errored(msg: impl Into<String>) -> Self {
392 Self::Errored {
393 message: msg.into(),
394 }
395 }
396}
397
398#[derive(Debug, Clone, Serialize)]
399#[serde(tag = "kind", rename_all = "snake_case")]
400pub enum LlmCallStatus {
401 Ok,
402 Errored { message: String },
403}
404
405impl LlmCallStatus {
406 pub fn errored(msg: impl Into<String>) -> Self {
407 Self::Errored {
408 message: msg.into(),
409 }
410 }
411}
412
413#[derive(Debug, Clone)]
414pub enum NodeEvent {
415 LlmChunk {
416 text: String,
417 cumulative_tokens: u64,
418 },
419 ThinkingChunk {
420 text: String,
421 },
422 LlmDone {
423 total_tokens: u64,
424 },
425 ToolStdoutLine {
426 line: String,
427 seq: u64,
428 },
429 ToolStderrLine {
430 line: String,
431 seq: u64,
432 },
433 ToolDone {
434 exit: i32,
435 },
436}
437
438pub struct Observable<T> {
439 pub output: crate::tool::BoxFut<'static, Result<T, crate::error::RuntimeError>>,
440 pub events: broadcast::Receiver<NodeEvent>,
441 pub cancel: CancellationToken,
442}
443
444#[derive(Default, Clone)]
445pub struct EventSink {
446 events: Arc<Mutex<Vec<Event>>>,
447 forwarder: Option<mpsc::UnboundedSender<Event>>,
448 seq_counter: Arc<std::sync::atomic::AtomicU64>,
449 redactor: Option<Arc<crate::redact::Redactor>>,
450 last_compact_at: Arc<Mutex<Option<chrono::DateTime<chrono::Utc>>>>,
451}
452
453impl EventSink {
454 pub fn new() -> Self {
455 Self::default()
456 }
457
458 pub fn with_forwarder(mut self, tx: mpsc::UnboundedSender<Event>) -> Self {
459 self.forwarder = Some(tx);
460 self
461 }
462
463 pub fn with_redactor(mut self, redactor: Arc<crate::redact::Redactor>) -> Self {
464 self.redactor = Some(redactor);
465 self
466 }
467
468 pub fn next_seq_peek(&self) -> u64 {
472 self.seq_counter.load(std::sync::atomic::Ordering::SeqCst) + 1
473 }
474
475 pub fn restore_seq(&self, last_seq: u64) {
476 self.seq_counter
477 .store(last_seq, std::sync::atomic::Ordering::SeqCst);
478 }
479
480 pub fn reserve_seq(&self) -> u64 {
484 self.seq_counter
485 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
486 + 1
487 }
488
489 pub fn emit_returning_seq(&self, event: Event) -> u64 {
490 let next = self
491 .seq_counter
492 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
493 + 1;
494 let mut event = event;
495 event.set_seq(next);
496 if let Some(tx) = &self.forwarder {
497 let _ = tx.send(event.clone());
498 }
499 self.events.lock().expect("event sink poisoned").push(event);
500 next
501 }
502
503 pub fn emit(&self, mut event: Event) {
504 let next = self
505 .seq_counter
506 .fetch_add(1, std::sync::atomic::Ordering::SeqCst)
507 + 1;
508 event.set_seq(next);
509 if let Some(tx) = &self.forwarder {
510 let _ = tx.send(event.clone());
511 }
512 self.events.lock().expect("event sink poisoned").push(event);
513 }
514
515 pub fn redactor(&self) -> Option<Arc<crate::redact::Redactor>> {
516 self.redactor.clone()
517 }
518
519 pub fn mark_compacted(&self) {
520 *self.last_compact_at.lock().expect("last_compact poisoned") = Some(chrono::Utc::now());
521 }
522
523 pub fn last_compact_ago_seconds(&self) -> Option<i64> {
524 self.last_compact_at
525 .lock()
526 .expect("last_compact poisoned")
527 .map(|t| (chrono::Utc::now() - t).num_seconds())
528 }
529
530 pub fn drain(&self) -> Vec<Event> {
531 std::mem::take(&mut *self.events.lock().expect("event sink poisoned"))
532 }
533
534 pub fn snapshot(&self) -> Vec<Event> {
535 self.events.lock().expect("event sink poisoned").clone()
536 }
537}
538
539#[cfg(test)]
540mod tests {
541 use super::*;
542
543 #[test]
544 fn flow_start_serializes_parent_linkage() {
545 let parent = FlowRunId::now();
546 let ev = Event::FlowStart {
547 seq: 42,
548 run_id: FlowRunId::now(),
549 flow_name: "child".into(),
550 parent_run_id: Some(parent.clone()),
551 parent_node_id: Some("stmt_3".into()),
552 ts: chrono::Utc::now(),
553 };
554 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
555 assert_eq!(v["type"], "flow_start");
556 assert_eq!(v["parent_run_id"], serde_json::json!(parent.0.to_string()));
557 assert_eq!(v["parent_node_id"], "stmt_3");
558 }
559
560 #[test]
561 fn flow_node_start_carries_parent_node_id() {
562 let ev = Event::FlowNodeStart {
563 seq: 7,
564 run_id: FlowRunId::now(),
565 node_id: "stmt_1.branch[0]".into(),
566 kind: crate::nodegraph::NodeKind::UserConfirm,
567 label: "fanout".into(),
568 parent_node_id: Some("stmt_1".into()),
569 ts: chrono::Utc::now(),
570 };
571 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
572 assert_eq!(v["type"], "flow_node_start");
573 assert_eq!(v["parent_node_id"], "stmt_1");
574 }
575
576 #[test]
577 fn tool_node_serializes_all_fields() {
578 let run_id = FlowRunId::now();
579 let ev = Event::ToolNode {
580 seq: 12,
581 run_id: run_id.clone(),
582 parent_node_id: "stmt_2".into(),
583 tool_use_id: "tu_abc".into(),
584 tool_name: "fs.read".into(),
585 args_preview: "{\"path\":\"a.rs\"}".into(),
586 ts: chrono::Utc::now(),
587 };
588 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
589 assert_eq!(v["type"], "tool_node");
590 assert_eq!(v["run_id"], run_id.0.to_string());
591 assert_eq!(v["parent_node_id"], "stmt_2");
592 assert_eq!(v["tool_use_id"], "tu_abc");
593 assert_eq!(v["tool_name"], "fs.read");
594 assert_eq!(v["args_preview"], "{\"path\":\"a.rs\"}");
595 }
596
597 #[test]
598 fn seq_and_set_seq_cover_tool_node() {
599 let mut ev = Event::ToolNode {
600 seq: 0,
601 run_id: FlowRunId::now(),
602 parent_node_id: "s".into(),
603 tool_use_id: "t".into(),
604 tool_name: "n".into(),
605 args_preview: "{}".into(),
606 ts: chrono::Utc::now(),
607 };
608 ev.set_seq(99);
609 assert_eq!(ev.seq(), 99);
610 }
611
612 #[test]
613 fn attachment_degraded_serializes_all_fields() {
614 let turn = TurnId::now();
615 let flow = FlowRunId::now();
616 let ev = Event::AttachmentDegraded {
617 seq: 55,
618 turn_id: Some(turn.clone()),
619 flow_run_id: Some(flow.clone()),
620 message_seq: 42,
621 part_index: 1,
622 file_basename: "photo.png".into(),
623 reason: "image_too_large".into(),
624 ts: chrono::Utc::now(),
625 };
626 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
627 assert_eq!(v["type"], "attachment_degraded");
628 assert_eq!(v["message_seq"], 42);
629 assert_eq!(v["part_index"], 1);
630 assert_eq!(v["file_basename"], "photo.png");
631 assert_eq!(v["reason"], "image_too_large");
632 assert_eq!(v["turn_id"], serde_json::json!(turn.0.to_string()));
633 assert_eq!(v["flow_run_id"], serde_json::json!(flow.0.to_string()));
634 }
635
636 #[test]
637 fn seq_and_set_seq_cover_attachment_degraded() {
638 let mut ev = Event::AttachmentDegraded {
639 seq: 0,
640 turn_id: None,
641 flow_run_id: None,
642 message_seq: 10,
643 part_index: 0,
644 file_basename: "x".into(),
645 reason: "y".into(),
646 ts: chrono::Utc::now(),
647 };
648 ev.set_seq(101);
649 assert_eq!(ev.seq(), 101);
650 }
651
652 #[test]
653 fn tool_pending_approval_round_trip() {
654 let ev = Event::ToolPendingApproval {
655 seq: 5,
656 run_id: FlowRunId::now(),
657 tool_use_id: "tu1".into(),
658 tool_name: "fs.write".into(),
659 args_preview: "{}".into(),
660 level: "approve".into(),
661 preview: None,
662 ts: chrono::Utc::now(),
663 };
664 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
665 assert_eq!(v["type"], "tool_pending_approval");
666 assert_eq!(v["tool_use_id"], "tu1");
667 assert_eq!(v["level"], "approve");
668 }
669
670 #[test]
671 fn seq_and_set_seq_cover_approval_variants() {
672 let rid = FlowRunId::now();
673 for mut ev in [
674 Event::ToolPendingApproval {
675 seq: 0,
676 run_id: rid.clone(),
677 tool_use_id: "t".into(),
678 tool_name: "n".into(),
679 args_preview: "{}".into(),
680 level: "approve".into(),
681 preview: None,
682 ts: chrono::Utc::now(),
683 },
684 Event::ToolApproved {
685 seq: 0,
686 run_id: rid.clone(),
687 tool_use_id: "t".into(),
688 decided_by: "user".into(),
689 ts: chrono::Utc::now(),
690 },
691 Event::ToolDenied {
692 seq: 0,
693 run_id: rid.clone(),
694 tool_use_id: "t".into(),
695 reason: "no".into(),
696 ts: chrono::Utc::now(),
697 },
698 ] {
699 ev.set_seq(77);
700 assert_eq!(ev.seq(), 77);
701 }
702 }
703
704 #[test]
705 fn compaction_summary_serializes_all_fields() {
706 let ev = Event::CompactionSummary {
707 seq: 9,
708 session_id: "sess".into(),
709 range_start: 2,
710 range_end: 8,
711 compacted_count: 7,
712 before_tokens: 1000,
713 after_tokens: 250,
714 summary: "gist".into(),
715 ts: chrono::Utc::now(),
716 };
717 let v: serde_json::Value = serde_json::to_value(&ev).unwrap();
718 assert_eq!(v["type"], "compaction_summary");
719 assert_eq!(v["session_id"], "sess");
720 assert_eq!(v["range_start"], 2);
721 assert_eq!(v["range_end"], 8);
722 assert_eq!(v["compacted_count"], 7);
723 assert_eq!(v["before_tokens"], 1000);
724 assert_eq!(v["after_tokens"], 250);
725 assert_eq!(v["summary"], "gist");
726 }
727
728 #[test]
729 fn seq_and_set_seq_cover_compaction_summary() {
730 let mut ev = Event::CompactionSummary {
731 seq: 0,
732 session_id: "sess".into(),
733 range_start: 0,
734 range_end: 1,
735 compacted_count: 2,
736 before_tokens: 10,
737 after_tokens: 3,
738 summary: String::new(),
739 ts: chrono::Utc::now(),
740 };
741 ev.set_seq(123);
742 assert_eq!(ev.seq(), 123);
743 }
744}