diff --git a/docs/serve/consumer-capability-checklist.md b/docs/serve/consumer-capability-checklist.md new file mode 100644 index 000000000..114e107e1 --- /dev/null +++ b/docs/serve/consumer-capability-checklist.md @@ -0,0 +1,326 @@ +# 新版消费端事件消费能力 Checklist + +> 消费端指:app-ai-native、cs-cloud Web UI、VS Code 插件等通过 SSE/REST 消费 csc serve 的客户端 +> 前提:csc serve 已完成事件标准化整改,输出 opencode canonical 格式 +> 关联:`docs/serve/serve-event-standardization-proposal.md`、`docs/serve/serve-event-standardization-todo.md` +> +> 状态标记说明: +> - ✅ app-ai-native 已实现 +> - ❌ app-ai-native 未处理(需新增实现) +> - ⚠️ 部分实现 / 存在差距 + +--- + +## 1. SSE 事件消费 + +### 1.1 消息生命周期 + +| 事件 | 数据结构 | 消费端应实现 | app-ai-native | +|---|---|---|---| +| `message.updated` | `{ sessionID, info: { id, role, modelID, providerID, cost, tokens, time, parentID, finish, error } }` | 创建/更新消息对象 | ✅ `device-session.tsx` → `useMessageUpdater` | +| `message.part.updated` | `{ sessionID, part: Part }` | 创建/更新 part,渲染对应 UI | ✅ `device-session.tsx` → part dispatch by type | +| `message.part.delta` | `{ sessionID, messageID, partID, field, delta }` | 追加 delta 到指定 part 字段(流式渲染) | ✅ `device-session.tsx` → `handlePartDelta` | +| `message.removed` | `{ sessionID, messageID }` | 移除指定消息(tombstone 场景) | ❌ 无 switch case 处理 | +| `message.attachment` | `{ sessionID, attachmentType, attachment }` | 展示附加信息(hook 结果、记忆、诊断等) | ❌ 无 switch case 处理 | + +### 1.2 Part 类型渲染 + +| Part 类型 | 渲染要求 | app-ai-native | +|---|---|---| +| `text` | 流式文本渲染(打字机效果),field `text` 的 delta 追加 | ✅ `createPacedValue` + Markdown 实时渲染 | +| `reasoning` | 可折叠的思考过程块,`redacted: true` 时显示「思考内容已隐藏」 | ✅ ReasoningPart 组件,默认折叠 | +| `tool` (pending) | 工具调用已开始,显示工具名 + 等待状态(spinner) | ✅ ToolPart 状态机 `pending` | +| `tool` (running) | 工具正在执行,显示工具名 + 输入参数摘要 + 运行中状态 | ✅ ToolPart 状态机 `running` | +| `tool` (completed) | 工具执行完成,显示标题 + 输出内容(可折叠) + 执行耗时 | ✅ ToolPart 状态机 `completed` | +| `tool` (error) | 工具执行失败,显示错误信息(红色标记) | ✅ ToolPart 状态机 `error` | +| `step-start` | 标记新一轮 LLM 调用开始(内部标记,不一定有独立 UI) | ✅ StepStartPart 组件 | +| `step-finish` | 标记一轮 LLM 调用结束,更新 step 级别的 cost/tokens/reason | ✅ StepFinishPart 组件 | +| `compaction` | 显示上下文压缩指示器(「对话已被压缩」) | ✅ CompactionPart 组件 | +| `subtask` | 显示子任务信息(prompt、agent、description) | ⚠️ 无独立 SubtaskPart 组件,可能以 tool part 显示 | + +### 1.3 Task 生命周期 + +| 事件 | 消费端应实现 | app-ai-native | +|---|---|---| +| `task.started` | 在任务面板创建任务条目,显示 description、taskType | ❌ 无 `task.*` 事件处理;通过 `todo.updated` + ToolPart 状态推断 | +| `task.progress` | 更新任务进度:description 变更、usage 累积、workflow 进度条 | ❌ 同上 | +| `task.completed` | 标记任务终态(completed/failed/stopped),显示 summary + usage | ❌ 同上 | + +> **注**: app-ai-native 目前通过 ToolPart (`tool === "task"`) 的状态转换 + `todo.updated` 来追踪任务进度,没有独立的 `task.*` 事件消费。新增 `task.*` 事件处理需要新建 TaskState Map 和相关 UI。 + +消费端应维护任务状态 Map: + +```typescript +interface TaskState { + taskID: string + status: 'running' | 'completed' | 'failed' | 'stopped' + description: string + taskType?: string + summary?: string + usage?: { total_tokens: number, tool_uses: number, duration_ms: number } + startTime: number + endTime?: number +} +``` + +### 1.4 Session 状态 + +| 事件 | 消费端应实现 | app-ai-native | +|---|---|---| +| `session.created` | 创建 session 对象 | ✅ `device-workspace.tsx` → `session.created` handler | +| `session.updated` | 更新 session 元信息(model, provider, status) | ✅ `device-workspace.tsx` → `session.updated` handler | +| `session.deleted` | 移除 session 对象 | ✅ `device-workspace.tsx` → `session.deleted` handler | +| `session.status` | 更新 busy/idle 状态指示器(spinner 切换) | ✅ `device-session.tsx` → `session.status` handler | +| `session.error` | 显示错误 banner,含重试倒计时(retryInMs) | ❌ SDK 中有类型定义但 app-ai-native 无 switch case 处理 | +| `session.warning` | 显示警告 banner(cache_warning 等) | ❌ 无处理 | +| `session.info` | 显示信息提示(informational system 消息) | ❌ 无处理 | +| `session.metrics` | 更新性能指标面板(turn_duration, TTFT 等) | ❌ 无处理 | +| `session.hook_summary` | 显示 Hook 执行汇总 | ❌ 无处理 | +| `session.diff` | 更新文件变更预览 | ✅ `device-session.tsx` → `session.diff` handler | + +### 1.5 权限与问答 + +| 事件 | 消费端应实现 | app-ai-native | +|---|---|---| +| `permission.asked` | 弹出权限确认对话框(工具名 + patterns + metadata) | ✅ `device-session.tsx` → `permission.asked` handler | +| `permission.replied` | 关闭对应权限对话框 | ✅ `device-session.tsx` → `permission.replied` handler | +| `question.asked` | 弹出问题对话框(header + options + multiple + custom) | ✅ `device-session.tsx` → `question.asked` handler | +| `question.replied` | 关闭对应问题对话框 | ✅ `device-session.tsx` → `question.replied` handler | +| `question.rejected` | 关闭对应问题对话框(取消状态) | ✅ `device-session.tsx` → `question.rejected` handler | + +### 1.6 工具进度 + +| 事件 | 消费端应实现 | app-ai-native | +|---|---|---| +| `tool.progress` | 关联到对应 tool part,显示实时进度(Bash 输出、Agent 进度等) | ❌ 无 `tool.progress` 事件处理 | + +### 1.7 基础设施事件 + +| 事件 | 消费端应实现 | app-ai-native | +|---|---|---| +| `server.connected` | SSE 连接建立,标记连接状态为 online | ⚠️ 通过 fetch 响应头判断连接状态,无独立事件 | +| `server.heartbeat` | 更新心跳时间戳,检测连接存活 | ⚠ SSE 注释行(`: heartbeat`)用作 keep-alive,无事件处理 | + +--- + +## 2. REST API 消费 + +### 2.1 消息历史 + +| 端点 | 消费端应实现 | app-ai-native | +|---|---|---| +| `GET /session/{id}/message` | 加载消息历史,parts-based 格式解析 | ✅ `device-client.ts` → `session.message.list()` + parts 解析 | +| `GET /session/{id}/todo` | 加载 TODO 列表,结构化展示 | ✅ `todo.updated` 事件驱动 + SDK `todo` 类型 | +| `GET /session/{id}/diff` | 加载文件 diff,代码变更预览 | ✅ `session.diff` 事件 + diff 渲染组件 | + +消息历史返回格式应为: + +```typescript +interface MessageResponse { + id: string + role: 'user' | 'assistant' + parts: Part[] + time: { created: number, completed?: number } + cost?: number + tokens?: { input: number, output: number, reasoning: number, cache: { read: number, write: number } } + modelID?: string + providerID?: string + parentID?: string + finish?: string + error?: { name: string, message: string } +} +``` + +### 2.2 Session 管理 + +| 端点 | 消费端应实现 | app-ai-native | +|---|---|---| +| `GET /session` | 列出所有 session,显示标题/时间/模型 | ✅ `device-client.ts` → `session.list()` | +| `POST /session` | 创建新 session | ✅ `device-client.ts` → `session.create()` | +| `GET /session/{id}` | 获取 session 详情 | ✅ `device-client.ts` → `session.get()` | +| `PATCH /session/{id}` | 更新 session(标题、tag 等) | ✅ `device-client.ts` → `session.update()` | +| `DELETE /session/{id}` | 删除 session | ✅ `device-client.ts` → `session.delete()` | +| `POST /session/{id}/prompt` | 发送 prompt | ✅ `device-client.ts` → `session.chat()` | +| `POST /session/{id}/abort` | 中止当前 turn | ✅ `device-client.ts` → `session.abort()` | +| `GET /session/status` | 获取所有 session 状态 | ✅ 通过 SSE `session.status` 事件 | + +### 2.3 Provider / Model + +| 端点 | 消费端应实现 | app-ai-native | +|---|---|---| +| `GET /provider/capabilities` | 获取可用模型列表及能力 | ⚠️ 通过 config/model API 获取,非标准端点 | + +--- + +## 3. 流式渲染能力 + +### 3.1 文本流式渲染 + +- [x] ✅ 实现 `createPacedValue` 或等价的打字机效果 +- [x] ✅ `message.part.delta { field: "text" }` → 追加到 text part 的 text 字段 +- [x] ✅ 支持 Markdown 实时渲染(文本未完成时不关闭代码块等) + +### 3.2 Reasoning 流式渲染 + +- [x] ✅ `message.part.delta { field: "text" }` → 追加到 reasoning part 的 text 字段 +- [x] ✅ reasoning part 默认折叠,点击展开 +- [x] ✅ `redacted: true` 时显示固定文案而非内容 + +### 3.3 Tool Input 流式渲染 + +- [x] ✅ `message.part.delta { field: "input" }` → 累积 tool input JSON +- [x] ✅ 实时解析部分 JSON 展示关键参数(如 file_path、command) + +### 3.4 流式场景的滚动行为 + +- [x] ✅ 自动滚动到底部(用户未手动上翻时) +- [x] ✅ 用户上翻时暂停自动滚动,新消息提示后恢复 + +--- + +## 4. 工具状态机渲染 + +### 4.1 状态转换 + +``` +pending → running → completed + → error +``` + +- [x] ✅ `pending`:显示工具名 + spinner + 输入参数预览 +- [x] ✅ `running`:显示工具名 + 执行中标记 + 已运行时间 +- [x] ✅ `completed`:显示标题 + 可折叠的输出内容 + 执行耗时 + cost +- [x] ✅ `error`:显示错误信息(红色)+ 可折叠的错误详情 + +### 4.2 特定工具的富渲染 + +| 工具 | 富渲染要求 | app-ai-native | +|---|---|---| +| `bash` / `powershell` | 命令行 + 输出(终端风格),tool.progress 实时追加输出 | ✅ BashToolRenderer(终端风格输出),❌ 无 `tool.progress` 实时追加 | +| `read` / `glob` / `grep` | 文件路径 + 匹配行数,可点击跳转 | ✅ ReadToolRenderer / GlobToolRenderer / GrepToolRenderer | +| `edit` / `fileedittool` | diff 预览(红色删除 / 绿色新增),含 filediff metadata | ✅ EditToolRenderer(diff 视图) | +| `write` | 新建文件标记 + 内容预览 | ✅ WriteToolRenderer | +| `agent` / `task` | 子任务进度指示器,关联 task.started/task.completed | ⚠️ ToolPart 状态机处理,无独立 task 事件关联 | +| `webfetch` / `websearch` | URL + 搜索结果摘要 | ✅ WebFetchToolRenderer | +| `ask_user_question` | 已由 question.asked 处理,tool part 显示为「等待用户回复」 | ✅ 由 question.asked 事件驱动 | + +--- + +## 5. 任务面板 + +### 5.1 数据模型 + +```typescript +interface TaskPanelState { + tasks: Map + activeTaskCount: number + totalCost: number + totalTokens: number +} +``` + +### 5.2 UI 要求 + +- ❌ 任务列表视图(运行中 / 已完成分组)— 当前通过 `todo.updated` + ToolPart 间接展示 +- ❌ 单任务详情展开(description、summary、usage)— 需新增 +- ⚠️ 后台任务计数 badge — 通过 ToolPart `tool === "task"` 部分推断 +- ❌ 任务完成通知(toast / desktop notification)— 需新增 +- ❌ task.progress 的 workflow 进度条(phase 级别,如有的话)— 需新增 + +> **总结**: app-ai-native 没有独立的任务面板。任务进度通过 ToolPart 状态机 + `todo.updated` 间接追踪。如需富任务 UI,需要新增 `task.*` 事件消费 + TaskPanelState 管理。 + +--- + +## 6. 错误处理与重试 + +### 6.1 错误展示 + +| 场景 | 展示方式 | app-ai-native | +|---|---|---| +| `session.error` (api_error) | 红色 banner,含错误消息 + 重试倒计时 | ❌ SDK 有类型但无 UI 处理 | +| `session.error` (api_retry) | 黄色 banner,「正在重试 (N/M)...」 | ❌ 同上 | +| tool part (error) | 工具卡片内红色错误信息 | ✅ ToolPart error 状态渲染 | +| message.updated (error 字段) | 消息级别的错误标记 | ⚠️ 部分处理,依赖 message info 结构 | +| result (subtype: error_max_turns) | 「已达到最大轮次限制」提示 | ❌ 无独立处理 | +| result (subtype: error_max_budget) | 「已达到预算上限」提示 | ❌ 无独立处理 | + +### 6.2 重试 UI + +- ❌ 倒计时显示(来自 `session.error.retryInMs`)— 需新增 +- ❌ 重试进度(`session.error.retryAttempt / maxRetries`)— 需新增 + +--- + +## 7. Cost / Token 追踪 + +### 7.1 数据来源 + +| 来源 | 字段 | 用途 | +|---|---|---| +| `step-finish` part | `cost`, `tokens.{input, output, reasoning, cache.read, cache.write}` | 每步成本 | +| `task.completed` | `usage.{total_tokens, tool_uses, duration_ms}` | 任务级成本 | +| `message.updated` | `info.cost`, `info.tokens` | 消息级成本 | +| `session.status` (idle) | 可触发 session 级汇总 | 会话总成本 | + +### 7.2 UI 要求 + +- ⚠️ 会话级总成本显示(header / sidebar)— 部分实现,通过 message cost 累加 +- ❌ 每步成本 tooltip(hover step-finish 区域)— 需新增 +- ❌ 后台任务成本汇总 — 需新增 +- ❌ Cache token 显示(read vs write)— 需新增 + +--- + +## 8. 向后兼容消费 + +### 8.1 旧事件兼容层(过渡期) + +如果消费端仍需支持旧版 csc serve(未整改版本),应同时处理: + +| 旧事件 | 处理方式 | app-ai-native | +|---|---|---| +| `session.message` (type: assistant) | 解析 content blocks,自建 parts | ❌ 不处理旧格式 | +| `session.message` (type: user) | 解析 content blocks | ❌ 不处理旧格式 | +| `session.stream_event` | 自行实现 stream → parts 转换(同旧 cs-cloud adapter) | ❌ 不处理旧格式 | +| `session.message` (type: system) | 检查 subtype 字段,按子类型分发 | ❌ 不处理旧格式 | +| `session.result` | 提取 cost/usage/subtype | ❌ 不处理旧格式 | +| `session.control_request` | 权限/问答处理 | ❌ 不处理旧格式 | + +> **注**: app-ai-native 是全新实现,只消费 opencode canonical 格式,不兼容旧 `session.*` 前缀事件。旧版兼容由 cs-cloud adapter 层负责。 + +### 8.2 兼容判断 + +- ❌ SSE 连接后检查首条事件的格式(canonical vs legacy)— 不需要,仅支持 canonical +- ❌ 或通过 `GET /health` 的 `version` 字段判断 csc 版本 — 不需要 +- ❌ canonical 模式下忽略旧事件,legacy 模式下走旧路径 — 不需要 + +--- + +## 9. 总结:app-ai-native 缺失项(需新增实现) + +### 高优先级(核心体验影响) + +| 缺失项 | 说明 | 建议实现位置 | +|---|---|---| +| `session.error` 处理 | API 错误、重试等无任何 UI 反馈 | `device-session.tsx` 新增 switch case | +| `message.removed` 处理 | tombstone 消息无法被移除 | `device-session.tsx` 新增 switch case | +| `tool.progress` 处理 | Bash 等工具无实时输出流 | `device-session.tsx` 新增 switch case + ToolPart 扩展 | + +### 中优先级(增强体验) + +| 缺失项 | 说明 | 建议实现位置 | +|---|---|---| +| `task.*` 事件消费 | 无独立任务面板/追踪 | 新建 `useTaskState` hook + TaskPanel 组件 | +| `session.warning` / `session.info` | 系统警告/信息无展示 | `device-session.tsx` 新增 switch case | +| `message.attachment` | hook 结果、记忆等附加信息无展示 | `device-session.tsx` 新增 switch case | +| Cost/Token 详细追踪 | 仅有消息级累加,无步骤级/缓存级展示 | 扩展 StepFinishPart 渲染 | + +### 低优先级(锦上添花) + +| 缺失项 | 说明 | 建议实现位置 | +|---|---|---| +| `session.metrics` | turn_duration/TTFT 性能面板 | 新建 MetricsPanel 组件 | +| `session.hook_summary` | Hook 执行汇总展示 | `device-session.tsx` 新增 switch case | +| Task 进度条 | workflow phase 级别进度 | TaskPanel 组件内 | +| 后台任务通知 | 任务完成 toast/desktop notification | 新建通知系统 | +| Cache token 展示 | cache read vs write 分离显示 | CostPanel 组件内 | diff --git a/docs/serve/serve-event-standardization-proposal.md b/docs/serve/serve-event-standardization-proposal.md new file mode 100644 index 000000000..f4bf00b35 --- /dev/null +++ b/docs/serve/serve-event-standardization-proposal.md @@ -0,0 +1,872 @@ +# csc serve 事件标准化整改提案 + +> 状态:提案 (2026-05-18) +> 范围:`src/server/sessionMessageRouter.ts`、`src/server/eventBus.ts`、`src/server/sessionHandle.ts` +> 依赖:cs-cloud `internal/agent/csc/adapter*.go`(Phase 3 可移除) + +--- + +## 1. 问题概述 + +### 1.1 现状 + +`csc serve` 采用子进程模式(`--print --output-format stream-json`),子进程的 stdout 消息由 `sessionMessageRouter.ts` 路由为 SSE 事件。当前路由只处理 7 种消息类型: + +| 类型 | 处理 | SSE 事件 | +|---|---|---| +| `assistant` | handleAssistantMessage | `session.message` | +| `user` | handleUserMessage | `session.message` | +| `result` | handleResultMessage | `session.result` | +| `control_request` | handleControlRequest | `session.control_request` | +| `control_cancel_request` | handleCancelRequest | (opencode events) | +| `stream_event` | 直接转发 | `session.stream_event` | +| `system` | 统一发出 | `session.message`(不区分子类型) | +| **其余所有** | **default: break** | **丢弃** | + +### 1.2 丢弃的消息类型 + +以下消息类型到达 router 后被静默丢弃或信息不完整: + +| 消息类型 | TUI 行为 | serve 现状 | 影响 | +|---|---|---|---| +| `attachment` (50+ 子类型) | 专用组件渲染 | **丢弃** | Hook 结果、记忆、任务状态、诊断等全部不可见 | +| `progress` | 关联到工具消息显示 | **丢弃** | Bash 实时输出、Agent 进度不可见 | +| `system.task_notification` | `UserAgentNotificationMessage` 渲染 | 统一为 `session.message` | 消费端无法识别任务完成事件 | +| `system.task_started` | 任务面板追踪 | 统一为 `session.message` | 后台任务启动不可见 | +| `system.task_progress` | 进度条 | 统一为 `session.message` | 实时进度不可见 | +| `system.api_error` | 错误+重试倒计时 | 统一为 `session.message` | API 错误/重试信息不可用 | +| `system.compact_boundary` | 压缩边界指示器 | 统一为 `session.message` | 上下文压缩不可追踪 | +| `system.stop_hook_summary` | Hook 汇总 | 统一为 `session.message` | Hook 执行结果不可见 | +| `tombstone` | 移除孤立消息 | **不到达 router** | 压缩后孤立消息永久残留 | +| `stream_event` 内容 | TUI 细粒度处理 | raw Anthropic 格式转发 | 需要 cs-cloud Go adapter 重度适配 | + +### 1.3 cs-cloud 适配层的代价 + +cs-cloud 的 CSC adapter 有 ~800 行 Go 代码专门做格式转换: + +| 文件 | 行数 | 功能 | +|---|---|---| +| `adapter_sse_stream.go` | ~300 | Anthropic stream_event → message.part.delta | +| `adapter_sse_message.go` | ~350 | CSC 消息 → parts 分解 | +| `adapter_parts.go` | ~200 | 工具/文本 part 构建 | +| `adapter_json.go` | ~300 | REST 响应适配 | + +且 adapter 仍有大量遗漏:所有 `system` 子类型被丢弃、`tombstone` 不处理、`attachment` 不处理、`tool_progress` 信息不完整。 + +### 1.4 根本原因 + +**csc serve 输出的是「子进程协议格式」,不是「消费端协议格式」**。TUI 内部通过 `handleMessageFromStream()` 做了一层转换(stream event → React state),但 serve 模式缺少这层转换,将原始格式直接暴露给外部消费者。 + +--- + +## 2. 目标 + +1. **csc serve 原生输出 opencode 兼容的 canonical 事件格式**,消除 cs-cloud adapter 的适配需求 +2. **补齐 task 生命周期事件**:`task.started` → `task.progress` → `task.completed` +3. **补齐 system 子类型事件**:api_error、compact_boundary、hook_summary 等 +4. **补齐 attachment/progress 消息路由** +5. **保持向后兼容**:新旧事件格式并存过渡期 + +--- + +## 3. Canonical 事件格式规范 + +### 3.1 SSE 事件命名规范 + +对标 opencode serve 的事件总线,csc serve 应输出以下事件: + +| 事件名 | 替代现有 | 数据结构 | +|---|---|---| +| `session.created` | 保留 | `{ session_id, status, created_at }` | +| `session.updated` | 新增(替代部分 session.ready) | `{ session_id, status, model, provider_id }` | +| `session.deleted` | 保留 | `{ session_id }` | +| `session.status` | 保留 | `{ sessionID, status: { type: "busy"/"idle" } }` | +| `session.diff` | 新增 | `{ sessionID, diff: FileDiff[] }` | +| `session.error` | 新增 | `{ sessionID, error }` | +| `message.updated` | 新增(替代 session.message) | `{ sessionID, info: MessageInfo }` | +| `message.part.updated` | 新增(替代 stream_event 部分) | `{ sessionID, part: Part }` | +| `message.part.delta` | 新增(替代 stream_event 部分) | `{ sessionID, messageID, partID, field, delta }` | +| `permission.asked` | 已有(opencode 事件) | 保留 | +| `permission.replied` | 已有 | 保留 | +| `question.asked` | 已有 | 保留 | +| `question.replied` | 已有 | 保留 | +| `task.started` | 新增 | 见 §3.3 | +| `task.progress` | 新增 | 见 §3.3 | +| `task.completed` | 新增 | 见 §3.3 | + +### 3.2 Part 类型规范 + +每个消息由一个或多个 Part 组成: + +```typescript +type Part = + | { type: "text", id: string, text: string, time?: { start: number, end?: number } } + | { type: "reasoning", id: string, text: string, time?: { start: number, end?: number } } + | { type: "tool", id: string, callID: string, tool: string, + state: { status: "pending", input: Record } + | { status: "running", input: Record, title?: string, time: { start: number } } + | { status: "completed", input: Record, output: string, title: string, time: { start: number, end: number } } + | { status: "error", input: Record, error: string, time: { start: number, end: number } } + } + | { type: "step-start", id: string } + | { type: "step-finish", id: string, reason: string, cost: number, tokens: { input: number, output: number, reasoning: number, cache: { read: number, write: number } } } + | { type: "compaction", id: string, auto: boolean } + | { type: "subtask", id: string, prompt: string, description: string, agent: string } +``` + +### 3.3 Task 生命周期事件 + +```typescript +// task.started +{ + sessionID: string, + taskID: string, + toolUseID?: string, + description: string, + taskType?: string, // "local_agent" | "local_shell" | "local_workflow" + workflowName?: string, + prompt?: string, +} + +// task.progress +{ + sessionID: string, + taskID: string, + description: string, + usage: { total_tokens: number, tool_uses: number, duration_ms: number }, + lastToolName?: string, + summary?: string, + workflowProgress?: SdkWorkflowProgress[], +} + +// task.completed +{ + sessionID: string, + taskID: string, + toolUseID?: string, + status: "completed" | "failed" | "stopped", + summary: string, + outputFile: string, + usage?: { total_tokens: number, tool_uses: number, duration_ms: number }, +} +``` + +### 3.4 System 子类型事件映射 + +| system.subtype | canonical 事件 | 备注 | +|---|---|---| +| `task_notification` | `task.completed` | 终态通知 | +| `task_started` | `task.started` | 启动事件 | +| `task_progress` | `task.progress` | 进度事件 | +| `api_error` / `api_retry` | `session.error` + `message.part.updated { type: "retry" }` | 重试信息 | +| `compact_boundary` | `message.part.updated { type: "compaction" }` | 压缩边界 | +| `stop_hook_summary` | `session.hook_summary` | Hook 汇总 | +| `turn_duration` | `session.metrics` | Turn 耗时 | +| `cache_warning` | `session.warning` | 缓存警告 | +| `informational` | `session.info` | 一般信息 | + +--- + +## 4. 改造方案 + +### 4.1 Phase 1:sessionMessageRouter 增强(P0) + +#### 4.1.1 新增 StreamStateTracker + +**文件**:`src/server/streamStateTracker.ts`(新建) + +追踪每个 session 的流式状态,将 Anthropic 原始 stream_event 分解为 canonical parts 事件: + +```typescript +interface StreamState { + sessionID: string + messageID: string + parentID: string + modelID: string + activeBlocks: Map + stepPartID: string // step-start 的 part ID + usage: { + inputTokens: number + outputTokens: number + cacheReadTokens: number + cacheWriteTokens: number + } + stopReason: string +} +``` + +关键方法: + +```typescript +class StreamStateTracker { + // 处理 Anthropic stream_event,返回 canonical 事件数组 + processEvent(event: AnthropicStreamEvent, ctx: MessageRouterCtx): CanonicalEvent[] +} +``` + +处理逻辑: + +| Anthropic 事件 | 输出 canonical 事件 | +|---|---| +| `message_start` | `message.updated` (assistant) + `message.part.updated { type: "step-start" }` | +| `content_block_start` (text) | `message.part.updated { type: "text" }` | +| `content_block_start` (thinking) | `message.part.updated { type: "reasoning" }` | +| `content_block_start` (tool_use) | `message.part.updated { type: "tool", status: "pending" }` | +| `content_block_start` (redacted_thinking) | `message.part.updated { type: "reasoning", redacted: true }` | +| `content_block_delta` (text_delta) | `message.part.delta { field: "text" }` | +| `content_block_delta` (thinking_delta) | `message.part.delta { field: "text" }` | +| `content_block_delta` (input_json_delta) | `message.part.delta { field: "input" }` | +| `content_block_stop` (tool_use) | `message.part.updated { type: "tool", status: "running" }` | +| `content_block_stop` (text/thinking) | _(finalize part timing)_ | +| `message_delta` | _(extract usage + stopReason, no event)_ | +| `message_stop` | `message.part.updated { type: "step-finish", cost, tokens, reason }` | + +#### 4.1.2 改造 sessionMessageRouter.ts + +**改动范围**:`src/server/sessionMessageRouter.ts` + +##### a) stream_event case 改用 StreamStateTracker + +```typescript +case 'stream_event': { + ctx.setLastActiveAt(Date.now()) + const tracker = getOrCreateTracker(ctx.sessionId) + const canonicalEvents = tracker.processEvent(msg.event, ctx) + for (const event of canonicalEvents) { + ctx.emitOpencodeEvent(event.type, event.properties) + } + // 保留原始 stream_event 向后兼容(过渡期后移除) + ctx.emitEvent('stream_event', msg) + break +} +``` + +##### b) system case 增加子类型分发 + +```typescript +case 'system': { + const subtype = msg.subtype as string + switch (subtype) { + case 'task_notification': + emitTaskCompleted(msg, ctx) + break + case 'task_started': + emitTaskStarted(msg, ctx) + break + case 'task_progress': + emitTaskProgress(msg, ctx) + break + case 'api_error': + case 'api_retry': + emitSessionError(msg, ctx) + break + case 'compact_boundary': + case 'microcompact_boundary': + emitCompactionEvent(msg, ctx) + break + case 'stop_hook_summary': + emitHookSummary(msg, ctx) + break + default: + break + } + // 向后兼容:同时发 session.message + ctx.emitEvent('message', msg) + break +} +``` + +##### c) 新增 attachment case + +```typescript +case 'attachment': { + ctx.emitOpencodeEvent('message.attachment', { + sessionID: ctx.sessionId, + attachmentType: (msg.attachment as Record)?.type, + attachment: msg.attachment, + }) + // 不发 session.message(attachment 不是对话消息) + break +} +``` + +##### d) 新增 progress case + +```typescript +case 'progress': { + const toolUseID = msg.toolUseID as string + ctx.emitOpencodeEvent('tool.progress', { + sessionID: ctx.sessionId, + toolUseID, + parentToolUseID: msg.parentToolUseID, + data: msg.data, + }) + break +} +``` + +##### e) handleResultMessage 增强 + +```typescript +function handleResultMessage(msg: StdoutMessage, ctx: MessageRouterCtx): void { + // ... 现有逻辑 ... + + const stopReason = msg.stop_reason as string | undefined + const modelUsage = msg.modelUsage as Record | undefined + + ctx.emitOpencodeEvent('session.status', { + sessionID: ctx.sessionId, + status: { type: 'idle' }, + }) + + // 如果有 diff 信息,发射 session.diff + // 如果有错误子类型,发射 session.error + if (msg.subtype !== 'success') { + ctx.emitOpencodeEvent('session.error', { + sessionID: ctx.sessionId, + error: { + subtype: msg.subtype, + is_error: msg.is_error, + errors: msg.errors, + }, + }) + } +} +``` + +#### 4.1.3 新增 tombstone 传递 + +**文件**:`src/QueryEngine.ts` + +当前 tombstone 在 QueryEngine 中被拦截不转发。需要在 serve 模式下将其转发到 stdout: + +```typescript +// QueryEngine.ts — yield 逻辑中 +if (message.type === 'tombstone') { + // 现有:skipped + // 新增:在 serve 模式下转发 + if (getIsNonInteractiveSession()) { + yield message as StdoutMessage + } + continue +} +``` + +**sessionMessageRouter.ts** 新增 case: + +```typescript +case 'tombstone': { + const targetUuid = (msg.message as { uuid: string })?.uuid + if (targetUuid) { + ctx.emitOpencodeEvent('message.removed', { + sessionID: ctx.sessionId, + messageID: targetUuid, + }) + } + break +} +``` + +### 4.2 Phase 2:assistant 消息 parts 分解(P1) + +当收到完整的 `assistant` 消息(非流式,如历史加载),需要在 router 层将其 content blocks 分解为 parts: + +```typescript +function handleAssistantMessage(msg: StdoutMessage, ctx: MessageRouterCtx): void { + // ... 现有逻辑 ... + + // 新增:发射 message.updated 事件 + ctx.emitOpencodeEvent('message.updated', { + sessionID: ctx.sessionId, + info: { + id: msg.uuid, + role: 'assistant', + modelID: msg.model, + providerID: msg.provider_id, + cost: 0, // 单条消息无独立 cost + tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } }, + time: { created: msg.timestamp ? new Date(msg.timestamp as string).getTime() : Date.now() }, + parentID: msg.parentUuid ?? null, + }, + }) + + // 新增:分解 content blocks 为 parts + const content = msg.message?.content as Array> ?? [] + for (const block of content) { + const part = buildPartFromContentBlock(block, msg) + if (part) { + ctx.emitOpencodeEvent('message.part.updated', { + sessionID: ctx.sessionId, + part, + }) + } + } +} + +function buildPartFromContentBlock(block: Record, msg: StdoutMessage): Part | null { + switch (block.type) { + case 'text': + return { type: 'text', id: uuid(), text: block.text as string } + case 'thinking': + return { type: 'reasoning', id: uuid(), text: block.thinking as string } + case 'redacted_thinking': + return { type: 'reasoning', id: uuid(), text: '', redacted: true } + case 'tool_use': + return { + type: 'tool', id: uuid(), + callID: block.id as string, + tool: normalizeToolName(block.name as string), + state: { + status: 'running', + input: block.input as Record, + time: { start: Date.now() }, + }, + } + default: + return null + } +} +``` + +### 4.3 Phase 3:向后兼容过渡(P1) + +#### 4.3.1 双写期 + +在 Phase 1/2 实施后,新旧事件同时发出: + +``` +session.message (旧) → 保留 +message.updated (新) → 新增 +session.stream_event (旧) → 保留 +message.part.delta (新) → 新增 +session.control_request (旧) → 保留 +permission.asked (新) → 已有 +``` + +#### 4.3.2 cs-cloud 逐步切换 + +cs-cloud adapter 逐步切换到新事件: + +1. **v1**:adapter 优先消费新事件,旧事件作为 fallback +2. **v2**:adapter 只消费新事件,移除 `adapter_sse_stream.go`、`adapter_sse_message.go`、`adapter_parts.go` +3. **v3**:csc serve 移除旧事件输出,cs-cloud adapter 变为 thin proxy + +#### 4.3.3 版本协商(可选) + +SSE 连接可通过 `Accept-Event-Version: 2` header 选择新格式: + +```typescript +// event route handler +const eventVersion = c.req.header('accept-event-version') +if (eventVersion === '2') { + // 只发 canonical 事件 +} else { + // 双写 +} +``` + +--- + +## 5. 新增事件处理函数 + +### 5.1 Task 生命周期 + +```typescript +// sessionMessageRouter.ts — 新增函数 + +function emitTaskStarted(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('task.started', { + sessionID: ctx.sessionId, + taskID: msg.task_id as string, + toolUseID: msg.tool_use_id as string | undefined, + description: msg.description as string, + taskType: msg.task_type as string | undefined, + workflowName: msg.workflow_name as string | undefined, + prompt: msg.prompt as string | undefined, + }) +} + +function emitTaskProgress(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('task.progress', { + sessionID: ctx.sessionId, + taskID: msg.task_id as string, + description: msg.description as string, + usage: msg.usage, + lastToolName: msg.last_tool_name as string | undefined, + summary: msg.summary as string | undefined, + workflowProgress: msg.workflow_progress, + }) +} + +function emitTaskCompleted(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('task.completed', { + sessionID: ctx.sessionId, + taskID: msg.task_id as string, + toolUseID: msg.tool_use_id as string | undefined, + status: msg.status as 'completed' | 'failed' | 'stopped', + summary: msg.summary as string, + outputFile: msg.output_file as string, + usage: msg.usage, + }) +} +``` + +### 5.2 System 子类型 + +```typescript +function emitSessionError(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('session.error', { + sessionID: ctx.sessionId, + error: { + subtype: msg.subtype, + level: msg.level ?? 'error', + message: msg.content ?? msg.error?.message, + retryInMs: msg.retry_in_ms ?? msg.retryInMs, + retryAttempt: msg.retry_attempt ?? msg.retryAttempt, + maxRetries: msg.max_retries ?? msg.maxRetries, + }, + }) +} + +function emitCompactionEvent(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('message.part.updated', { + sessionID: ctx.sessionId, + part: { + type: 'compaction', + id: msg.uuid ?? randomUUID(), + auto: msg.subtype === 'microcompact_boundary' || (msg.compact_metadata as any)?.trigger === 'auto', + }, + }) +} + +function emitHookSummary(msg: StdoutMessage, ctx: MessageRouterCtx): void { + ctx.emitOpencodeEvent('session.hook_summary', { + sessionID: ctx.sessionId, + hookLabel: msg.hook_label ?? msg.hookLabel, + hookCount: msg.hook_count ?? msg.hookCount, + hookErrors: msg.hook_errors ?? msg.hookErrors, + preventedContinuation: msg.prevented_continuation ?? msg.preventedContinuation, + totalDurationMs: msg.total_duration_ms ?? msg.totalDurationMs, + }) +} +``` + +--- + +## 6. StreamStateTracker 详细设计 + +### 6.1 状态管理 + +```typescript +// src/server/streamStateTracker.ts + +import { randomUUID } from 'crypto' + +interface BlockState { + type: string + partID: string + toolUseID?: string + toolName?: string + inputJson: string + startTime: number +} + +interface SessionStreamState { + messageID: string + parentID: string + modelID: string + activeBlocks: Map + stepStartPartID: string + assistantPartEmitted: boolean + usage: { + inputTokens: number + outputTokens: number + cacheReadTokens: number + cacheWriteTokens: number + } + stopReason: string +} + +const states = new Map() + +function getOrCreate(sessionID: string): SessionStreamState { + let state = states.get(sessionID) + if (!state) { + state = { + messageID: randomUUID(), + parentID: '', + modelID: '', + activeBlocks: new Map(), + stepStartPartID: randomUUID(), + assistantPartEmitted: false, + usage: { inputTokens: 0, outputTokens: 0, cacheReadTokens: 0, cacheWriteTokens: 0 }, + stopReason: '', + } + states.set(sessionID, state) + } + return state +} + +function reset(sessionID: string): void { + states.delete(sessionID) +} +``` + +### 6.2 事件处理 + +```typescript +type CanonicalEvent = { + type: string + properties: Record +} + +export function processStreamEvent( + sessionID: string, + event: Record, +): CanonicalEvent[] { + const results: CanonicalEvent[] = [] + const state = getOrCreate(sessionID) + const eventType = event.type as string + + switch (eventType) { + case 'message_start': { + const msg = (event as any).message + state.messageID = randomUUID() + state.modelID = msg?.model ?? '' + state.usage.inputTokens = msg?.usage?.input_tokens ?? 0 + state.usage.cacheReadTokens = msg?.usage?.cache_read_input_tokens ?? 0 + state.usage.cacheWriteTokens = msg?.usage?.cache_creation_input_tokens ?? 0 + + if (!state.assistantPartEmitted) { + results.push({ + type: 'message.updated', + properties: { + sessionID, + info: { + id: state.messageID, + role: 'assistant', + modelID: state.modelID, + time: { created: Date.now() }, + }, + }, + }) + state.assistantPartEmitted = true + } + + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { type: 'step-start', id: state.stepStartPartID, sessionID, messageID: state.messageID }, + }, + }) + break + } + + case 'content_block_start': { + const block = (event as any).content_block + const index = (event as any).index as number + const partID = randomUUID() + const blockState: BlockState = { + type: block.type, + partID, + inputJson: '', + startTime: Date.now(), + } + + state.activeBlocks.set(index, blockState) + + if (block.type === 'text') { + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { type: 'text', id: partID, sessionID, messageID: state.messageID, text: '' }, + }, + }) + } else if (block.type === 'thinking' || block.type === 'redacted_thinking') { + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { + type: 'reasoning', id: partID, sessionID, messageID: state.messageID, + text: '', redacted: block.type === 'redacted_thinking', + }, + }, + }) + } else if (block.type === 'tool_use') { + blockState.toolUseID = block.id + blockState.toolName = block.name + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { + type: 'tool', id: partID, sessionID, messageID: state.messageID, + callID: block.id, tool: block.name, + state: { status: 'pending', input: {}, raw: '' }, + }, + }, + }) + } + break + } + + case 'content_block_delta': { + const index = (event as any).index as number + const delta = (event as any).delta + const blockState = state.activeBlocks.get(index) + + if (!blockState) break + + if (delta.type === 'text_delta') { + results.push({ + type: 'message.part.delta', + properties: { sessionID, messageID: state.messageID, partID: blockState.partID, field: 'text', delta: delta.text }, + }) + } else if (delta.type === 'thinking_delta') { + results.push({ + type: 'message.part.delta', + properties: { sessionID, messageID: state.messageID, partID: blockState.partID, field: 'text', delta: delta.thinking }, + }) + } else if (delta.type === 'input_json_delta') { + blockState.inputJson += delta.partial_json + results.push({ + type: 'message.part.delta', + properties: { sessionID, messageID: state.messageID, partID: blockState.partID, field: 'input', delta: delta.partial_json }, + }) + } + break + } + + case 'content_block_stop': { + const index = (event as any).index as number + const blockState = state.activeBlocks.get(index) + if (!blockState) break + + if (blockState.type === 'tool_use') { + let parsedInput: Record = {} + try { parsedInput = JSON.parse(blockState.inputJson || '{}') } catch {} + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { + type: 'tool', id: blockState.partID, sessionID, messageID: state.messageID, + callID: blockState.toolUseID, tool: blockState.toolName, + state: { status: 'running', input: parsedInput, time: { start: blockState.startTime } }, + }, + }, + }) + } + break + } + + case 'message_delta': { + const delta = (event as any).delta + const usage = (event as any).usage + if (delta?.stop_reason) state.stopReason = delta.stop_reason + if (usage?.output_tokens) state.usage.outputTokens = usage.output_tokens + break + } + + case 'message_stop': { + results.push({ + type: 'message.part.updated', + properties: { + sessionID, + part: { + type: 'step-finish', id: randomUUID(), sessionID, messageID: state.messageID, + reason: state.stopReason || 'stop', + cost: 0, + tokens: { + input: state.usage.inputTokens, + output: state.usage.outputTokens, + reasoning: 0, + cache: { read: state.usage.cacheReadTokens, write: state.usage.cacheWriteTokens }, + }, + }, + }, + }) + reset(sessionID) + break + } + } + + return results +} +``` + +--- + +## 7. 测试计划 + +### 7.1 单元测试 + +| 测试文件 | 覆盖范围 | +|---|---| +| `src/server/__tests__/streamStateTracker.test.ts` | 所有 Anthropic stream event → canonical event 转换 | +| `src/server/__tests__/sessionMessageRouter.test.ts` | system 子类型分发、attachment/progress/tombstone 路由 | + +### 7.2 集成测试 + +- 启动 csc serve,发送 prompt,验证 SSE 事件流包含: + - `message.updated` + `message.part.updated` (text/reasoning/tool parts) + - `message.part.delta` (流式 delta) + - `task.started` / `task.completed`(后台 agent 场景) + - `session.error`(API 错误场景) + - `message.part.updated { type: "compaction" }`(压缩场景) + +### 7.3 兼容性测试 + +- cs-cloud adapter 在双写期仍能正常工作(消费旧事件) +- cs-cloud adapter 切换到新事件后功能不降级 + +--- + +## 8. 实施计划 + +| Phase | 内容 | 预估工时 | 依赖 | +|---|---|---|---| +| Phase 1a | StreamStateTracker 实现 + stream_event 分解 | 3-4 天 | 无 | +| Phase 1b | system 子类型分发(task_*、api_error、compact_boundary) | 2 天 | 无 | +| Phase 1c | attachment / progress / tombstone 路由 | 1-2 天 | 无 | +| Phase 1d | handleResultMessage 增强 | 0.5 天 | 无 | +| Phase 2 | assistant 消息 parts 分解(历史加载场景) | 2 天 | Phase 1a | +| Phase 3a | cs-cloud adapter 优先消费新事件 | 3-4 天 | Phase 1+2 | +| Phase 3b | 移除 cs-cloud adapter 重度适配代码 | 1 天 | Phase 3a | +| Phase 3c | csc serve 移除旧事件(下个 major 版本) | 0.5 天 | Phase 3b | + +--- + +## 9. 风险与缓解 + +| 风险 | 缓解措施 | +|---|---| +| 双写期 SSE 事件量翻倍影响性能 | StreamStateTracker 产出的事件替代而非追加;旧事件可在 Phase 3 移除 | +| StreamStateTracker 状态泄漏 | message_stop 后 reset;session 关闭时清理 | +| 子进程协议变更影响非 serve 消费者 | 只改 sessionMessageRouter(serve 专属),不改子进程 stdout 协议 | +| tool input JSON 解析失败 | try/catch 包裹,fallback 为 `{}` | +| 多 step(tool_use 循环)中 messageID 管理错误 | 每次 message_start 生成新 messageID,step-finis/step-start 配对 | + +--- + +## 10. 参考 + +- `src/server/sessionMessageRouter.ts` — 当前消息路由 +- `src/server/eventBus.ts` — SSE 事件总线 +- `src/server/sessionHandle.ts` — 子进程管理 +- `src/utils/handleMessageFromStream()` — TUI 的流式事件处理(`src/utils/messages.ts:3262`) +- `packages/@ant/model-provider/src/types/message.ts` — 消息类型定义 +- `src/utils/sdkEventQueue.ts` — SDK 事件队列(task_started/task_progress/task_notification) +- cs-cloud `internal/agent/csc/adapter_sse_stream.go` — Go 层流式适配(对标参考) +- cs-cloud `internal/agent/csc/adapter_sse_message.go` — Go 层消息适配(对标参考) +- opencode `packages/opencode/src/session/message-v2.ts` — canonical Part 类型定义 +- opencode `packages/opencode/src/session/processor.ts` — 流式事件处理参考 diff --git a/docs/serve/serve-event-standardization-todo.md b/docs/serve/serve-event-standardization-todo.md new file mode 100644 index 000000000..c39e68392 --- /dev/null +++ b/docs/serve/serve-event-standardization-todo.md @@ -0,0 +1,198 @@ +# csc serve 事件标准化 — 实施任务清单 + +> 关联文档:`docs/serve/serve-event-standardization-proposal.md` +> 原则:**不做双写**,csc serve 直接输出 canonical 格式,旧 cs-cloud passthrough 兼容 + +--- + +## Phase 1:sessionMessageRouter 路由增强 + +### 1.1 新建 StreamStateTracker + +- [x] 创建 `src/server/streamStateTracker.ts` +- [x] 实现 `SessionStreamState` 状态结构(messageID、activeBlocks、usage、stopReason) +- [x] 实现 `processStreamEvent()` — Anthropic stream_event → canonical 事件转换 + - [x] `message_start` → `message.updated` + `step-start` part + - [x] `content_block_start` (text) → `message.part.updated { type: "text" }` + - [x] `content_block_start` (thinking) → `message.part.updated { type: "reasoning" }` + - [x] `content_block_start` (redacted_thinking) → `message.part.updated { type: "reasoning", redacted: true }` + - [x] `content_block_start` (tool_use) → `message.part.updated { type: "tool", status: "pending" }` + - [x] `content_block_delta` (text_delta) → `message.part.delta { field: "text" }` + - [x] `content_block_delta` (thinking_delta) → `message.part.delta { field: "text" }` + - [x] `content_block_delta` (input_json_delta) → `message.part.delta { field: "input" }` + - [x] `content_block_stop` (tool_use) → `message.part.updated { status: "running" }` + - [x] `content_block_stop` (text/thinking) → finalize timing + - [x] `message_delta` → 提取 usage.outputTokens + stopReason(不输出事件) + - [x] `message_stop` → `message.part.updated { type: "step-finish" }` + reset state +- [x] 工具名规范化函数(`Bash` → `bash`、`Read` → `read` 等,复用现有 `toPermissionKey` 映射) +- [x] tool input JSON 累积 + try/catch 解析 fallback +- [x] 多 step(tool_use 循环)场景:每次 `message_start` 生成新 messageID + +### 1.2 改造 stream_event 路由 + +- [x] `sessionMessageRouter.ts` — `case 'stream_event'` 改用 StreamStateTracker +- [x] 移除旧 `ctx.emitEvent('stream_event', msg)` — 不再做 raw 转发 +- [x] Tracker 产出的 canonical 事件通过 `ctx.emitOpencodeEvent()` 发出 + +### 1.3 system 子类型分发 + +- [x] `case 'system'` — 替换当前的统一 `ctx.emitEvent('message', msg)` +- [x] `task_notification` → `ctx.emitOpencodeEvent('task.completed', { taskID, status, summary, usage, ... })` +- [x] `task_started` → `ctx.emitOpencodeEvent('task.started', { taskID, description, taskType, ... })` +- [x] `task_progress` → `ctx.emitOpencodeEvent('task.progress', { taskID, usage, summary, workflowProgress, ... })` +- [x] `api_error` / `api_retry` → `ctx.emitOpencodeEvent('session.error', { error: { subtype, message, retryInMs, ... } })` +- [x] `compact_boundary` / `microcompact_boundary` → `ctx.emitOpencodeEvent('message.part.updated', { part: { type: "compaction", auto, overflow } })` +- [x] `stop_hook_summary` → `ctx.emitOpencodeEvent('session.hook_summary', { ... })` +- [x] `turn_duration` → `ctx.emitOpencodeEvent('session.metrics', { ... })` +- [x] `cache_warning` → `ctx.emitOpencodeEvent('session.warning', { ... })` +- [x] `informational` → `ctx.emitOpencodeEvent('session.info', { ... })` +- [x] `post_turn_summary` → `ctx.emitOpencodeEvent('session.info', { ... })` +- [x] `session_state_changed` → `ctx.emitOpencodeEvent('session.status', { ... })` +- [x] `status` → `ctx.emitOpencodeEvent('session.status', { ... })` +- [x] `default` → `ctx.emitOpencodeEvent('session.info', { subtype, ...msg })` + +### 1.4 新增 attachment 路由 + +- [x] `case 'attachment'` — 当前 default:break 丢弃 +- [x] 转发为 `ctx.emitOpencodeEvent('message.attachment', { attachmentType, attachment })` +- [x] 关键 attachment 子类型至少覆盖: + - [x] `hook_success` / `hook_error` / `hook_cancelled` + - [x] `relevant_memories` / `nested_memory` + - [x] `task_status` / `task_reminder` + - [x] `diagnostics` + - [x] `token_usage` / `budget_usd` + - [x] `invoked_skills` + +> 注:attachment 路由不区分子类型,统一通过 `attachmentType` 字段透传。消费端按需过滤。 + +### 1.5 新增 progress 路由 + +- [x] `case 'progress'` — 当前 default:break 丢弃 +- [x] 转发为 `ctx.emitOpencodeEvent('tool.progress', { toolUseID, parentToolUseID, data })` + +### 1.6 新增 tombstone 路由 + +- [x] `src/QueryEngine.ts` — serve 模式下将 tombstone 消息 yield 到 stdout +- [x] `sessionMessageRouter.ts` — `case 'tombstone'` +- [x] 转发为 `ctx.emitOpencodeEvent('message.removed', { messageID })` + +### 1.7 handleResultMessage 增强 + +- [x] 提取 `stop_reason` 传播到 `session.result` 的 reason 字段 +- [x] 提取 `usage` / `cost_usd` 透传到 `session.result` 事件 +- [x] `subtype !== 'success'` 时发射 `session.error` 事件 +- [x] 区分 result 子类型:`success` / `error_max_turns` / `error_max_budget` / `error_doom_loop` / `error` + +### 1.8 handleAssistantMessage 增强(非流式 / 历史加载场景) + +- [x] 完整 assistant 消息到达时,发射 `message.updated` 事件(含 role/modelID/cost/tokens/parentID) +- [x] 遍历 content blocks,为每个 block 发射 `message.part.updated` + - [x] `text` block → text part + - [x] `thinking` block → reasoning part + - [x] `redacted_thinking` block → reasoning part (redacted) + - [x] `tool_use` block → tool part (status: running) + +--- + +## Phase 2:REST API 适配 + +### 2.1 消息历史端点标准化 + +- [x] `GET /session/{id}/message` — 返回 parts-based 格式(对标 opencode) + - [x] assistant 消息:content blocks 分解为 parts 数组 + - [x] user 消息:text/tool_result 分解为 parts + - [x] 新增 `format=parts` query parameter 消费端按需选择 + - [ ] system 消息:compaction_boundary → compaction part + - [ ] attachment 消息:保留原始数据但增加 part 包装 +- [ ] `GET /session/{id}/todo` — 适配响应格式 + +--- + +## Phase 3:EventBus / SSE 层调整 + +### 3.1 emitEvent → emitOpencodeEvent 统一 + +- [x] 审计 `sessionHandle.ts` 中所有 `ctx.emitEvent()` 调用 +- [x] `emitEvent('message', ...)` → `emitOpencodeEvent('message.updated', ...)` +- [x] `emitEvent('result', ...)` → `emitOpencodeEvent('session.result', ...)` +- [x] `emitEvent('ready', ...)` → `emitOpencodeEvent('session.updated', ...)` +- [x] `emitEvent('deleted', ...)` → `emitOpencodeEvent('session.deleted', ...)` +- [x] `emitEvent('stream_event', ...)` → 移除(Phase 1.2 已由 StreamStateTracker 替代) +- [x] `emitEvent('control_request', ...)` → 保留(已有独立 opencode 事件) +- [x] `emitEvent('permission_replied', ...)` / `emitEvent('question_replied', ...)` → 保留(低级别确认) + +### 3.2 SSE 事件名映射 + +- [x] 确认 `emitOpencodeEvent` 输出的 SSE event name 格式 +- [x] 全部统一为 opencode 风格(无 `session.` 前缀) + +--- + +## Phase 4:测试 + +### 4.1 单元测试 + +- [x] `src/server/__tests__/streamStateTracker.test.ts` (18 tests, 35 assertions) + - [x] message_start → message.updated + step-start + - [x] text content_block 完整流程 (start → delta → stop) + - [x] thinking content_block 完整流程 + - [x] redacted_thinking 处理 + - [x] tool_use 完整流程 (start → input_delta → stop → running) + - [x] message_delta 提取 usage/stopReason + - [x] message_stop → step-finish + reset + - [x] 多 step 场景(tool_use 循环产生多个 message_start/stop) + - [x] input_json_delta 累积 + 损坏 JSON fallback +- [x] `src/server/__tests__/sessionMessageRouter.test.ts` (29 tests, 57 assertions) + - [x] system.task_notification → task.completed + - [x] system.task_started → task.started + - [x] system.task_progress → task.progress + - [x] system.api_error → session.error + - [x] system.compact_boundary → compaction part + - [x] attachment 路由 + - [x] progress 路由 + - [x] tombstone → message.removed + - [x] result 子类型区分 + - [x] assistant 非流式消息 → parts 分解 + - [x] init → session.updated + - [x] user → message.updated + +### 4.2 集成测试 + +- [ ] 启动 csc serve,发送 prompt,验证 SSE 流包含完整 canonical 事件序列 +- [ ] 后台 agent 场景:验证 task.started → task.progress → task.completed +- [ ] API 错误场景:验证 session.error 事件 +- [ ] 压缩场景:验证 compaction part +- [ ] 历史消息加载:验证 parts-based 格式 + +### 4.3 兼容性测试 + +- [ ] 旧版 cs-cloud + 新 csc serve:确认旧 adapter 正常处理保留的事件名 +- [ ] 旧版 cs-cloud passthrough 新事件:消费端不报错 +- [ ] 新版 cs-cloud + 新 csc serve:确认 adapter 切换到新事件后功能不降级 + +--- + +## Phase 5:cs-cloud Adapter 精简(cs-cloud 侧) + +### 5.1 切换到新事件源 + +- [ ] cs-cloud adapter 优先消费 `message.part.updated` / `message.part.delta` 而非 `session.stream_event` +- [ ] cs-cloud adapter 消费 `task.started` / `task.progress` / `task.completed` 而非从 `session.message` 推断 +- [ ] cs-cloud adapter 消费 `session.error` 而非丢弃 system 子类型 + +### 5.2 移除适配代码 + +- [ ] 移除 `adapter_sse_stream.go`(StreamStateTracker 已在 csc 侧完成) +- [ ] 移除 `adapter_sse_message.go`(parts 分解已在 csc 侧完成) +- [ ] 移除 `adapter_parts.go`(part 构建已在 csc 侧完成) +- [ ] 保留 `adapter_json.go`(REST API 响应仍需适配,直到 Phase 2 完成) +- [ ] adapter 降级为 thin proxy(仅 SSE passthrough + REST 路由) + +--- + +## 不在本次范围 + +- Event Sourcing / CQRS(架构变更过大) +- 多 Workspace 实例隔离(cs-cloud 多进程模型已解决) +- File snapshot / patch 追踪(依赖文件系统监控基础设施) +- Session fork / share 功能实现(需要 csc 核心逻辑支持)