补充DirectProject事件订阅规范

补充Thread Manager全局事件队列、subscriber cursor与consume契约

补充持久化顺序、队头回收、过期重订阅和历史锚点规则

新增Thread Manager里程碑规范与小步实施计划
This commit is contained in:
2026-09-15 16:50:14 +08:00
parent 3353906e6f
commit 9d2ca40369
3 changed files with 154 additions and 0 deletions
@@ -0,0 +1,37 @@
# 【实施计划】DirectProject Thread Manager 事件订阅
| 字段 | 值 |
| --- | --- |
| Milestone | `docs/project-memory/plans/【里程碑】DirectProject Thread Manager事件订阅-2026-09-15.md` |
| Status | ready |
| Owner | Codex |
## 修改边界
- 允许修改:DirectProject Rust Thread Manager 深模块、app-server 事件适配、Tauri command/event 桥接、DirectProject 前端订阅/reducer/历史加载、对应测试和主规范。
- 明确不修改:SpacetimeDB、HTTP API、非 DirectProject Runtime、Codex app-server durable thread、用户可见 JSONL 细节。
- 保持已有 `.env` 未提交修改,不触碰个人配置。
## 实现顺序
1. 先新增独立 Rust queue/subscriber 深模块,只承载事件追加、逻辑队头回收、subscriber cursor 锁和纯单测。
2. 将 app-server 公开事件安全标准化后接入 Thread Manager;在 item 完成持久化成功后追加完成事件,并追加 turn 生命周期事件。
3. 增加 Tauri `subscribe/consume/readHistory` 命令与 notify 事件,固定错误和 bootstrap 原子边界。
4. 前端改为 subscriptionId 驱动的 raw event reducer;重进/过期时先 bootstrap,完成后原子替换;历史按 itemId 懒加载。
5. 移除 DirectProject legacy conversation 读取分支,补齐契约、并发、恢复和失败关闭测试。
6. 每个独立切片分别运行定向验证并形成中文小提交;里程碑验收后再清理临时计划。
## 验证命令
1. `cargo test --manifest-path apps/ai-game-creator-shell/src-tauri/Cargo.toml` 的 DirectProject/Thread Manager 定向测试。
2. 相关前端 Vitest 与类型检查。
3. `npm run check:encoding`
4. `npm run check:doc-index`
5. `git diff --check`
## 风险与回滚点
- 现有 app-server 事件模型与公开 raw event envelope 不完全一致:先在适配层收口,不让协议细节泄漏到前端。
- 单 Vec 队列不能中间删除;unfinished item 长时间不结束可能暂时 pin 住队头,必须保留可观测上限和测试。
- Tauri command 无传输层断开回调,subscriber 只通过 queue eviction 失效;测试不能依赖 unsubscribe 或连接断开清理。
- legacy 删除属于 breaking history 行为;失败关闭测试必须确认不会 fallback 或迁移。
@@ -0,0 +1,52 @@
# 【里程碑】DirectProject Thread Manager 事件订阅
| 字段 | 值 |
| --- | --- |
| Version | 1.0 |
| Status | ready |
| Date | 2026-09-15 |
| Parent Spec | `docs/technical/【技术方案】DirectProject Codex原始历史与异常恢复-2026-09-04.md` |
## 目标
让 DirectProject 对话在页面离开、重进和短暂断线后仍能由前端重建运行态;运行态事件由 Tauri 进程级 Thread Manager 管理,已完成 item 继续以 `project.jsonl` 为持久化事实源。
## 范围
- 每 thread 一个全局 seq 和 append-only replay queue。
- 每 subscriber 独立的 Rust 内部 cursor、并发安全消费和 notify 唤醒。
- `subscribe` bootstrap、`consume``SUBSCRIPTION_EXPIRED` 和历史 item 锚点。
- app-server 公开事件的安全标准化、item 持久化先于完成事件转发。
- 前端 raw event reducer、历史懒加载和过期重订阅。
- 删除本链路 legacy conversation 格式支持,不提供 fallback 或 migration。
## 不在范围内
- SpacetimeDB、HTTP API、Codex thread durable recovery。
- 新的 item durable/status/pendingInteraction 字段或持久化确认事件。
- 前端访问 JSONL 路径、格式或持久化细节。
- 多 active turn;同一 thread 仍只有一个 active turn。
## 依赖与前置条件
- DirectProject 现有 app-server 事件解析和 `project.jsonl` 读写。
- 当前 Tauri command/event 注册入口。
- 现有前端 DirectProject 聊天 reducer 与历史加载入口。
## 验收标准
- [ ] 页面离开后 app-server 回合继续,重进页面能通过 subscribe 重建 unfinished item。
- [ ] 同一 thread 的多个 subscriber 各自消费,不互相覆盖或重复推进 cursor。
- [ ] `consume` 返回 cursor 之后的全局 raw events,通知不携带 payload。
- [ ] queue eviction 只清理队头;落后 subscriber 得到 `SUBSCRIPTION_EXPIRED` 并可重新 subscribe。
- [ ] item 完成先持久化,成功后才进入完成事件队列;失败不发送正常完成事件。
- [ ] `turn.completed` 由 app-server 终态进入 raw queue,前端据此结束运行态。
- [ ] subscribe 返回 item 历史锚点而不是完整 history;前端可按 itemId 懒加载。
- [ ] 未完成 item 的每个 delta 可从 `item.started` 开始重放;不截断 active item。
- [ ] legacy conversation 行直接失败关闭,无 fallback、无迁移。
## 证据要求
- 自动化:queue/cursor/eviction 并发单测、事件标准化和持久化顺序测试、Tauri command 测试、前端 reducer 与重订阅测试。
- 运行时:关闭/切页后重进 DirectProject;并发 item;短暂断线 consumecursor 过期重订阅。
- 边界:持久化失败、未知 subscription、queue 超限、多个 subscriber、turn 无 item 间隙、legacy 行拒绝。
@@ -62,3 +62,68 @@ DirectProject 的浏览器层只负责显示和乐观状态,不再调用通用
写入使用 `write_all + flush`。读取时允许丢弃文件末尾一条不完整 JSON 行;非 `response_item` 行和无法投影的 item 直接失败,不做数据迁移或 fallback。
该失败有专门恢复提示,并按不可重试处理:同一份历史文件每次读都会得到同一结论,重试不会改变结果,因此不会向用户显示「可直接重试」。
## Thread Manager 运行态事件订阅
DirectProject 的页面不是回合执行的所有者。Tauri 进程内的 Thread Manager 按 thread 维护运行态事件,并允许同一 thread 存在多个独立 subscriber。事件队列只服务运行期间和短期断线恢复,不替代 `project.jsonl` 历史事实源。
### 公开契约
概念接口如下:
```ts
subscribe(threadId) -> {
subscriptionId,
lastCompletedItemId: string | null,
events: RawEvent[],
}
consume(subscriptionId) -> {
events: RawEvent[],
}
notify -> { subscriptionId }
readHistory(threadId, { beforeItemId?, limit }) -> {
items: CompletedItem[],
hasMore: boolean,
}
```
`subscribe` 不返回完整历史。`lastCompletedItemId` 只是历史读取锚点,前端自行按 item ID 懒加载需要的历史切片。`events` 是当前运行态重建所需的未完成 item 原始事件,以及当前 turn 的生命周期锚点;前端用同一个 reducer 重放 bootstrap 和后续事件。Rust 不保存或理解前端 reducer state。
`consume` 不接收或返回 cursor。每个 subscriber 在 Rust 内部持有自己的 cursor,并在加锁的临界区内完成过期判断、读取和 cursor 前进。前端只持有 `subscriptionId` 与 reducer state。并发 `consume` 不重复返回同一批事件。
`notify` 只负责唤醒,不携带事件、cursor 或持久化状态。前端收到通知后调用 `consume`;通知可合并、重复或丢失,事件完整性由 `consume` 保证。
### 事件和顺序
Thread 内所有公开事件共用一个单调递增 seq;seq 允许跳号,前端不要求连续。事件 envelope 至少包含:
```ts
{
seq: number,
type: string,
turnId: string,
itemId?: string,
payload: unknown,
}
```
进入 Thread Manager 的是已经完成安全过滤和协议标准化的公开 raw event,不是未经审查的 app-server JSON。事件可交错包含多个并发 item:`item.started``item.delta``item.completed`、approval/request/resolved 事件,以及 `turn.started``turn.completed` 生命周期事件。前端按 `turnId` / `itemId` 分发并 reduce,不需要 item 级 cursor 或第二套 reducer。
一个 thread 同时最多有一个 active turn;一个 turn 内允许多个并发 item。`turn.completed` 必须在该 turn 的完成 item 均成功持久化后进入队列,前端据此结束运行态;不能用“不存在 unfinished item”猜测 turn 是否完成。
### 队列、subscriber 和回收
每个 thread 一个 Vec-based append-only replay queue,使用逻辑 head 偏移清理前缀,不做中间删除。完成 item 的事件在持久化成功后才可进入普通 replay 回收流程;unfinished item 的事件必须保留到 item 完成,不能被普通上限截断。
队列有内部最大事件数和最大序列化字节数。超限时先标记长期落后的 subscriber 为 expired,并将其移出有效 subscriber 的最小 cursor 计算;随后只能清理队头连续、已无有效 subscriber 需要且所属 item 已持久化的事件。没有 subscriber 时,已持久化完成 item 的事件副本可以直接清理。未完成 item 的事件仍保留。
subscriber 不依赖 `unsubscribe` 或传输层断开清理。每次 `subscribe` 都创建新的独立 subscription;同一 thread 的其它 subscriber 不受影响。旧 subscription 只有在 queue eviction 后才失效,调用 `consume` 返回统一错误 `SUBSCRIPTION_EXPIRED`。前端保留旧 reducer state,重新 subscribe 完成 bootstrap 后再原子替换。
### Bootstrap 原子性和恢复
`subscribe` 必须在同一个 Thread Manager 边界注册 subscriber、捕获 queue 尾部、确定历史锚点和当前运行态事件;bootstrap 期间产生的新事件由该 subscriber 的内部 cursor 继续通过 `consume` 获取,不能丢失。
断线恢复优先调用 `consume(subscriptionId)`。subscription 仍有效时只返回该 subscriber 尚未消费的 queue 事件;subscription 已过期或 Thread Manager 重启后统一走新的 `subscribe`,再由前端按 `lastCompletedItemId` 从历史懒加载。Rust 不提供 `getItemSnapshot(itemId)`,已完成 item 始终通过历史读取。