经代码核验发现 9 份设计文档的状态标注严重滞后(标'待实施'但 实际已完整落地),本次批量同步: 已落地(核验确认): - 局部编辑工具:三层防御+三模式完整,仅文本不支持二进制 - 密钥迁移健壮性:空 key 不覆盖+即时迁移补密钥+阻断保存 - AST 符号解析:符号读取工具已注册+基线测试守护 - 查询能力补全:任务/项目/灵感均多维动态查询 - 条件表达式引擎:手写求值器+JSON Path+执行器集成+前端入口 - 工作流脚本边界:命令白名单/黑名单+危险关键词告警 - 消息拆分存储:消息表+全量迁移+读写全部切换 - 消息级溯源:消息 ID+四场景溯源+切读全部完成 部分落地: - 全局事件总线:基建+20 余个发射点就位,消费者未接(空转) 归档不实施: - 规格契约自检:核心价值已被求助协议+自审闸门覆盖,过度设计
149 lines
9.5 KiB
Markdown
149 lines
9.5 KiB
Markdown
# 全局事件数据总线设计
|
|
|
|
> 2026-06-21 · 专项设计 · 状态:🟡 部分落地(2026-06-28 核验:基建(AI/工作流两路总线)+ 20+ emit 点双写就位,但真实订阅方/前端镜像/tunnel 透传均未接入空转。阶段 2-6 待真实痛点驱动)
|
|
> 关联:[[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)
|
|
```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 背压(请求-回复)
|