Files
DevFlow/docs/02-架构设计/专项设计/全局事件数据总线-2026-06-21.md
绝尘 bd6a41fe6e 新增: 批次工作落地(推进链/评估闭环/事件总线/并发/加固) + 技术债清理 + 文档整理
后端:
- 工作流推进链(D-03):advance_task/状态机/闸门走 df-nodes Node trait,conditions 条件引擎扩展
- 想法评估闭环:启发式评分+对抗评估,df-ideas/scoring + df-storage/idea_eval_repo + idea 前端打通
- 全局事件数据总线:df-ai/context+context_helpers+augmentation 跨模块解耦
- AI planner/plan_hint/intent:aichat B 路线并行多轮基础
- patch_file 加固(TD-03/04):读改写整体锁防 lost update,expected_hash 合约闭环
- 压缩超时兜底(F-15 卡死根治)
- F-09 多会话并发:LlmConcurrency per-conv + streamingGuard 前端守护 + verify 脚本
- 知识注入 DRY/skills/audit 扩展

清理:
- aichat 技术债(误报 allow/死导入/过时注释 30 项)
- URGENT.md 删除(11 项加急全解决/迁 todo)
- 文档整理(todo/待决策/待审查/ARCHITECTURE/INDEX + 总线/技术债审查新文档)
2026-06-21 20:51:26 +08:00

9.4 KiB

全局事件数据总线设计

2026-06-21 · 专项设计 · 状态:构想定稿待评审 关联:cross-end-rust-backend 三层跨端 / devflow-product-positioning ai-working 可扩展 / F-260620-01 跨端小程序

背景与动机

ai-working(u-work)产品定位要求能力可插拔(编码 wedge → 扩展文档/数据/分析/办公自动化)+ 跨端(df-tunnel/relay/miniapp)。当前模块间直接函数调用/import 耦合,阻碍扩展:

  • 新能力接入需改现有模块(import 调用方)
  • 跨端模块(小程序/移动)无法直接调本地模块
  • 长时操作(压缩/工具执行)用同步 await,死等体验差(F-15 压缩卡死根因之一)

全局事件数据总线统一解决:模块经总线通信(pub-sub + request-reply + 流式),无直接耦合,响应式,跨端透传。

核心洞察(异步事件 + 同步等待 = 同步调用跨模块解耦):

  • 发请求事件(带 correlation_id)+ 等 reply 事件(oneshot await)→ 调用方语义同步,底层跨模块解耦(无 import 依赖)
  • 流式 reply(chunk 流)→ 响应式渐进(压缩/工具输出实时),根治死等

现状分析(碎片)

现有 范围 缺陷
df-workflow EventBus(crates/df-workflow/src/eventbus.rs,tokio broadcast) 工作流节点间 WorkflowEvent,非通用
Tauri emit/listen(@tauri-apps/api/event) 后端→前端单向 无 request-reply;前端→后端靠 IPC 命令(非事件)
前端专用事件(ai-drain-queue / ai-approval-clear-timers / ai-conversation-changed / ai-tool-slow-toast) 前端破环/联动 碎片化,无统一规范
AiChatEvent 通道(AiCompleted/AiError/AiCompressing/AiTextDelta...) 后端→前端 AI 流 单向,多监听器按 type 分发,无 reply
AR-11 df-data-changed 数据变更联动 雏形,未泛化(仅 AI 工具触发)

缺口:无统一总线 / 无 request-reply(同步语义跨模块)/ 无流式 reply / 无跨端透传。

概念

全局事件数据总线(Unified Event Bus):贯穿后端 + 前端 + 跨端的单一事件通道,支持三类语义:

1. pub-sub(发布订阅)

数据变更/状态广播。发布者 publish(event),订阅者 subscribe(filter)。泛化 AR-11(任务/知识/项目/灵感 增删改统一)。

