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 },
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 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 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 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}