diff --git a/docs/architecture/message-lifecycle-and-ai-flow.md b/docs/architecture/message-lifecycle-and-ai-flow.md index f937f38f..b0c1bd92 100644 --- a/docs/architecture/message-lifecycle-and-ai-flow.md +++ b/docs/architecture/message-lifecycle-and-ai-flow.md @@ -137,6 +137,38 @@ ActionCableListener.OnEvent (internal/wsevent/bridge_listener.go:51) └── 前端 Vue SPA WebSocket 连接接收 → Vuex/Pinia store 更新 → UI 刷新 ``` +#### Durable realtime 投递语义 + +配置 WorkerPool 后,`EventPublisher` 先把 account 与 pubsub token 目标分别写入 +`background_jobs`,再由 `realtime:event_publish` job 发布到 Redis、Hub 和 SSE。 +幂等键避免同一条 `message.created`/目标重复入队,job claim 避免并发执行同一条 +记录;它们不覆盖发布副作用与 job 完成状态之间的崩溃窗口: + +1. job 已完成 Redis 发布,并可能已完成本地 SSE 投递; +2. 进程在将 `background_jobs.status` 更新为 `completed` 前退出; +3. stale-job 恢复将该记录重新置为可执行并再次发布。 + +因此,对外契约是 **durable at-least-once publication**,不是 exactly-once +delivery。进程恢复后,在线 WebSocket/SSE 消费者可能再次收到同一事件;account +与 pubsub token 是独立 job,重试和到达顺序也彼此独立。Redis Pub/Sub、Hub 和 +SSE 不保存离线或慢消费者的确认状态,所以该契约保证发布尝试可恢复,不保证每个 +客户端至少接收一次。 + +消费者必须把 realtime 事件作为可重放通知处理: + +- `message.created` 使用订阅作用域、事件名和 payload 的稳定 message `id` 去重或 + upsert;不得使用未出现在 wire payload 中的 background job ID。 +- `*.updated` 按资源 `id` 幂等 upsert,并在 payload 提供 `updated_at`/版本时拒绝 + 旧更新;不要仅按资源 ID 永久丢弃后续合法更新。 +- 重连或发现序列缺口时,以 REST API 返回的持久化资源为准。会触发声音、通知、 + SDK callback 或其他非幂等副作用的外部消费者必须在执行副作用前自行去重。 + +当前风险可接受:复用的 dashboard 会按 message ID 替换重复消息,并按 +conversation `updated_at` 忽略旧更新;widget 消息仓库也按 message ID upsert。 +这与上游 Chatwoot 的异步 `ActionCableBroadcastJob` 行为一致。只有当业务要求 +跨进程崩溃的 exactly-once 副作用时,才应另行引入稳定 transport event ID 与 +消费端 inbox/ack;当前协议不承诺该能力。 + ### 2.5 通知创建 ```