2. request-reply(请求-回复)— 同步调用跨模块解耦

  • 请求方:reply = bus.request(RequestEvent{payload, deadline}).await(oneshot 等 reply)
  • 响应方:订阅 RequestEvent 类型,处理后 bus.reply(reply_to, ResponseEvent)
  • 语义同步(await),底层异步事件,跨模块无 import
  • 超时/取消:request 带 deadline,reply oneshot 超时 → Err(软超时,非死等)
  • 适用:压缩请求 / 工具执行 / 审批(用户决策 reply)

3. 流式 reply(响应式渐进)— 根治死等

  • 请求方:stream = bus.request_stream(RequestEvent).await; while let Some(chunk) = stream.next().await
  • 响应方:reply_tx 发多个 chunk(渐进),结束发 Done
  • chunk 持续 = 数据活性,流断即知(非 60s 死等);无数据 N 秒判 hang(软超时)
  • 适用:压缩流式(LLM chunk)/ 工具执行流式(命令输出实时)/ AI 对话流(AiTextDelta 雏形)

设计

架构

┌─────────────── 后端(tokio)───────────────┐
│  UnifiedEventBus                            │
│   ├─ broadcast<DomainEvent>(pub-sub)        │
│   ├─ 类型路由 HashMap<ReqType, mpsc::Sender<Request>>(request-reply)│
│   ├─ oneshot reply + mpsc stream chunk      │
│   └─ df-tunnel adapter(跨端桥接)           │
│                                              │
│  模块(AI/任务/知识/工作流/工具/新能力)       │
│   订阅 + 发布 + reply,无相互 import         │
└──────────────────────────────────────────────┘
            ▲ Tauri event bridge(后端↔前端)
            ▼
┌─────────────── 前端(Vue3)─────────────────┐
│  UnifiedEventBus 镜像(Tauri listen 转发)    │
│   composables 订阅/发布,store 响应式更新    │
└──────────────────────────────────────────────┘
            ▲ df-tunnel/relay 透传
            ▼
┌─────────────── 跨端(小程序/移动)──────────┐
│  UnifiedEventBus 远程端(relay 转发)         │
│   远程模块订阅/请求,经 relay 透传到本地     │
└──────────────────────────────────────────────┘

事件类型

  • DomainEvent(pub-sub):
    • EntityChanged{entity, id, op: Create/Update/Delete} — 泛化 AR-11(任务/知识/项目/灵感/工作流 CRUD 统一)
    • StateChanged{entity, id, from, to} — 任务状态机/审批态/会话态
    • Lifecycle{kind, id} — 会话开关/Provider 变更
  • RequestEvent(request-reply):{type, payload, reply_to: oneshot::Sender, deadline}
    • CompressRequest / ToolExecRequest / ApprovalRequest / KnowledgeExtractRequest
  • StreamChunk(流式):{stream_id, chunk: String/Bytes, done: bool} — LLM/工具输出渐进

request-reply 实现雏形(Rust)

struct UnifiedEventBus {
    domain_tx: broadcast::Sender<DomainEvent>,
    req_handlers: HashMap<ReqType, mpsc::Sender<Request>>,
}
impl UnifiedEventBus {
    pub fn publish(&self, e: DomainEvent) { let _ = self.domain_tx.send(e); }
    pub fn subscribe(&self) -> broadcast::Receiver<DomainEvent> { self.domain_tx.subscribe() }

    // request-reply(同步语义,跨模块解耦)
    pub async fn request<R: Request>(&self, req: R) -> Result<R::Reply> {
        let (tx, rx) = oneshot::channel();
        self.dispatch(req.into(tx)).await?;
        tokio::time::timeout(req.deadline(), rx).await??
    }
    // 流式 reply(响应式渐进)
    pub async fn request_stream(&self, req: R) -> Result<mpsc::Receiver<StreamChunk>> { ... }
    // 响应方注册 handler
    pub async fn handle<R: Request>(&self, handler: impl Handler<R>) { ... }
}

