docs: add serve event standardization and consumer capability docs
This commit is contained in:
parent
c63d8c60ce
commit
600170ecf3
326
docs/serve/consumer-capability-checklist.md
Normal file
326
docs/serve/consumer-capability-checklist.md
Normal file
|
|
@ -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<string, TaskState>
|
||||
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 组件内 |
|
||||
872
docs/serve/serve-event-standardization-proposal.md
Normal file
872
docs/serve/serve-event-standardization-proposal.md
Normal file
|
|
@ -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<string,any> }
|
||||
| { status: "running", input: Record<string,any>, title?: string, time: { start: number } }
|
||||
| { status: "completed", input: Record<string,any>, output: string, title: string, time: { start: number, end: number } }
|
||||
| { status: "error", input: Record<string,any>, 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<number, {
|
||||
type: string // "text" | "thinking" | "tool_use" | "redacted_thinking"
|
||||
partID: string
|
||||
toolUseID?: string
|
||||
toolName?: string
|
||||
inputJson: string // 累积的 tool input JSON
|
||||
}>
|
||||
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<string, unknown>)?.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<string, unknown> | 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<Record<string, unknown>> ?? []
|
||||
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<string, unknown>, 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<string, unknown>,
|
||||
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<number, BlockState>
|
||||
stepStartPartID: string
|
||||
assistantPartEmitted: boolean
|
||||
usage: {
|
||||
inputTokens: number
|
||||
outputTokens: number
|
||||
cacheReadTokens: number
|
||||
cacheWriteTokens: number
|
||||
}
|
||||
stopReason: string
|
||||
}
|
||||
|
||||
const states = new Map<string, SessionStreamState>()
|
||||
|
||||
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<string, unknown>
|
||||
}
|
||||
|
||||
export function processStreamEvent(
|
||||
sessionID: string,
|
||||
event: Record<string, unknown>,
|
||||
): 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<string, unknown> = {}
|
||||
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` — 流式事件处理参考
|
||||
198
docs/serve/serve-event-standardization-todo.md
Normal file
198
docs/serve/serve-event-standardization-todo.md
Normal file
|
|
@ -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 核心逻辑支持)
|
||||
Loading…
Reference in New Issue
Block a user