Skip to main content

roder_core/
bus.rs

1use std::sync::{
2    Arc,
3    atomic::{AtomicU64, Ordering},
4};
5
6use roder_api::events::{EventEnvelope, EventSource, RoderEvent, ThreadId, TurnId};
7use time::OffsetDateTime;
8use tokio::sync::broadcast;
9
10#[derive(Debug, Clone, Default)]
11pub struct EventFilter {
12    pub thread_id: Option<ThreadId>,
13    pub turn_id: Option<TurnId>,
14    pub kinds: Vec<String>,
15    pub sources: Vec<EventSource>,
16}
17
18impl EventFilter {
19    pub fn matches(&self, envelope: &EventEnvelope) -> bool {
20        if let Some(thread_id) = &self.thread_id
21            && envelope.thread_id.as_ref() != Some(thread_id)
22        {
23            return false;
24        }
25        if let Some(turn_id) = &self.turn_id
26            && envelope.turn_id.as_ref() != Some(turn_id)
27        {
28            return false;
29        }
30        if !self.kinds.is_empty() && !self.kinds.iter().any(|kind| kind == &envelope.kind) {
31            return false;
32        }
33        if !self.sources.is_empty() && !self.sources.iter().any(|source| source == &envelope.source)
34        {
35            return false;
36        }
37        true
38    }
39}
40
41#[derive(Clone)]
42pub struct EventBus {
43    sender: broadcast::Sender<EventEnvelope>,
44    next_seq: Arc<AtomicU64>,
45}
46
47impl EventBus {
48    pub fn new(capacity: usize) -> Self {
49        let (sender, _) = broadcast::channel(capacity);
50        Self {
51            sender,
52            next_seq: Arc::new(AtomicU64::new(1)),
53        }
54    }
55
56    pub fn subscribe(&self) -> broadcast::Receiver<EventEnvelope> {
57        self.sender.subscribe()
58    }
59
60    pub fn emit(&self, event: RoderEvent) -> EventEnvelope {
61        let envelope = EventEnvelope {
62            event_id: uuid::Uuid::new_v4().to_string(),
63            seq: self.next_seq.fetch_add(1, Ordering::SeqCst),
64            timestamp: OffsetDateTime::now_utc(),
65            source: event.source(),
66            kind: event.kind().to_string(),
67            thread_id: event.thread_id().cloned(),
68            turn_id: event.turn_id().cloned(),
69            event,
70        };
71        let _ = self.sender.send(envelope.clone());
72        envelope
73    }
74}