//! 事件总线 — 基于 tokio::sync::broadcast 的发布/订阅 use tokio::sync::broadcast; use df_types::events::WorkflowEvent; /// 事件总线 #[derive(Debug)] pub struct EventBus { sender: broadcast::Sender, } /// 事件订阅者 pub type EventSubscriber = broadcast::Receiver; /// 默认事件通道容量 const DEFAULT_CAPACITY: usize = 256; impl EventBus { /// 创建事件总线 pub fn new() -> Self { let (sender, _) = broadcast::channel(DEFAULT_CAPACITY); Self { sender } } /// 创建指定容量的事件总线 pub fn with_capacity(capacity: usize) -> Self { let (sender, _) = broadcast::channel(capacity); Self { sender } } /// 发送事件 pub async fn send(&self, event: WorkflowEvent) { // broadcast::send 是同步的,忽略接收者已关闭的错误 let _ = self.sender.send(event); } /// 订阅事件 pub fn subscribe(&self) -> EventSubscriber { self.sender.subscribe() } } impl Default for EventBus { fn default() -> Self { Self::new() } } impl Clone for EventBus { fn clone(&self) -> Self { Self { sender: self.sender.clone(), } } }