//! 全局事件数据总线 — pub-sub 骨架(载荷 serde_json::Value 透传) //! //! 关联设计:docs/02-架构设计/专项设计/全局事件数据总线-2026-06-21.md //! 跨端透传:docs/02-架构设计/已编号方案/F-260622-01-跨端AIChat-Phase3联调设计-2026-06-22.md(§2.1 D2/D3) //! //! 本模块是 **统一事件数据总线(Unified Event Bus)** 的后端骨架,聚焦 **pub-sub**(发布订阅) //! 语义。背景:AI 事件经 `app_handle.emit("ai-chat-event", ...)` 散布在各处(55 处/7 文件), //! 无统一总线,跨模块直接耦合,阻碍跨端透传(df-tunnel)。 //! //! ## 载荷决策(Phase3 D2 = 全 19 变体透传) //! 载荷用 `serde_json::Value`(即 AiChatEvent 序列化 JSON),**不用**强类型 enum 子集。 //! 理由:跨端透传需 AiChatEvent 全 19 变体(含 ApprovalRequired/Compressed 等),早期骨架的 //! AiBusEvent 8 变体强类型子集不足以覆盖,且总线在透传场景不该解析业务语义(传输层天职)。 //! 改 Value 后:tunnel subscriber 收 Value 原样转发 miniapp,零变体丢失,零业务耦合。 //! //! ## 本批范围(L3 骨架,纯新增无行为变化) //! - `EventBus` 结构(基于 tokio::sync::broadcast,容量 256,载荷 `serde_json::Value`) //! - `subscribe()` / `publish()` 基础接口 //! - `EVENT_BUS_ENABLED` 开关(预留,默认 on —— 总线骨架自带,接入事件源时此开关控制是否经总线) //! - pub/sub 基础单元测试 //! //! ## 不做(留后续批次) //! - 不接入 55 处 emit 点(批3 双写:publish + 原 app.emit 并行,关开关回退) //! - 不接 df-tunnel subscriber(Phase3 阶段2:tunnel 注册 subscriber,收 Value 转发 miniapp) //! - 不实现 request-reply / 流式 reply(设计阶段2/5,本批仅 pub-sub) //! //! ## 开关与兜底(渐进可回退) //! - `EVENT_BUS_ENABLED: bool`(默认 true):接入事件源后,关闭则 publish 静默丢弃(不影响原 emit 路径)。 //! 骨架阶段未接入,无实际效果,预留作接入期的快速回退开关。 //! - publish 返回 `usize`(接收者数量),调用方可忽略;无接收者时不报错(broadcast 语义)。 //! - 慢消费者(broadcast 满):`send` 返 `Err(SendError)`,本骨架 `publish` 静默丢弃并返回 0, //! 不 panic(对齐设计待决策6「背压/限流」的保守默认,正式背压策略留后续)。 //! //! ## 与 df-workflow EventBus 关系 //! df-workflow 的 `EventBus`(`crates/df-workflow/src/eventbus.rs`)是工作流节点间专用总线 //! (仅 `WorkflowEvent`),本总线是面向 AI 域的通用 pub-sub(透传 AiChatEvent 全量 Value)。 //! AppState 分两个字段:`event_bus`(工作流)+ `ai_event_bus`(AI 域),互不干扰。 use tokio::sync::broadcast; use super::AiChatEvent; // ============================================================ // 开关(EVENT_BUS_ENABLED,预留,默认 on) // // 机制:接入事件源后,关闭则 EventBus::publish 静默丢弃事件(原 app.emit 路径不受影响, // 双写桥接由接入批负责,关闭总线仅丢总线副本)。骨架阶段未接入,开关预留作接入期快速回退。 // // 兜底:即便开关误关,总线静默丢事件不影响原 emit 路径(AiChatEvent 经 app.emit 仍正常推送前端)。 // ============================================================ /// 事件总线开关(预留,默认 on)。 /// /// 接入实际事件源后,`EventBus::publish` 在 `false` 时静默丢弃(返回 0),原 `app.emit` /// 路径不受影响。骨架阶段未接入,无实际效果。 /// /// dead_code 说明:本批为骨架,EVENT_BUS_ENABLED 暂无消费方(未接入 emit 点); /// 标 allow 保留作批3 接入期的快速回退开关(零调用方≠垃圾,预留保留)。 #[allow(dead_code)] pub const EVENT_BUS_ENABLED: bool = true; /// 默认事件通道容量(对齐 df-workflow EventBus 默认 256)。 /// /// 容量权衡:太小 → 慢消费者丢事件(订阅者消费不及时);太大 → 内存占用。 /// 256 是 broadcast 常见默认值,覆盖典型 AI 流场景(文本片段 + 工具调用 + 心跳混合)。 /// /// dead_code 说明:骨架阶段 EventBus::new() 内联用字面量路径常量,但常量本身供 /// 外部自定义容量场景引用,标 allow 保留(预留)。 #[allow(dead_code)] pub const DEFAULT_BUS_CAPACITY: usize = 256; // ============================================================ // EventBus — 基于 tokio::sync::broadcast 的 pub-sub 总线(载荷 serde_json::Value) // // 选 broadcast 而非 mpsc: // - pub-sub 多订阅者(mpsc 单消费者不满足「多模块订阅同一事件」场景) // - broadcast 容量固定,慢消费者丢老事件而非阻塞发布者(对齐流式场景:宁可丢老 chunk 不阻塞 LLM 流) // // 载荷 serde_json::Value(Phase3 D2):透传 AiChatEvent 全 19 变体原样 JSON,subscriber(tunnel) // 收到后零解析转发 miniapp。总线不持业务类型知识(传输层天职),扩展 AiChatEvent 变体无需改总线。 // // Clone 语义:broadcast::Sender 内部 Arc,clone 共享同一通道(多持有者 publish 到同一总线)。 // 对齐设计 §架构「模块订阅 + 发布,无相互 import」—— 各模块持 clone 的 EventBus 即可 pub/sub。 // ============================================================ /// AI 事件总线(pub-sub 域,载荷 `serde_json::Value`,L3 骨架)。 /// /// 基于 `tokio::sync::broadcast`,多订阅者多发布者共享同一通道。 /// - `subscribe()`:订阅,返回 [`broadcast::Receiver`] recv `serde_json::Value` /// - `publish(payload)`:发布,所有活跃订阅者收到(broadcast 语义) /// /// 载荷是 AiChatEvent 序列化后的 JSON(Phase3 D2 全 19 变体透传),subscriber 不解析业务语义。 /// /// 容量默认 256(`DEFAULT_BUS_CAPACITY`),可经 [`EventBus::with_capacity`] 自定义。 /// /// 兜底: /// - 无订阅者时 publish 不报错(broadcast 语义,事件丢弃) /// - 慢消费者(容量满)publish 返 0(静默丢老事件,不 panic) /// - `EVENT_BUS_ENABLED=false` 时 publish 静默丢弃(预留开关,接入期回退用) /// /// dead_code 说明:骨架阶段 emit 点未双写(publish 无调用方);标 allow 保留作批3 接入 /// (零调用方≠垃圾,总线核心结构,批3 接入即消费)。 #[allow(dead_code)] #[derive(Debug, Clone)] pub struct EventBus { sender: broadcast::Sender, } /// 事件订阅者(broadcast::Receiver,serde_json::Value 接收端)。 /// /// dead_code 说明:骨架阶段无外部消费方(批3 接入 tunnel subscriber),标 allow 保留。 #[allow(dead_code)] pub type EventSubscriber = broadcast::Receiver; // dead_code 说明(impl 块):骨架阶段 EventBus publish 无调用方(emit 点未双写) // (批3 接入 emit 点 + tunnel subscriber 后立即消费)。标 allow 覆盖 new/with_capacity/ // subscribe/publish/subscriber_count 全部关联项的 never-used 警告。 // 对齐零调用方原则:预留保留不盲删(批3 接入即消除)。 #[allow(dead_code)] impl EventBus { /// 创建默认容量(256)的事件总线。 pub fn new() -> Self { Self::with_capacity(DEFAULT_BUS_CAPACITY) } /// 创建指定容量的事件总线。 /// /// capacity 为 0 会 panic(broadcast::channel 要求 capacity ≥ 1)。 pub fn with_capacity(capacity: usize) -> Self { let (sender, _) = broadcast::channel(capacity); Self { sender } } /// 订阅事件总线,返回 Receiver。 /// /// 订阅后仅收到订阅时刻之后的 publish(历史事件不补发)。可多次 subscribe 获多个独立 Receiver, /// 每个 Receiver 各自维护消费进度(broadcast 多消费者语义)。 pub fn subscribe(&self) -> EventSubscriber { self.sender.subscribe() } /// 发布 AiChatEvent 序列化 JSON 到总线。 /// /// 返回值:成功送达的活跃订阅者数量(0 = 无订阅者或容量满丢老事件)。 /// `EVENT_BUS_ENABLED=false` 时静默丢弃返回 0(预留开关,骨架阶段未接入)。 /// /// 注:broadcast::send 是同步方法(非 async),与 df-workflow EventBus::send 的 async 签名 /// 不同 —— 本骨架 publish 同步返回更直接(broadcast::send 内部无 await 点),df-workflow 的 /// async 签名是为对齐 trait 抽象,本总线无此约束故同步。 pub fn publish(&self, payload: serde_json::Value) -> usize { if !EVENT_BUS_ENABLED { return 0; } // broadcast::send 返 Result: // Ok(n) = 送达 n 个活跃订阅者 // Err(SendError(_)) = 无订阅者(事件丢弃) // 两种情况都不 panic,Err 视为 0(对齐「无订阅者静默丢弃」兜底)。 self.sender.send(payload).unwrap_or(0) } /// 便捷发布:接收强类型 AiChatEvent,内部序列化为 Value 透传(D2 全 19 变体透传)。 /// /// 调用方传 AiChatEvent(业务侧已有强类型对象),本方法内部 `serde_json::to_value` /// 序列化为 Value 后调 [`publish`](Self::publish)。供 emit 点双写(publish + app.emit), /// 经 EventBus 透传给 tunnel subscriber(阶段2 接入)。 /// /// 性能优化(治空转):无订阅者时跳过序列化直接返 0。当前 tunnel subscriber 尚未接入, /// 20+ emit 点双写调本方法但无消费者,跳过序列化消除热路径开销。接入消费者后自动生效。 /// /// 序列化失败兑底返 0(丢弃,对齐 publish 静默丢弃语义;AiChatEvent 序列化稳定不触发,防御性)。 pub fn publish_event(&self, event: AiChatEvent) -> usize { // 无订阅者时跳过序列化(治空转开销)。broadcast::receiver_count 是同步原子读,代价极低。 // 有订阅者后才做 serde_json::to_value 序列化 + send。 if self.sender.receiver_count() == 0 { return 0; } match serde_json::to_value(&event) { Ok(v) => self.publish(v), Err(_) => 0, } } /// 当前活跃订阅者数量(broadcast::receiver_count)。 /// /// 诊断/监控用:接入期可观测订阅者是否到位。骨架阶段无实际调用,预留作可观测性扩展。 pub fn subscriber_count(&self) -> usize { self.sender.receiver_count() } } impl Default for EventBus { fn default() -> Self { Self::new() } } // ============================================================ // 单元测试 — pub/sub 基础语义(载荷 serde_json::Value) // // 覆盖: // 1. subscribe 后 publish,receiver 收到 Value(基础 pub-sub) // 2. 多订阅者都收到(broadcast 多消费者) // 3. 无订阅者时 publish 不 panic(兜底) // 4. publish 返回值 = 活跃订阅者数量 // 5. AiChatEvent 形态 Value 透传(含 type tag,模拟真实 AiChatEvent JSON) // // 不覆盖(留后续批次接入时): // - 慢消费者容量满丢老事件(需灌满容量场景,接入期补) // - 跨模块解耦实际效果(需接入真实 emit 点,批3) // ============================================================ #[cfg(test)] mod tests { use super::*; /// 基础 pub-sub:subscribe 后 publish,receiver 收到 Value 原样。 #[tokio::test] async fn test_basic_pub_sub() { let bus = EventBus::new(); let mut rx = bus.subscribe(); // 模拟 AiChatEvent::AiTextDelta 序列化 JSON(type tag discriminator) let payload = serde_json::json!({ "type": "text_delta", "conversation_id": "conv-1", "delta": "hello" }); let delivered = bus.publish(payload.clone()); assert_eq!(delivered, 1, "publish 应送达 1 个订阅者"); let received = rx.recv().await.expect("receiver 应收到事件"); assert_eq!(received, payload, "Value 应原样透传(零解析)"); assert_eq!(received["type"], "text_delta"); } /// 多订阅者都收到(broadcast 多消费者语义)。 #[tokio::test] async fn test_multiple_subscribers() { let bus = EventBus::new(); let mut rx1 = bus.subscribe(); let mut rx2 = bus.subscribe(); let payload = serde_json::json!({"type": "heartbeat", "conversation_id": null}); let delivered = bus.publish(payload.clone()); assert_eq!(delivered, 2, "publish 应送达 2 个订阅者"); // 两个 receiver 各自收到(独立消费进度) let r1 = rx1.recv().await.expect("rx1 应收到"); let r2 = rx2.recv().await.expect("rx2 应收到"); assert_eq!(r1["type"], "heartbeat"); assert_eq!(r2["type"], "heartbeat"); } /// 无订阅者时 publish 不 panic,返回 0(兜底)。 #[test] fn test_publish_no_subscribers() { let bus = EventBus::new(); // 无 subscribe 直接 publish let delivered = bus.publish(serde_json::json!({"type": "error", "error": "no one listening"})); assert_eq!(delivered, 0, "无订阅者 publish 应返回 0"); } /// publish 返回值 = 活跃订阅者数量(subscriber_count 对齐)。 #[test] fn test_delivered_count_matches_subscribers() { let bus = EventBus::new(); assert_eq!(bus.subscriber_count(), 0); let _rx1 = bus.subscribe(); assert_eq!(bus.subscriber_count(), 1); let _rx2 = bus.subscribe(); assert_eq!(bus.subscriber_count(), 2); // delivered 应等于 subscriber_count let delivered = bus.publish(serde_json::json!({"type": "agent_round", "round": 1})); assert_eq!(delivered, bus.subscriber_count()); } /// EventBus Clone 后共享同一通道(同一总线多持有者)。 #[tokio::test] async fn test_clone_shares_channel() { let bus = EventBus::new(); let mut rx = bus.subscribe(); let bus_clone = bus.clone(); // 经 clone 的实例 publish,原实例订阅的 receiver 也应收到(共享通道) bus_clone.publish(serde_json::json!({ "type": "completed", "total_tokens": 100, "prompt_tokens": 50, "completion_tokens": 50, "conversation_id": null })); let received = rx.recv().await.expect("clone 共享通道,receiver 应收到"); assert_eq!(received["type"], "completed"); } /// 全 19 变体形态 Value 均可透传(不限于子集,验证 D2 全透传不丢变体)。 #[test] fn test_full_variants_passthrough() { let bus = EventBus::new(); let mut rx = bus.subscribe(); // 覆盖早期 AiBusEvent 子集外的变体(ApprovalRequired/Compressed 等),证明全透传 let events = vec![ serde_json::json!({"type": "approval_required", "conversation_id": "c1"}), serde_json::json!({"type": "compressed", "conversation_id": "c1"}), serde_json::json!({"type": "dir_auth_required", "conversation_id": "c1"}), serde_json::json!({"type": "context_cleared", "conversation_id": "c1"}), ]; for event in &events { bus.publish(event.clone()); } for expected in &events { let received = rx.try_recv().expect("应收到每个变体"); assert_eq!(received, expected.clone(), "每个变体应原样透传(全透传不丢)"); } } }