跨端透传(df-tunnel/relay)

  • 本地 bus adapter:DomainEvent/Request 经 df-tunnel 序列化透传到远程端
  • 远程端(小程序)bus:订阅本地事件 + 发请求(经 relay 回本地处理)
  • 统一抽象:本地/远程模块代码一致(经 bus,不知对端远近)— 契合 ai-working 跨端定位

应用场景

  1. 响应式压缩(替 F-15 超时死等):AI 发 CompressRequest → 压缩 handler 订阅 + provider.stream 流式 reply chunk → AI request_stream().await 渐进收 + 前端订阅 chunk 显进度 → chunk 持续(非死等),无数据软超时(非 60s 死等)
  2. 工具执行经总线:AI 发 ToolExecRequest → 工具 handler + reply;path_auth 未授权 → ApprovalRequest(总线 request-reply,用户审批 reply)
  3. 模块可插拔(ai-working):新能力(文档/数据/分析)订阅总线 Request,不改 AI/任务模块(零侵入扩展)
  4. 数据联动泛化:AR-11 df-data-changed → DomainEvent::EntityChanged 泛化(任务/知识/项目/灵感/工作流 CRUD 统一,前端 store 订阅刷新)
  5. 跨端 AI Chat:小程序经总线订阅本地 AiTextDelta 流 + 发消息请求(F-260620-01)

迁移路径(分阶段,每阶段独立可发布/可回退,配 feature flag)

  1. 阶段1 后端 UnifiedEventBus 基础:broadcast + request-reply(oneshot)+ 流式(mpsc)。df-workflow EventBus 包装为通用 DomainEvent。单测。
  2. 阶段2 首个应用:响应式压缩:CompressRequest + 流式 reply(替 compress_via_llm await 死等,撤销 F-15 60s 硬 timeout)。验证 request-reply/流式可行。
  3. 阶段3 前端镜像:Tauri bridge(后端 DomainEvent → 前端总线,前端订阅)。ai-drain-queue/ai-conversation-changed/ai-approval-clear-timers 等碎片迁移统一。
  4. 阶段4 数据联动泛化:AR-11 → DomainEvent::EntityChanged(任务/知识/项目/灵感 CRUD 统一,前端 store 订阅刷新)。
  5. 阶段5 工具/审批经总线:ToolExecRequest / ApprovalRequest(path_auth/risk 审批统一为 bus request-reply)。
  6. 阶段6 跨端透传:df-tunnel adapter + 小程序远程端(F-260620-01)。

与现有关系

  • cross-end-rust-backend 三层(df-tunnel/relay/miniapp):总线跨端透传基础
  • devflow-product-positioning ai-working 可扩展:模块可插拔经总线(编码特化 → 通用工作能力)
  • df-workflow EventBus:升级/包装为通用 bus 一部分(不废弃,复用 broadcast 基础)
  • AR-11 df-data-changed:泛化为 DomainEvent::EntityChanged
  • F-260620-01 跨端小程序:总线远程端首个跨端应用
  • F-15 上下文压缩:阶段2 响应式压缩经总线,替超时死等

待决策

  1. bus crate 位置:新建 df-bus?df-workflow eventbus 升级为通用?df-ai-core?
  2. 事件 schema:强类型 enum(编译期安全,扩展需改 enum)vs 动态 JSON(灵活,弱类型)→ 倾向强类型 + 开放式(未知类型透传)
  3. request-reply 路由:类型路由 HashMap<ReqType, Sender>(强类型)vs topic 字符串(灵活)
  4. 跨端序列化:serde JSON(可读,调试友好)vs bincode(紧凑,性能)
  5. 现有事件迁移节奏:AiChatEvent/前端专用事件一次性迁移 vs 渐进(新功能用 bus,旧的逐步迁)
  6. 背压/限流:broadcast 慢消费者丢弃 vs mpsc 背压(请求-回复)