diff --git a/.github/workflows/validate.yml b/.github/workflows/validate.yml index e426a27..5c05b01 100644 --- a/.github/workflows/validate.yml +++ b/.github/workflows/validate.yml @@ -47,3 +47,22 @@ jobs: - name: Build WebUI and Electron run: bun run build + + - name: Setup Node + uses: actions/setup-node@v4 + with: + node-version: "20" + + - name: Install Chromium and system dependencies + run: npx playwright install --with-deps chromium + + - name: Run Chromium motion smoke (web-only) + run: node scripts/verify-ui.cjs --web-only + + - name: Upload QA artifacts on failure + if: failure() + uses: actions/upload-artifact@v4 + with: + name: qa-artifacts + path: artifacts/qa/ + if-no-files-found: ignore diff --git a/README.md b/README.md index a057e7c..519ff78 100644 --- a/README.md +++ b/README.md @@ -1,205 +1,195 @@ -

-

ThreadCove

-

面向个人深度研究的 AI 信息分析工作台

-

- -

- Bun - TypeScript strict - React - Electron - tests -

+
---- +# ThreadCove -ThreadCove 把复杂研究问题拆成多个**并行任务**(搜索 → 阅读 → 分析 → 汇总),每个任务一个独立 Session 与工作区;按任务阶段选模型(抓取整理求快省、分析综合求强推理);外部能力(搜索 API / 网页抓取 / 本地文档 / MCP 工具)统一接入。Electron 桌面端是主力工作台,WebUI 是同一套逻辑的第二个载体——**传输层只写一遍**。 +### 让零散线索,长成自己的洞见。 -> 个人项目 · 工程实践风格:类型严格、测试护住核心机制、CI 常驻、提交历史干净。不是企业级系统——没有数据库/RBAC/监控/容灾。 +从一个问题出发,连接资料、追问与推演,留下一份可以继续生长的研究档案。 -## ✨ 核心能力 +*A personal AI research workbench. Follow the question. Keep the thread.* -| | 能力 | 说明 | -|---|---|---| -| 🔌 | **模型后端解耦** | `AgentBackend` 统一 `chat()` 接口,DeepSeek / Claude / Pi 三后端可换,切换 = 改一行配置 | -| 🗂 | **任务隔离** | 一个研究任务一个 Session + 独立工作目录,并行任务的产物物理分家,互不污染 | -| 🌊 | **统一事件流** | `AgentEvent` 可辨识联合类型——文本增量/工具/权限/结构化错误/完成,双后端出同一种事件,UI 按 type 分流渲染 | -| 🔑 | **统一能力接入** | MCP / API / 本地数据源统一转成 Agent 可调用工具,凭据加密集中管理,调用时自动注入 | -| 🖥 | **多端复用** | 每条 RPC 通道显式分类 LOCAL_ONLY / REMOTE_ELIGIBLE(穷举测试强制),桌面直连与 Web 代理共用一套协议 | +[![Validate](https://github.com/Bluuok/ThreadCove/actions/workflows/validate.yml/badge.svg)](https://github.com/Bluuok/ThreadCove/actions/workflows/validate.yml) +[![Bun](https://img.shields.io/badge/Bun-%E2%89%A51.3.10-526f85?logo=bun&logoColor=white)](https://bun.sh) +[![React](https://img.shields.io/badge/React-18-526f85?logo=react&logoColor=white)](packages/ui) +[![Electron](https://img.shields.io/badge/Electron-39-526f85?logo=electron&logoColor=white)](apps/electron) +[![License: MIT](https://img.shields.io/badge/License-MIT-526f85)](LICENSE) -## 🏗 架构 +[产品界面](#产品界面) · [研究流程](#研究流程) · [快速开始](#快速开始) · [项目文档](#项目文档) -```text -研究任务 = Session(会话上下文 + 工具上下文 + 独立工作目录) -└── Workspace(rootPath:偏好默认值 / sources / skills / sessions 存储) - └── AgentBackend(deepseek=OpenAI兼容流式 / anthropic=进程内 SDK / pi=进程外 JSONL 子进程) - └── AgentEvent 统一词表 × 多后端适配 × EventQueue 桥 × typed_error 恢复 - └── Session/Workspace 隔离 · Source→Credential→Tool 能力接入(运行时汇合) - └── 传输层 CHANNEL_MAP + WsRpc + 路由穷举 → 多端 UI -``` +ThreadCove 研究首页:雾蓝色工作台、研究图谱与问题输入框 -### 包结构 +**FIELDNOTES LAB / 个人研究现场** -```text -packages/ -├── core/ 纯类型层:AgentEvent / Message / Session / Workspace + win32 路径可移植工具 -└── shared/ - ├── protocol/ channels · dto · events · routing(穷举测试) · types · codec - ├── agent/ backend 抽象 / factory / claude / pi / deepseek / EventQueue / 事件适配器 - ├── sessions/ session.jsonl 存储 · 持久化队列 · 路径可移植 · 防穿越 - ├── workspaces/ Workspace CRUD + defaults - ├── sources/ SourceServerBuilder(三层header优先级) · api-tools · storage - ├── credentials/ AES-256-GCM 机器绑定加密库 - ├── mcp/ McpClientPool 中央池 + 代理工具命名 - ├── skills/ SKILL.md 加载 + requiredSources 校验 - └── config/ MODEL_REGISTRY(能力字段) · 思考档 · 权限模式 -packages/pi-agent-server/ Pi 子进程:stdin/stdout JSONL 协议循环 -apps/ -├── electron/ 主进程(内嵌 WS server + handlers + SessionManager) · transport · preload · React renderer -└── webui/ 复用同一 transport + CHANNEL_MAP,只覆写 LOCAL_ONLY 方法 -``` +
-### 一条消息的完整旅程 +## 为什么做 ThreadCove -```text -UI 调用 api.sendMessage() - → buildClientApi 生成的代理 → RoutedClient(REMOTE_ELIGIBLE → WsRpcClient) - → WS envelope(codec JSON+base64) → WsRpcServer → handlers → SessionManager - → AgentBackend.chat() 产出 AsyncGenerator - ├─ deepseek:OpenAI 兼容 /chat/completions SSE 流 - ├─ anthropic:Claude Agent SDK query() 流式消息 → 事件适配器 - └─ pi:spawn 子进程 → JSONL init/prompt → 原生事件 → PiEventAdapter → EventQueue 桥 - → 每个事件实时广播 session:event 到所有已连接客户端 + 按 parity 契约落盘 session.jsonl +一个值得深入的问题,往往会带来许多网页、几段对话和零散的笔记。材料越来越多,思路却容易在窗口切换中断掉。 + +ThreadCove 把这些过程放回同一个工作台:提出问题,让 AI 协助探索资料,在对话中持续追问,再把结论和文件留在当前任务里。下次回来,仍然能找到上次思考的起点。 + +> 一个研究任务,一段连续的对话,一处独立的资料空间。 + +## 研究流程 + +```mermaid +flowchart LR + A[提出问题] --> B[搜索与阅读] + B --> C[追问与整理] + C --> D[形成结论] + D --> E[保存与归档] + E -. 继续探索 .-> A + classDef paper fill:#edf2f5,stroke:#7893a6,color:#294356; + class A,B,C,D,E paper; ``` -## 🚀 快速开始 +- **开始**:记录一个疑问,或从首页的线索节点获得提问灵感。 +- **深入**:由模型按需调用搜索、网页读取和子任务工具,在同一段对话里综合资料。 +- **留下**:查看研究回复与任务文件,归档阶段性成果,随时重新打开并继续追问。 + +这是使用路径,不是固定的执行阶段或进度条;具体工具调用由模型与所选后端决定。 + +## 产品界面 -**环境要求**:[Bun](https://bun.sh) ≥ 1.3(Windows/macOS/Linux 均可,本仓库在 win32 上开发) +### 从一个好问题开始 + +雾蓝纸面、档案索引与节点连线,组成一张可以回应你的研究手记。首页保持安静,聚焦线索时,再用轻微的描线反馈引导探索。 + +### 让思考有上下文,让成果有去处 + +| 研究对话 | 资料与成果 | +| :---: | :---: | +| [![研究对话:问题、回复与运行状态](docs/images/research-conversation.png)](docs/images/research-conversation.png) | [![研究手记:任务文件列表与内容预览](docs/images/research-notes.png)](docs/images/research-notes.png) | +| **沿着线索,让思考逐渐清晰。**
提问、流式回复与运行反馈,留在同一条研究线上。 | **将探索留下,成为下一次研究的起点。**
打开任务文件,在侧栏查看与当前研究关联的内容。 | + +截图来自实际运行的界面;对话和文件为人工准备的演示内容,不代表真实在线研究结果。点击双栏图片可查看原图。 + +## 核心能力 + +| 能力 | 在工作台里做什么 | +| --- | --- | +| **连续研究对话** | 提问、追问、查看流式回复;需要时停止当前运行,保留已收到的内容。 | +| **搜索与网页阅读** | Pi 后端可按需检索资料、读取公开网页,并将来源产物保存在任务目录。 | +| **有界子任务协作** | 将独立问题交给子任务研究,由主任务汇总;宿主管理并发、次数、超时和取消。 | +| **独立研究档案** | 每个任务有独立会话与文件目录,支持历史恢复、归档和继续研究。 | +| **资料与执行侧栏** | 查看工具活动、权限请求和任务文件;大文件提供有上限的内容预览。 | +| **桌面与浏览器共用** | Electron 与 Web 共用工作台界面和事件协议,连接同一套会话服务。 | +| **模型与工具可接入** | 提供 Pi、DeepSeek、Claude 后端适配,以及 MCP/API 资料源接入层。具体能力因后端而异。 | + +## 快速开始 + +需要 [Bun](https://bun.sh) **≥ 1.3.10**。浏览器验收脚本另需 Node.js。 + +### 1. 获取项目并配置模型 ```bash git clone https://github.com/Bluuok/ThreadCove.git cd ThreadCove bun install +``` -# 配置 LLM 密钥(任选其一,写入 .env 或环境变量) -cp .env.example .env # 编辑 .env,填写自己的密钥 -# ANTHROPIC_API_KEY=sk-ant-... # Claude +将 `.env.example` 复制为 `.env`(PowerShell 可使用 `Copy-Item .env.example .env`),再填写自己的凭据。仓库示例使用 Pi + OpenCode Go: -# 验证工程链路 -bun run typecheck # 全仓 strict tsc 零错误 -bun run lint -bun test # 无密钥时自动跳过 5 个 live 测试 -bun run build # WebUI 和 Electron 生产构建 +```dotenv +THREADCOVE_PROVIDER=pi +THREADCOVE_PI_PROVIDER=opencode-go +THREADCOVE_MODEL=deepseek-v4.1-flash +OPENCODE_GO_API_KEY=你的密钥 +``` -# 桌面端 -bun run dev:electron +也可使用本机加密配置,步骤见 [OpenCode Go 接入说明](docs/OPENCODE-GO.md)。模型请求使用所配置服务的额度;密钥不要提交到仓库。 -# WebUI + headless 服务(演示同一套逻辑服务两种载体) -bun run server:headless 8787 -bun run dev:webui # 打开 http://localhost:5173 -``` +### 2. 打开桌面工作台 -WebUI 启动后,使用 headless 终端输出的带 fragment 引导链接连接;随机 token 只保存在当前浏览器会话。服务默认监听本机,数据位于 `.threadcove-workspace/`。模型可通过 `THREADCOVE_PROVIDER` / `THREADCOVE_MODEL` 配置,任务设置支持持久化模型覆盖。 +```bash +bun run dev:electron +``` -本次评审整改、兼容边界和验证步骤见 [整改说明](docs/IMPLEMENTATION-REVIEW.md);OpenDesign 规范与实际前端移植见 [前端设计](docs/FRONTEND-DESIGN.md)。[OpenCode Go 接入说明](docs/OPENCODE-GO.md) 提供本机加密配置与真实 Pi SDK 验证方法。OpenCode Go / DeepSeek V4.1 Flash 已通过真实 SDK、宿主工具往返及 Web/Electron 联调;DeepSeek 直连接口和 Claude API 本次未验证。 +### 3. 或在浏览器中使用 -**验证流式对话**(需要真实密钥,无需 UI): +分别在两个终端运行: ```bash -DEEPSEEK_API_KEY=sk-... bun run scripts/smoke-deepseek.ts +# 终端一:启动本地会话服务 +bun run server:headless 8787 ``` -## 🎨 三个值得看的设计 +```bash +# 终端二:启动 Web 界面 +bun run dev:webui +``` -### 1. redirect 双分支——统一接口不假装所有后端一样 +使用 headless 终端输出的带 fragment 引导链接连接。服务默认监听本机,数据默认保存在 `.threadcove-workspace/`;连接 token 保存在当前浏览器会话中。 -```typescript -// pi-agent.ts — 有原生 steering:消息注入当前流,事件继续走 -redirect(message: string): boolean { - if (!this.isProcessing()) { this.forceAbort(AbortReason.Redirect); return false; } - this.send({ type: 'steer', message }); - return true; // ← true:会话层什么都不用做 -} +## 项目结构 -// deepseek-agent.ts — 无原生 steering:内部排队,会话层决定重发时机 +```text +ThreadCove/ +├── apps/ +│ ├── electron/ 桌面宿主、会话服务与 RPC 处理 +│ └── webui/ 浏览器入口与传输适配 +├── packages/ +│ ├── ui/ 两端共用的研究工作台 +│ ├── core/ 会话、消息与事件类型 +│ ├── shared/ 后端适配、工具、存储与通信协议 +│ └── pi-agent-server/ Pi SDK 子进程与宿主工具桥接 +├── scripts/ 开发与验收脚本 +└── docs/ 接入说明、设计文档与界面截图 ``` -接口签名一个 `boolean`,把「各后端能力不同」这个事实诚实地暴露给上层,而不是用一个统一假象掩盖它。 +界面通过统一 RPC 调用会话服务;模型后端产出统一事件,由宿主保存并推送到客户端。任务存储、模型适配和界面展示各自独立,便于更换后端或调整交互。 -### 2. EventQueue——不是 async generator 换皮 +## 开发与验证 -Claude 是同步 `for-await` 消费 SDK 流;Pi 子进程事件是**异步回调**到达。`EventQueue`(`enqueue` 推入唤醒 / `drain` 等待产出 / `complete` 收流 / `reset` 翻新)把 push 变 pull,让两种截然不同的喂法产出**同一个** `AsyncGenerator` 消费面。六个测试锁住:顺序保证、空队列等待、提前 complete 缓冲投递、reset 复用、交错生产。 +```bash +bun run typecheck +bun run lint +bun test +bun run build -### 3. 路由穷举——每条通道都被显式决策过 +# 安装浏览器后,运行 Web 界面验收 +bunx playwright install chromium +bun run test:ui:web -```typescript -// packages/shared/tests/routing.test.ts -test('every channel is classified exactly once', () => { - for (const ch of all) { - const inLocal = LOCAL_ONLY_CHANNELS.has(ch); - const inRemote = REMOTE_ELIGIBLE_CHANNELS.has(ch); - if (!inLocal && !inRemote) throw new Error(`Channel "${ch}" 未分类`); - if (inLocal && inRemote) throw new Error(`Channel "${ch}" 双重分类`); - } -}); +# 本地完整验收,包含 Electron 窗口 +bun run test:ui ``` -新增通道不分类 → CI 直接红。「文件对话框只能在桌面走,会话内容两端都要走」不是约定俗成,是测试强制的架构约束。 +CI 执行类型检查、lint、自动测试、双端构建与 Chromium 界面验收。真实模型调用需单独配置凭据;自动界面验收使用本地测试服务,不证明真实上游服务可用。 -## 📡 传输协议(摘要) +## 使用边界 -```text -client → { type:'handshake', protocolVersion:'1.0', token?, workspaceId? } -server → { type:'handshake_ack', clientId, registeredChannels } +- 面向个人、本地优先的研究场景,目前没有多用户权限系统、向量知识库或跨任务长期记忆。 +- 搜索与子任务编排属于 Pi 研究工具链;其他后端不能视为具备相同能力。 +- 网页读取面向公开 HTML、文本和 JSON;登录页面、PDF 和依赖 JavaScript 渲染的内容不在当前读取范围内。 +- 中断的子任务会记录状态,不会自动恢复其 SDK 内存上下文或重新执行。来源内容与模型结论仍需核对。 -request { id, type:'request', channel:'sessions:sendMessage', args:[...] } -response { id, type:'response', result:{...} } | error { code:'CHANNEL_NOT_FOUND', message } -event { id, type:'event', channel:'session:event', args:[{ sessionId, event:AgentEvent }] } -``` +## 项目文档 -- 二进制走 `{"__tcRpcType":"u8","base64":"..."}` 标注编码,两端无损还原 -- 错误码:`HANDLER_ERROR` · `CHANNEL_NOT_FOUND` · `AUTH_FAILED` · `PROTOCOL_VERSION_UNSUPPORTED` · `REQUEST_TIMEOUT`(类标识不跨线,接收端按 `err.code` 分支) -- 非 localhost 明文 `ws://` 一律拒绝(token 明文防护) -- 完整通道分类表与 envelope 细节见 **[docs/PROTOCOL.md](docs/PROTOCOL.md)** - -## 🔒 安全模型 - -- **凭据**:AES-256-GCM 加密落盘(`credentials.enc`,64B 头 + PBKDF2(机器硬件 ID ∧ 随机盐, 100k 轮))——文件拷到别的机器解不开。机器标识是**密钥派生输入**,不是密钥本身 -- **进程边界**:MCP/API 源凭据只在宿主进程;Pi 子进程执行源工具时发 `tool_execute_request` 回宿主,**凭据永不进入子进程环境** -- **渲染进程**:preload 只有引导逻辑,无手工 IPC 分发表;`contextIsolation: true`;未分类通道过不了路由穷举测试 -- 细节与已知边界见 **[docs/SECURITY.md](docs/SECURITY.md)** - -## 🧪 测试地图 - -| 套件 | 护住什么 | -|---|---| -| `routing.test` | 通道分类穷举:完备/总数相等/零交集/非空 | -| `codec.test` | envelope 往返含二进制标注编码、畸形包拒绝 | -| `event-queue.test` | push→pull 桥:顺序/等待/收流/复位/交错 | -| `factory.test` | provider→backend 路由、未知 provider 抛错、生命周期面 | -| `claude/pi-event-adapter.test` | 双后端原生事件 → 统一词表快照 + 未知事件容错 | -| `pi-protocol.test` | 真实子进程:init→ready 握手、事件流、tool_execute 往返、坏 JSONL 容错、abort、redirect 双分支 | -| `session-event-message-parity.test` | 事件字段 ↔ 持久化字段双端一致(防字段漂移) | -| `sessions.test` | 双会话物理隔离、路径穿越、归档/删除、并发持久化、`{{SESSION_PATH}}` 往返(含 Windows 转义) | -| `sources.test` | 三层 header 优先级、缺 token→null、Authorization 三态、AES 加解密往返+换机失败、代理命名映射、SKILL.md 校验 | -| `transport/bootstrap.test` | 真实 WS 收发、token 鉴权、广播、buildClientApi、明文 ws:// 拒绝 | -| `integration.test` | 端到端:建会话→发消息→收流→落盘、多会话隔离、后端切换、typed error code | -| `deepseek.test` | 无密钥 typed_error;**有密钥时跑 live**:流式/用量/中断/标题生成 | - -## 📚 文档 - -- **[docs/REPRODUCTION_SPEC.md](docs/REPRODUCTION_SPEC.md)** — 完整复现规格:范围/深度/边界/不做清单 -- **[docs/PROTOCOL.md](docs/PROTOCOL.md)** — 传输协议:envelope 格式 + 通道分类表 -- **[docs/SECURITY.md](docs/SECURITY.md)** — 凭据加密模型 + 安全设计动机 - -## 🗺 边界(诚实交底) - -- **不做**:RAG/向量库(主打实时抓取+多源核对)、多 Agent 编排(多任务=多 Session)、Memory/跨任务知识复用、完整 OAuth(凭据手动粘贴,只留 `getToken` 钩子) -- **MCP 只到统一转换层**:官方 SDK Client + 代理工具组装,不自研 JSON-RPC/握手 -- WebUI 定位是「验证同一套逻辑能否复用」的接续查看面,不是全功能第二产品 -- Pi 后端已接入真实 `createAgentSession`,默认注册免密钥搜索、公开网页读取及最多 3 个并发子任务;详见[研究工具与预算](docs/RESEARCH-TOOLS.md) -- 已有真实 OpenCode Go API / Pi SDK / Web / Electron 验收,以及两个真实子任务搜索、读取、来源落盘和刷新恢复验证。自动测试不使用个人 API 密钥,DeepSeek 直连 live 测试默认跳过;1M 上下文尚未实测满窗口 +| 文档 | 内容 | +| --- | --- | +| [研究工具与子任务](docs/RESEARCH-TOOLS.md) | 搜索、网页读取、编排约束与来源保存 | +| [OpenCode Go 接入](docs/OPENCODE-GO.md) | API 配置、本机凭据与 Pi SDK 验证 | +| [前端设计](docs/FRONTEND-DESIGN.md) | 设计方向与共享界面实现 | +| [通信协议](docs/PROTOCOL.md) | RPC、事件与多端路由 | +| [安全模型](docs/SECURITY.md) | 凭据管理、进程边界与已知限制 | +| [实现与验收记录](docs/IMPLEMENTATION-REVIEW.md) | 实现调整与验证步骤 | +| [项目范围与复现规格](docs/REPRODUCTION_SPEC.md) | 设计范围、实现深度与边界 | + +## 参与项目 + +欢迎通过 [Issues](https://github.com/Bluuok/ThreadCove/issues) 反馈问题或提出想法。提交改动时,请附上复现步骤与相关验证,并保持一个改动对应一个明确范围。 ## License -MIT +[MIT](LICENSE) + +--- + +
+ +**不只获得一个回答,也留下通往答案的线索。** + +ThreadCove · FIELDNOTES LAB + +
diff --git a/apps/electron/src/server/handlers.ts b/apps/electron/src/server/handlers.ts index c9d0fd1..66098ec 100644 --- a/apps/electron/src/server/handlers.ts +++ b/apps/electron/src/server/handlers.ts @@ -23,9 +23,10 @@ import { resolveSessionFilePath, updateSessionConfig, validateSessionId, + readFilePreview, } from '@threadcove/shared/sessions'; import { listSources } from '@threadcove/shared/sources'; -import { readFileSync, writeFileSync, readdirSync, existsSync, statSync } from 'fs'; +import { writeFileSync, readdirSync, existsSync, statSync } from 'fs'; import { join } from 'path'; import type { ModelProvider } from '@threadcove/shared/config'; import type { SessionManager } from './session-manager.ts'; @@ -229,7 +230,7 @@ export function registerHandlers(server: RpcServerLike, ctx: HandlerContext): vo // getSessionPath imported at module top (defense-in-depth sanitize inside) const file = resolveSessionFilePath(root, String(sessionId), String(subPath)); if (!existsSync(file)) throw new Error(`File not found: ${String(subPath)}`); - return { content: readFileSync(file, 'utf-8') }; + return readFilePreview(file); }); server.handle(RPC_CHANNELS.files.WRITE, (...args: unknown[]) => { diff --git a/docs/images/research-conversation.png b/docs/images/research-conversation.png new file mode 100644 index 0000000..77f93f6 Binary files /dev/null and b/docs/images/research-conversation.png differ diff --git a/docs/images/research-home.png b/docs/images/research-home.png new file mode 100644 index 0000000..2d2c33b Binary files /dev/null and b/docs/images/research-home.png differ diff --git a/docs/images/research-notes.png b/docs/images/research-notes.png new file mode 100644 index 0000000..f7f835a Binary files /dev/null and b/docs/images/research-notes.png differ diff --git a/package.json b/package.json index 4cf0b38..fd5dec7 100644 --- a/package.json +++ b/package.json @@ -12,6 +12,7 @@ "scripts": { "test": "bun test", "test:ui": "node scripts/verify-ui.cjs", + "test:ui:web": "node scripts/verify-ui.cjs --web-only", "typecheck": "bun run typecheck:core && bun run typecheck:shared && bun run typecheck:pi-agent-server && bun run typecheck:electron && bun run typecheck:webui && bun run --cwd packages/ui typecheck", "typecheck:core": "cd packages/core && bun run typecheck", "typecheck:shared": "cd packages/shared && bun run typecheck", diff --git a/packages/shared/src/client/types.ts b/packages/shared/src/client/types.ts index db37f1c..c1103f4 100644 --- a/packages/shared/src/client/types.ts +++ b/packages/shared/src/client/types.ts @@ -50,7 +50,7 @@ export interface ElectronAPI { // Files (workspace content) listFiles(workspaceId: string, sessionId: string, subPath?: string): Promise<{ entries: unknown[] }>; - readFile(workspaceId: string, sessionId: string, subPath: string): Promise<{ content: string }>; + readFile(workspaceId: string, sessionId: string, subPath: string): Promise; writeFile(workspaceId: string, sessionId: string, subPath: string, content: string): Promise<{ success: boolean }>; // LOCAL_ONLY — native dialog (Electron only; WebUI overrides with input[type=file]) @@ -71,3 +71,11 @@ export interface SessionEventPayload { sessionId: string; event: AgentEvent | { type: 'user_message'; message: unknown }; } + +/** Result shape for bounded file reading. */ +export interface FileReadResult { + content: string; + truncated?: boolean; + totalBytes?: number; + previewBytes?: number; +} diff --git a/packages/shared/src/sessions/index.ts b/packages/shared/src/sessions/index.ts index 48fc3d0..29e389c 100644 --- a/packages/shared/src/sessions/index.ts +++ b/packages/shared/src/sessions/index.ts @@ -4,3 +4,4 @@ export * from './persistence-queue.ts'; export * from './validation.ts'; export * from './slug-generator.ts'; export * from './file-path.ts'; +export * from './preview.ts'; diff --git a/packages/shared/src/sessions/preview.ts b/packages/shared/src/sessions/preview.ts new file mode 100644 index 0000000..5ccc45c --- /dev/null +++ b/packages/shared/src/sessions/preview.ts @@ -0,0 +1,92 @@ +import { openSync, fstatSync, readSync, closeSync } from 'node:fs'; + +export const MAX_FULL_PREVIEW_BYTES = 1024 * 1024; // 1 MiB +export const TRUNCATED_PREVIEW_BYTES = 256 * 1024; // 256 KiB + +export interface FilePreviewResult { + content: string; + truncated?: boolean; + totalBytes: number; + previewBytes: number; +} + +/** + * Truncates a buffer to end at a valid UTF-8 character boundary. + * Prevents mangled replacement characters (\uFFFD) when slicing multibyte sequences. + */ +export function truncateToValidUtf8(buf: Buffer): Buffer { + const len = buf.length; + if (len === 0) return buf; + + let i = len - 1; + let continuations = 0; + while (i >= 0 && i >= len - 4) { + const byte = buf[i]!; + if ((byte & 0x80) === 0) { + // ASCII byte, sequence complete + break; + } + if ((byte & 0xc0) === 0x80) { + // Continuation byte (10xxxxxx) + continuations++; + i--; + continue; + } + // Found leading byte + let expected = 0; + if ((byte & 0xe0) === 0xc0) expected = 1; + else if ((byte & 0xf0) === 0xe0) expected = 2; + else if ((byte & 0xf8) === 0xf0) expected = 3; + + if (expected > 0 && continuations < expected) { + // Incomplete trailing sequence: strip lead byte and trailing continuation bytes + return buf.subarray(0, i); + } + // Complete sequence + return buf; + } + + // Trailing orphan continuation bytes without a leading byte + if (continuations > 0 && i < len - continuations) { + return buf.subarray(0, len - continuations); + } + + return buf; +} + +/** + * Reads a bounded file preview using a single opened file descriptor. + * - Rejects non-regular files + * - Bounded reading: <= 1 MiB full read, > 1 MiB reads at most 256 KiB + * - UTF-8 character truncation never forges or mangles trailing bytes + */ +export function readFilePreview(filePath: string): FilePreviewResult { + const fd = openSync(filePath, 'r'); + try { + const stat = fstatSync(fd); + if (!stat.isFile()) { + throw new Error(`Not a regular file: ${filePath}`); + } + const totalBytes = stat.size; + const isTruncated = totalBytes > MAX_FULL_PREVIEW_BYTES; + const maxToRead = isTruncated ? TRUNCATED_PREVIEW_BYTES : Math.min(totalBytes, MAX_FULL_PREVIEW_BYTES); + + const buf = Buffer.alloc(maxToRead); + let bytesRead = 0; + if (maxToRead > 0) { + bytesRead = readSync(fd, buf, 0, maxToRead, 0); + } + + const rawSlice = buf.subarray(0, bytesRead); + const validSlice = isTruncated ? truncateToValidUtf8(rawSlice) : rawSlice; + + return { + content: validSlice.toString('utf-8'), + truncated: isTruncated, + totalBytes, + previewBytes: validSlice.length, + }; + } finally { + closeSync(fd); + } +} diff --git a/packages/shared/tests/preview.test.ts b/packages/shared/tests/preview.test.ts new file mode 100644 index 0000000..f4b8845 --- /dev/null +++ b/packages/shared/tests/preview.test.ts @@ -0,0 +1,126 @@ +import { describe, it, expect } from 'bun:test'; +import { mkdtempSync, writeFileSync, statSync, readFileSync, mkdirSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { + readFilePreview, + truncateToValidUtf8, + MAX_FULL_PREVIEW_BYTES, + TRUNCATED_PREVIEW_BYTES, +} from '../src/sessions/preview.ts'; +import { resolveSessionFilePath } from '../src/sessions/file-path.ts'; + +describe('readFilePreview', () => { + it('reads small files completely (<= threshold)', () => { + const dir = mkdtempSync(join(tmpdir(), 'tc-preview-test-')); + const file = join(dir, 'small.txt'); + const text = 'Hello, ThreadCove!'.repeat(50); + writeFileSync(file, text, 'utf-8'); + + const result = readFilePreview(file); + expect(result.truncated).toBe(false); + expect(result.content).toBe(text); + expect(result.totalBytes).toBe(Buffer.byteLength(text)); + expect(result.previewBytes).toBe(Buffer.byteLength(text)); + }); + + it('bounds reading to at most 256KiB for files > 1MiB (> threshold)', () => { + const dir = mkdtempSync(join(tmpdir(), 'tc-preview-test-')); + const file = join(dir, 'large.txt'); + const count = 1200; // 1200 KiB > 1024 KiB + const fd = Buffer.alloc(count * 1024, 65); // All 'A's + writeFileSync(file, fd); + + const statBefore = statSync(file); + const result = readFilePreview(file); + + expect(result.truncated).toBe(true); + expect(result.totalBytes).toBe(statBefore.size); + expect(result.previewBytes).toBe(TRUNCATED_PREVIEW_BYTES); + expect(result.content.length).toBe(TRUNCATED_PREVIEW_BYTES); + expect(Buffer.byteLength(result.content, 'utf-8')).toBe(TRUNCATED_PREVIEW_BYTES); + }); + + it('safely truncates multibyte UTF-8 characters without corruption', () => { + // Chinese character '中' is 3 bytes: 0xE4, 0xB8, 0xAD + // Emoji '🎉' is 4 bytes: 0xF0, 0x9F, 0x8E, 0x89 + // Let's test truncateToValidUtf8 directly with partial sequences + const char3 = Buffer.from('中', 'utf-8'); // 3 bytes + expect(char3.length).toBe(3); + + // Buffer ending with only 1 byte of a 3-byte sequence + const partial1 = Buffer.concat([Buffer.from('ASCII'), char3.subarray(0, 1)]); + const fixed1 = truncateToValidUtf8(partial1); + expect(fixed1.toString('utf-8')).toBe('ASCII'); + expect(fixed1.toString('utf-8').includes('\uFFFD')).toBe(false); + + // Buffer ending with 2 bytes of a 3-byte sequence + const partial2 = Buffer.concat([Buffer.from('ASCII'), char3.subarray(0, 2)]); + const fixed2 = truncateToValidUtf8(partial2); + expect(fixed2.toString('utf-8')).toBe('ASCII'); + expect(fixed2.toString('utf-8').includes('\uFFFD')).toBe(false); + + // Buffer ending with full 3-byte sequence + const complete3 = Buffer.concat([Buffer.from('ASCII'), char3]); + const fixed3 = truncateToValidUtf8(complete3); + expect(fixed3.toString('utf-8')).toBe('ASCII中'); + + // Emoji 4-byte partial + const emoji = Buffer.from('🎉', 'utf-8'); // 4 bytes + const partialEmoji = Buffer.concat([Buffer.from('Hello'), emoji.subarray(0, 3)]); + const fixedEmoji = truncateToValidUtf8(partialEmoji); + expect(fixedEmoji.toString('utf-8')).toBe('Hello'); + expect(fixedEmoji.toString('utf-8').includes('\uFFFD')).toBe(false); + + // File-level test: place a 3-byte character right at 256KiB boundary + const dir = mkdtempSync(join(tmpdir(), 'tc-preview-test-')); + const file = join(dir, 'multibyte-boundary.txt'); + // 262143 bytes of 'A', then '中' (3 bytes), followed by padding up to 1.1 MiB + const prefix = Buffer.alloc(TRUNCATED_PREVIEW_BYTES - 1, 65); // 262143 bytes + const middle = Buffer.from('中', 'utf-8'); // 3 bytes, starts at byte 262143, ends at 262146 + const rest = Buffer.alloc(MAX_FULL_PREVIEW_BYTES, 66); + writeFileSync(file, Buffer.concat([prefix, middle, rest])); + + const res = readFilePreview(file); + expect(res.truncated).toBe(true); + expect(res.content.includes('\uFFFD')).toBe(false); + expect(res.previewBytes).toBe(TRUNCATED_PREVIEW_BYTES - 1); + expect(res.content).toBe(prefix.toString('utf-8')); + }); + + it('rejects directory paths', () => { + const dir = mkdtempSync(join(tmpdir(), 'tc-preview-test-')); + expect(() => readFilePreview(dir)).toThrow(/regular file/i); + }); + + it('ensures target file is never modified by preview read', () => { + const dir = mkdtempSync(join(tmpdir(), 'tc-preview-test-')); + const file = join(dir, 'read-only-test.txt'); + const content = 'UNTOUCHED_CONTENT_TEST_'.repeat(100); + writeFileSync(file, content, 'utf-8'); + + const beforeStat = statSync(file); + const res = readFilePreview(file); + const afterStat = statSync(file); + + expect(res.content).toBe(content); + expect(afterStat.size).toBe(beforeStat.size); + expect(afterStat.mtimeMs).toBe(beforeStat.mtimeMs); + expect(readFileSync(file, 'utf-8')).toBe(content); + }); + + it('preserves path containment and security via resolveSessionFilePath', () => { + const root = mkdtempSync(join(tmpdir(), 'tc-containment-test-')); + const sessionDir = join(root, 'sessions', 'sess-1', 'data'); + mkdirSync(sessionDir, { recursive: true }); + writeFileSync(join(sessionDir, 'artifact.txt'), 'clean artifact'); + + // Valid path resolves under data + const valid = resolveSessionFilePath(root, 'sess-1', 'artifact.txt'); + expect(valid).toBe(join(sessionDir, 'artifact.txt')); + + // Traversal attempts throw + expect(() => resolveSessionFilePath(root, 'sess-1', '../session.jsonl')).toThrow(); + expect(() => resolveSessionFilePath(root, 'sess-1', '../../secret.txt')).toThrow(); + }); +}); diff --git a/packages/ui/src/components/Conversation.tsx b/packages/ui/src/components/Conversation.tsx index 73ad4ea..f001a51 100644 --- a/packages/ui/src/components/Conversation.tsx +++ b/packages/ui/src/components/Conversation.tsx @@ -1,3 +1,4 @@ +import { memo, useCallback } from 'react'; import type { StoredMessage } from '@threadcove/core/types'; import { RunIndicator } from './RunIndicator.tsx'; @@ -13,14 +14,34 @@ type Props = { onCopyError: (error: unknown) => void; }; +type MessageItemProps = { + message: StoredMessage; + onCopyError: (error: unknown) => void; +}; + +export const MessageItem = memo(function MessageItem({ message, onCopyError }: MessageItemProps) { + const handleCopy = useCallback(() => { + void navigator.clipboard.writeText(message.content).catch(onCopyError); + }, [message.content, onCopyError]); + + return ( +
+
+ {message.type === 'user' ? '你 / 提问' : message.type === 'assistant' ? 'THREADCOVE / 研究回复' : message.type === 'tool' ? `${message.toolName || '工具'} / 结果` : '执行记录'} + +
+
{message.content}
+ {message.type === 'assistant' && } +
+ ); +}); + export function Conversation({ messages, status, online, waiting, completion, sessionId, scrollRef, onScroll, onCopyError }: Props) { return
- {messages.map(message =>
-
{message.type === 'user' ? '你 / 提问' : message.type === 'assistant' ? 'THREADCOVE / 研究回复' : message.type === 'tool' ? `${message.toolName || '工具'} / 结果` : '执行记录'}
-
{message.content}
- {message.type === 'assistant' && } -
)} + {messages.map(message => ( + + ))}
; diff --git a/packages/ui/src/components/DetailsPanel.tsx b/packages/ui/src/components/DetailsPanel.tsx index 4de2bea..9c77755 100644 --- a/packages/ui/src/components/DetailsPanel.tsx +++ b/packages/ui/src/components/DetailsPanel.tsx @@ -7,7 +7,13 @@ type Props = { active?: SessionDto; sources: SourceDto[]; files: Array<{ name: string; type: string }>; - filePreview?: { name: string; content: string }; + filePreview?: { + name: string; + content: string; + truncated?: boolean; + totalBytes?: number; + previewBytes?: number; + }; activity: string[]; permission?: Permission; onClose: () => void; @@ -20,6 +26,6 @@ export function DetailsPanel({ active, sources, files, filePreview, activity, pe {permission &&

等待你的允许

{permission.description}

}

执行记录 01

{active?.isProcessing ? '正在研究' : runNames[active?.lastRun?.status ?? ''] ?? '等待问题'}
{activity.map((item, index) =>

{item}

)}{!activity.length &&

工具调用发生时,会在这里显示。

}

资料源 {String(sources.length).padStart(2, '0')}

{sources.map(source =>
{source.name}{source.enabled ? '已配置' : '已停用'}
)}{!sources.length &&

尚未配置资料源。

}
-

任务文件 {String(files.length).padStart(2, '0')}

{files.map(file => )}{!files.length &&

任务资料目录中还没有文件。

}{filePreview &&
{filePreview.name}
{filePreview.content}
}
+

任务文件 {String(files.length).padStart(2, '0')}

{files.map(file => )}{!files.length &&

任务资料目录中还没有文件。

}{filePreview &&
{filePreview.name}{filePreview.truncated &&
文件较大,仅显示前 {Math.round((filePreview.previewBytes ?? 262144) / 1024)} KiB
}
{filePreview.content}
}
; } diff --git a/packages/ui/src/components/RunIndicator.tsx b/packages/ui/src/components/RunIndicator.tsx index 012046a..e0af13a 100644 --- a/packages/ui/src/components/RunIndicator.tsx +++ b/packages/ui/src/components/RunIndicator.tsx @@ -3,25 +3,89 @@ import './conversation-motion.css'; export type RunFeedback = { sessionId: string; runId: string; status: string; completion?: string }; type Props = { status?: string; online: boolean; waiting: boolean; completion?: string }; -const labels: Record = { running: '正在研究…', completed: '研究已完成', failed: '研究失败', cancelled: '研究已取消', interrupted: '研究已中断' }; -/** Only a new live completion token can animate; mounting historical state stays quiet. */ +export const runIndicatorLabels: Record = { + running: '正在研究…', + completed: '研究已完成', + failed: '研究失败', + cancelled: '研究已取消', + interrupted: '研究已中断', +}; + +export const TERMINAL_STATUSES = new Set(['completed', 'cancelled', 'failed', 'interrupted']); + +/** + * Pure state calculation for RunIndicator: + * - Terminal states (completed, cancelled, failed, interrupted) remain preserved when offline. + * - Only unfinished active research shows 'disconnected' ("连接中断,执行状态待同步") when offline. + * - Idle/empty tasks without active runs do not show disconnect banners. + */ +export function computeRunIndicatorState( + status?: string, + online = true, + waiting = false, +): { state?: string; label?: string } { + let state: string | undefined; + if (status && TERMINAL_STATUSES.has(status)) { + state = status; + } else if (!online) { + if (status === 'running' || waiting) { + state = 'disconnected'; + } else { + state = undefined; + } + } else if (waiting) { + state = 'waiting'; + } else { + state = status; + } + + if (!state || (state !== 'disconnected' && state !== 'waiting' && !runIndicatorLabels[state])) { + return {}; + } + + const label = state === 'disconnected' + ? '连接中断,执行状态待同步' + : state === 'waiting' + ? '等待你的确认' + : runIndicatorLabels[state]; + + return { state, label }; +} + +/** Only a new live completion token can animate; mounting historical state stays quiet; reconnecting does not celebrate. */ export function RunIndicator({ status, online, waiting, completion }: Props) { - const previous = useRef(completion); + const animatedCompletions = useRef>(new Set()); const [finishing, setFinishing] = useState(false); - const state = !online ? 'disconnected' : waiting ? 'waiting' : status; + + const { state, label } = computeRunIndicatorState(status, online, waiting); + useEffect(() => { - const fresh = Boolean(completion && completion !== previous.current); - previous.current = completion; - if (state !== 'completed' || !fresh) { setFinishing(false); return; } + if (state !== 'completed' || !completion || animatedCompletions.current.has(completion)) { + setFinishing(false); + return; + } + animatedCompletions.current.add(completion); setFinishing(true); const timer = setTimeout(() => setFinishing(false), 260); return () => clearTimeout(timer); }, [completion, state]); - if (!state || (online && !waiting && !labels[state])) return null; - const label = state === 'disconnected' ? '连接中断,执行状态待同步' : state === 'waiting' ? '等待你的确认' : labels[state]; - return
-
; + + if (!state || !label) return null; + + return ( +
+
+ ); } diff --git a/packages/ui/src/streamBatcher.ts b/packages/ui/src/streamBatcher.ts new file mode 100644 index 0000000..4e85cd0 --- /dev/null +++ b/packages/ui/src/streamBatcher.ts @@ -0,0 +1,101 @@ +import type { StoredMessage } from '@threadcove/core/types'; + +export const BATCH_INTERVAL_MS = 40; + +export interface PendingStreamDelta { + messageId: string; + turnId?: string; + timestamp?: number; + snapshot?: string; + deltaText: string; + completeText?: string; +} + +export interface StreamEventPayload { + type: 'text_delta' | 'text_complete'; + messageId?: string; + text: string; + textSnapshot?: string; + turnId?: string; +} + +/** + * Accumulates a text_delta or text_complete event into the pending batch map. + * Returns the resolved message ID. + */ +export function accumulateStreamDelta( + pending: Map, + event: StreamEventPayload, + fallbackId: string, +): string { + const id = event.messageId ?? fallbackId; + const existing = pending.get(id) ?? { + messageId: id, + turnId: event.turnId, + timestamp: Date.now(), + deltaText: '', + }; + if (event.turnId) existing.turnId = event.turnId; + + if (event.type === 'text_complete') { + existing.completeText = event.text; + } else if (event.textSnapshot !== undefined) { + existing.snapshot = event.textSnapshot; + existing.deltaText = ''; + } else { + existing.deltaText += event.text; + } + + pending.set(id, existing); + return id; +} + +/** + * Applies all pending message updates to the current message list. + * Preserves object identity for unchanged messages so React.memo works optimally. + * Returns next message array (or current if no changes). + */ +export function applyPendingBatch( + current: StoredMessage[], + pending: Map, +): StoredMessage[] { + if (pending.size === 0) return current; + + let changed = false; + + const next = current.map(message => { + const update = pending.get(message.id); + if (!update) return message; + + const base = update.snapshot !== undefined ? update.snapshot : message.content; + const nextContent = update.completeText !== undefined ? update.completeText : (base + update.deltaText); + + if (nextContent === message.content && (update.turnId === undefined || update.turnId === message.turnId)) { + return message; + } + + changed = true; + return { + ...message, + content: nextContent, + turnId: update.turnId ?? message.turnId, + }; + }); + + for (const [id, update] of pending.entries()) { + if (!current.some(m => m.id === id)) { + changed = true; + const base = update.snapshot !== undefined ? update.snapshot : ''; + const content = update.completeText !== undefined ? update.completeText : (base + update.deltaText); + next.push({ + id, + type: 'assistant', + content, + timestamp: update.timestamp ?? Date.now(), + turnId: update.turnId, + }); + } + } + + return changed ? next : current; +} diff --git a/packages/ui/src/useWorkbench.ts b/packages/ui/src/useWorkbench.ts index 40f2bbf..3cdcf5d 100644 --- a/packages/ui/src/useWorkbench.ts +++ b/packages/ui/src/useWorkbench.ts @@ -3,10 +3,24 @@ import type { AgentEvent, StoredMessage } from '@threadcove/core/types'; import type { SessionDto, SourceDto } from '@threadcove/shared/protocol'; import type { ElectronAPI, TransportConnectionState } from '@threadcove/shared/client'; import type { RunFeedback } from './components/RunIndicator.tsx'; +import { + accumulateStreamDelta, + applyPendingBatch, + BATCH_INTERVAL_MS, + type PendingStreamDelta, +} from './streamBatcher.ts'; export const connectionNames: Record = { idle: '等待连接', connecting: '正在连接', connected: '已连接', reconnecting: '正在重连', disconnected: '已断开', failed: '连接失败' }; export const runNames: Record = { running: '研究中', completed: '已完成', failed: '执行失败', cancelled: '已停止', interrupted: '执行已中断' }; +export interface FilePreviewState { + name: string; + content: string; + truncated?: boolean; + totalBytes?: number; + previewBytes?: number; +} + type RunEvent = { type: 'run_status'; status: string; error?: string; runId: string }; type Incoming = { workspaceId?: string; sessionId: string; event: AgentEvent | RunEvent | { type: 'user_message'; message: StoredMessage } }; @@ -28,7 +42,7 @@ export function useWorkbench(api: ElectronAPI) { const [submitting, setSubmitting] = useState(false); const [sources, setSources] = useState([]); const [files, setFiles] = useState>([]); - const [filePreview, setFilePreview] = useState<{ name: string; content: string }>(); + const [filePreview, setFilePreview] = useState(); const [activity, setActivity] = useState([]); const [permission, setPermission] = useState['request']>(); const [runFeedback, setRunFeedback] = useState(); @@ -40,6 +54,8 @@ export function useWorkbench(api: ElectronAPI) { const selectedRef = useRef(selected); const workspaceRef = useRef(workspace); const revision = useRef(0); + const pendingDeltas = useRef>(new Map()); + const batchTimer = useRef>(); const pending = useRef<{ session: string; text: string; id: string }>(); const scroll = useRef(null); const stick = useRef(true); @@ -71,6 +87,28 @@ export function useWorkbench(api: ElectronAPI) { } }, [api, report]); + const flushBatch = useCallback(() => { + if (batchTimer.current) { + clearTimeout(batchTimer.current); + batchTimer.current = undefined; + } + if (pendingDeltas.current.size === 0) return; + const batch = pendingDeltas.current; + pendingDeltas.current = new Map(); + revision.current++; + if (mounted.current) { + setMessages(current => applyPendingBatch(current, batch)); + } + }, []); + + const cancelBatch = useCallback(() => { + if (batchTimer.current) { + clearTimeout(batchTimer.current); + batchTimer.current = undefined; + } + pendingDeltas.current = new Map(); + }, []); + useEffect(() => { mounted.current = true; let stopped = false; @@ -88,6 +126,7 @@ export function useWorkbench(api: ElectronAPI) { const off = api.onConnectionStateChanged(state => { connected.current = state === 'connected'; observedRun.current = undefined; setRunFeedback(undefined); + if (state !== 'connected') flushBatch(); setConnection(state); if (state === 'connected') void load().catch(report); }); const offEvents = api.onSessionEvent(raw => { @@ -102,22 +141,28 @@ export function useWorkbench(api: ElectronAPI) { if (sessionId !== selectedRef.current) return; revision.current++; if (event.type === 'user_message') { + flushBatch(); if (connected.current) { observedRun.current = event.message.runId; setRunFeedback({ sessionId, runId: event.message.runId ?? event.message.id, status: 'running' }); } setMessages(current => current.some(message => message.id === event.message.id) ? current : [...current, event.message]); setActivity([]); setError(''); - } else if (event.type === 'text_delta' || event.type === 'text_complete') { - const id = event.messageId ?? liveId.current ?? `live-${crypto.randomUUID()}`; - liveId.current = event.type === 'text_complete' ? undefined : id; - setMessages(current => { - const index = current.findIndex(message => message.id === id); - const previous = current[index]; - const next: StoredMessage = { id, type: 'assistant', content: event.type === 'text_complete' ? event.text : event.textSnapshot ?? (previous?.content ?? '') + event.text, timestamp: previous?.timestamp ?? Date.now(), turnId: event.turnId }; - return index === -1 ? [...current, next] : current.map((message, i) => i === index ? next : message); - }); + } else if (event.type === 'text_delta') { + const id = accumulateStreamDelta(pendingDeltas.current, event, liveId.current ?? `live-${crypto.randomUUID()}`); + liveId.current = id; + if (!batchTimer.current) { + batchTimer.current = setTimeout(() => { + batchTimer.current = undefined; + flushBatch(); + }, BATCH_INTERVAL_MS); + } + } else if (event.type === 'text_complete') { + accumulateStreamDelta(pendingDeltas.current, event, liveId.current ?? `live-${crypto.randomUUID()}`); + liveId.current = undefined; + flushBatch(); } else if (event.type === 'run_status') { + flushBatch(); if (connected.current) { const completion = event.status === 'completed' && observedRun.current === event.runId ? event.runId : undefined; // Reject a late terminal event from a previous run for presentation purposes. @@ -129,16 +174,24 @@ export function useWorkbench(api: ElectronAPI) { setPermission(undefined); if (event.error) setError(event.error); void loadMessages(sessionId).catch(report); - } else if (event.type === 'typed_error') setError(event.error.message); - else if (event.type === 'error') setError(event.message); - else if (event.type === 'permission_request') { setPermission(event.request); setDetails(true); } - else if (event.type === 'tool_start' || event.type === 'tool_result') setActivity(current => [...current.slice(-29), `${event.type === 'tool_start' ? '正在执行' : '已返回'} · ${event.toolName}`]); + } else if (event.type === 'typed_error') { + flushBatch(); + setError(event.error.message); + } else if (event.type === 'error') { + flushBatch(); + setError(event.message); + } else if (event.type === 'permission_request') { + setPermission(event.request); setDetails(true); + } else if (event.type === 'tool_start' || event.type === 'tool_result') { + setActivity(current => [...current.slice(-29), `${event.type === 'tool_start' ? '正在执行' : '已返回'} · ${event.toolName}`]); + } }); void api.ready().then(load).catch(report); - return () => { stopped = true; mounted.current = false; off(); offEvents(); }; - }, [api, loadMessages, refresh, report]); + return () => { stopped = true; mounted.current = false; cancelBatch(); off(); offEvents(); }; + }, [api, cancelBatch, flushBatch, loadMessages, refresh, report]); useEffect(() => { + cancelBatch(); selectedRef.current = selected; liveId.current = undefined; observedRun.current = undefined; setRunFeedback(undefined); @@ -149,7 +202,7 @@ export function useWorkbench(api: ElectronAPI) { setDraft(sessionStorage.getItem(`threadcove-draft:${workspace}:${selected}`) ?? ''); stick.current = true; if (selected && workspace) void loadMessages(selected).catch(report); - }, [selected, workspace, loadMessages, report]); + }, [selected, workspace, cancelBatch, loadMessages, report]); useEffect(() => { if (active?.model) setModel(active.model); }, [active?.model]); useEffect(() => { if (stick.current && scroll.current) scroll.current.scrollTop = scroll.current.scrollHeight; @@ -167,12 +220,14 @@ export function useWorkbench(api: ElectronAPI) { setDraft(value); sessionStorage.setItem(`threadcove-draft:${workspace}:${selected}`, value); } async function createTask() { + flushBatch(); if (!online || !workspace || submitting) return; setSubmitting(true); setError(''); try { const session = await api.createSession(workspace); await refresh(); setSelected(session.id); setSidebar(false); } catch (error) { report(error); } finally { setSubmitting(false); } } async function send() { + flushBatch(); const text = draft.trim(); if (!text || busy || !online) return; const submissionWorkspace = workspace; @@ -211,7 +266,13 @@ export function useWorkbench(api: ElectronAPI) { setFilePreview(undefined); try { const result = await api.readFile(requestWorkspace, requestSession, name); - if (isCurrent()) setFilePreview({ name, content: result.content }); + if (isCurrent()) setFilePreview({ + name, + content: result.content, + truncated: result.truncated, + totalBytes: result.totalBytes, + previewBytes: result.previewBytes, + }); } catch (error) { if (isCurrent()) report(error); } } async function saveModel() { @@ -229,6 +290,7 @@ export function useWorkbench(api: ElectronAPI) { } catch (error) { report(error); } } async function cancel() { + flushBatch(); try { await api.cancelProcessing(workspace, selected); } catch (error) { report(error); } } diff --git a/packages/ui/src/workbench.css b/packages/ui/src/workbench.css index db51dba..1d33b49 100644 --- a/packages/ui/src/workbench.css +++ b/packages/ui/src/workbench.css @@ -868,6 +868,12 @@ h2{ font-size:12px } +.file-preview-notice{ + margin:6px 0; + font-size:11px; + color:var(--muted); +} + .file-preview pre{ white-space:pre-wrap; overflow-wrap:anywhere diff --git a/packages/ui/tests/RunIndicator.test.tsx b/packages/ui/tests/RunIndicator.test.tsx new file mode 100644 index 0000000..5788bda --- /dev/null +++ b/packages/ui/tests/RunIndicator.test.tsx @@ -0,0 +1,69 @@ +import { describe, it, expect } from 'bun:test'; +import { + computeRunIndicatorState, + TERMINAL_STATUSES, +} from '../src/components/RunIndicator.tsx'; + +describe('RunIndicator offline and terminal states', () => { + it('preserves completed status when offline', () => { + const { state, label } = computeRunIndicatorState('completed', false, false); + expect(state).toBe('completed'); + expect(label).toBe('研究已完成'); + }); + + it('preserves cancelled status when offline', () => { + const { state, label } = computeRunIndicatorState('cancelled', false, false); + expect(state).toBe('cancelled'); + expect(label).toBe('研究已取消'); + }); + + it('preserves failed status when offline', () => { + const { state, label } = computeRunIndicatorState('failed', false, false); + expect(state).toBe('failed'); + expect(label).toBe('研究失败'); + }); + + it('preserves interrupted status when offline', () => { + const { state, label } = computeRunIndicatorState('interrupted', false, false); + expect(state).toBe('interrupted'); + expect(label).toBe('研究已中断'); + }); + + it('shows disconnected label for running research when offline', () => { + const { state, label } = computeRunIndicatorState('running', false, false); + expect(state).toBe('disconnected'); + expect(label).toBe('连接中断,执行状态待同步'); + }); + + it('shows disconnected label for waiting research when offline', () => { + const { state, label } = computeRunIndicatorState(undefined, false, true); + expect(state).toBe('disconnected'); + expect(label).toBe('连接中断,执行状态待同步'); + }); + + it('returns empty state when offline and no active research exists', () => { + const { state, label } = computeRunIndicatorState(undefined, false, false); + expect(state).toBeUndefined(); + expect(label).toBeUndefined(); + }); + + it('renders running state when online', () => { + const { state, label } = computeRunIndicatorState('running', true, false); + expect(state).toBe('running'); + expect(label).toBe('正在研究…'); + }); + + it('renders waiting state when online and permission is requested', () => { + const { state, label } = computeRunIndicatorState('running', true, true); + expect(state).toBe('waiting'); + expect(label).toBe('等待你的确认'); + }); + + it('verifies all terminal statuses are accounted for', () => { + expect(TERMINAL_STATUSES.has('completed')).toBe(true); + expect(TERMINAL_STATUSES.has('cancelled')).toBe(true); + expect(TERMINAL_STATUSES.has('failed')).toBe(true); + expect(TERMINAL_STATUSES.has('interrupted')).toBe(true); + expect(TERMINAL_STATUSES.has('running')).toBe(false); + }); +}); diff --git a/packages/ui/tests/streamBatcher.test.ts b/packages/ui/tests/streamBatcher.test.ts new file mode 100644 index 0000000..c5cda2a --- /dev/null +++ b/packages/ui/tests/streamBatcher.test.ts @@ -0,0 +1,113 @@ +import { describe, it, expect } from 'bun:test'; +import type { StoredMessage } from '@threadcove/core/types'; +import { + accumulateStreamDelta, + applyPendingBatch, + type PendingStreamDelta, + BATCH_INTERVAL_MS, +} from '../src/streamBatcher.ts'; + +describe('streamBatcher', () => { + it('accumulates deltas for a new assistant message', () => { + const pending = new Map(); + accumulateStreamDelta(pending, { type: 'text_delta', text: 'Hello' }, 'msg-1'); + accumulateStreamDelta(pending, { type: 'text_delta', text: ', world!' }, 'msg-1'); + + const result = applyPendingBatch([], pending); + expect(result).toHaveLength(1); + expect(result[0].id).toBe('msg-1'); + expect(result[0].content).toBe('Hello, world!'); + expect(result[0].type).toBe('assistant'); + }); + + it('handles snapshots overriding prior deltas', () => { + const pending = new Map(); + accumulateStreamDelta(pending, { type: 'text_delta', text: 'Stale chunk' }, 'msg-1'); + accumulateStreamDelta(pending, { type: 'text_delta', text: '', textSnapshot: 'Full snapshot text' }, 'msg-1'); + accumulateStreamDelta(pending, { type: 'text_delta', text: ' and more' }, 'msg-1'); + + const result = applyPendingBatch([], pending); + expect(result[0].content).toBe('Full snapshot text and more'); + }); + + it('batches updates across multiple message IDs in a single pass', () => { + const pending = new Map(); + accumulateStreamDelta(pending, { type: 'text_delta', text: 'Message A chunk' }, 'msg-a'); + accumulateStreamDelta(pending, { type: 'text_delta', text: 'Message B chunk' }, 'msg-b'); + accumulateStreamDelta(pending, { type: 'text_delta', text: ' more A' }, 'msg-a'); + + const initial: StoredMessage[] = [ + { id: 'user-1', type: 'user', content: 'Prompt', timestamp: 100 }, + ]; + const result = applyPendingBatch(initial, pending); + + expect(result).toHaveLength(3); + expect(result[0]).toBe(initial[0]); // Referential equality preserved + expect(result[1].id).toBe('msg-a'); + expect(result[1].content).toBe('Message A chunk more A'); + expect(result[2].id).toBe('msg-b'); + expect(result[2].content).toBe('Message B chunk'); + }); + + it('applies text_complete directly and overrides partial deltas', () => { + const pending = new Map(); + accumulateStreamDelta(pending, { type: 'text_delta', text: 'Incomplete' }, 'msg-1'); + accumulateStreamDelta(pending, { type: 'text_complete', text: 'Final complete text.' }, 'msg-1'); + + const result = applyPendingBatch([], pending); + expect(result[0].content).toBe('Final complete text.'); + }); + + it('preserves references for messages not touched by the batch', () => { + const existing: StoredMessage[] = [ + { id: 'msg-0', type: 'user', content: 'Question', timestamp: 1 }, + { id: 'msg-1', type: 'assistant', content: 'Answer 1', timestamp: 2 }, + { id: 'msg-2', type: 'assistant', content: 'Answer 2', timestamp: 3 }, + ]; + const pending = new Map(); + accumulateStreamDelta(pending, { type: 'text_delta', text: ' updated' }, 'msg-2'); + + const next = applyPendingBatch(existing, pending); + expect(next[0]).toBe(existing[0]); + expect(next[1]).toBe(existing[1]); + expect(next[2]).not.toBe(existing[2]); + expect(next[2].content).toBe('Answer 2 updated'); + }); + + it('demonstrates refresh count reduction when batching 20 stream deltas', () => { + // Simulate 20 deltas arriving within a 40ms burst + let flushes = 0; + const pending = new Map(); + let batchTimer: boolean = false; + + const onDelta = (text: string) => { + accumulateStreamDelta(pending, { type: 'text_delta', text }, 'live-msg'); + if (!batchTimer) { + batchTimer = true; + } + }; + + const flush = (current: StoredMessage[]) => { + batchTimer = false; + flushes++; + const res = applyPendingBatch(current, pending); + pending.clear(); + return res; + }; + + let state: StoredMessage[] = []; + // 20 chunks in a single batch interval + for (let i = 0; i < 20; i++) { + onDelta(` chunk-${i}`); + } + + // Single flush at the end of the batch window + state = flush(state); + + expect(flushes).toBe(1); + expect(state[0].content).toContain('chunk-0'); + expect(state[0].content).toContain('chunk-19'); + expect(BATCH_INTERVAL_MS).toBeGreaterThanOrEqual(32); + expect(BATCH_INTERVAL_MS).toBeLessThanOrEqual(50); + }); +}); diff --git a/scripts/verify-ui.cjs b/scripts/verify-ui.cjs index 1071a11..d3dee25 100644 --- a/scripts/verify-ui.cjs +++ b/scripts/verify-ui.cjs @@ -6,7 +6,7 @@ const { join, resolve } = require('node:path'); const { spawn } = require('node:child_process'); const { chromium, _electron } = require('playwright'); const root = resolve(__dirname, '..'); -const output = resolve(process.env.THREADCOVE_QA_OUTPUT || join(root, 'artifacts/qa')); +const output = resolve(process.env.THREADCOVE_QA_OUTPUT || join(root, 'artifacts/qa/stability')); mkdirSync(output, { recursive: true }); const data = mkdtempSync(join(tmpdir(), 'threadcove-ui-')); const requests = []; @@ -94,9 +94,16 @@ async function idle(page) { await page.getByRole('button', { name: '开始研究 assert(firstAnswer); await require('./verify-run-motion.cjs')(page, event => injectEvent(firstSession, event), firstAnswer); writeFileSync(join(data, 'sessions', firstSession, 'data', 'evidence.txt'), 'SESSION_A_FILE_CONTENT'); + const largePreviewData = Buffer.alloc(1200 * 1024, 65); // 1.2 MiB of 'A' + writeFileSync(join(data, 'sessions', firstSession, 'data', 'large.txt'), largePreviewData); await page.getByRole('button', { name:'执行详情' }).click(); await page.getByRole('button', { name:/evidence.txt/ }).click(); await page.getByText('SESSION_A_FILE_CONTENT', {exact:true}).waitFor(); + await page.getByRole('button', { name:/large\.txt/ }).click(); + await page.getByText('文件较大,仅显示前 256 KiB', {exact:false}).waitFor(); + const previewLen = await page.locator('.file-preview pre').evaluate(el => el.textContent.length); + assert.equal(previewLen, 256 * 1024, 'preview text must be bounded to 256KiB'); + assert.equal(readFileSync(join(data, 'sessions', firstSession, 'data', 'large.txt')).length, 1200 * 1024, 'full file on disk untouched'); await screenshot(page, 'web-desktop-conversation.png'); holdFileReads = true; await page.getByRole('button', { name:/evidence.txt/ }).click(); @@ -140,11 +147,24 @@ async function idle(page) { await page.getByRole('button', { name: '开始研究 assert(!await page.getByText('这段内容不应该在取消后出现。', { exact:false }).count()); await page.reload(); await page.getByText('慢速研究已经开始。', { exact:true }).waitFor(); + // 1. Verify terminal + offline: cancelled run remains 'cancelled' when offline await stopChild(backend); - await page.waitForFunction(() => document.querySelector('.run-indicator')?.dataset.state === 'disconnected'); + await page.waitForFunction(() => !document.querySelector('.connection i.online')); + assert.equal(await page.evaluate(() => document.querySelector('.run-indicator')?.dataset.state), 'cancelled', 'terminal state must be retained when offline'); backend = await startBackend(rpcPort, env); await page.waitForFunction(() => Boolean(document.querySelector('.connection i.online')), { timeout:20000 }); assert.equal(await page.evaluate(() => window.__runAnimations.filter(n => n === 'research-finish').length), 0, 'reconnecting must not celebrate restored history'); + + // 2. Verify running + offline: unfinished run shows 'disconnected' when offline + await send(page, '慢速断线状态验收测试'); + await page.getByText('慢速研究已经开始。', { exact:false }).waitFor(); + await page.waitForFunction(() => document.querySelector('.run-indicator')?.dataset.state === 'running'); + await stopChild(backend); + await page.waitForFunction(() => document.querySelector('.run-indicator')?.dataset.state === 'disconnected'); + assert.equal(await page.evaluate(() => document.querySelector('.run-indicator')?.textContent?.trim()), '连接中断,执行状态待同步'); + backend = await startBackend(rpcPort, env); + await page.waitForFunction(() => Boolean(document.querySelector('.connection i.online')), { timeout:20000 }); + assert.equal(await page.evaluate(() => window.__runAnimations.filter(n => n === 'research-finish').length), 0, 'reconnecting must not celebrate interrupted run'); await send(page, '我之前的口令是什么?'); await idle(page); await page.waitForFunction(() => window.__runAnimations.filter(n => n === 'research-finish').length === 1); @@ -192,26 +212,45 @@ async function idle(page) { await page.getByRole('button', { name: '开始研究 await page.locator('#qa-long-error').evaluate(el => el.remove()); await browser.close(); browser = undefined; await stopChild(backend); - const electronExe = require(require.resolve('electron', { paths:[join(root,'apps/electron')] })); - electron = await _electron.launch({ executablePath:electronExe, args:[join(root,'apps/electron')], env, timeout:25000 }); - const desktop = await electron.firstWindow(); - desktop.on('pageerror', error => faults.push(error.message)); - await desktop.getByRole('button',{name:'新建研究任务'}).waitFor({timeout:20000}); - await desktop.getByRole('button',{name:'新建研究任务'}).click(); - await desktop.locator('.task.active').waitFor(); - await send(desktop,'Electron 实际窗口发送验收'); - await idle(desktop); - await desktop.getByText('完整回答:历史与模型配置均已接入真实应用。',{exact:false}).first().waitFor(); - await screenshot(desktop,'electron-conversation.png'); - await send(desktop,'Electron 慢速取消验收'); - await desktop.getByText('慢速研究已经开始。',{exact:true}).waitFor(); - await desktop.getByRole('button',{name:'停止',exact:false}).click(); - await idle(desktop); - await desktop.reload(); - await desktop.getByText('慢速研究已经开始。',{exact:true}).waitFor(); + const isWebOnly = process.argv.includes('--web-only') || process.env.THREADCOVE_WEB_ONLY === '1'; + if (!isWebOnly) { + const electronExe = require(require.resolve('electron', { paths:[join(root,'apps/electron')] })); + electron = await _electron.launch({ executablePath:electronExe, args:[join(root,'apps/electron')], env, timeout:25000 }); + const desktop = await electron.firstWindow(); + desktop.on('pageerror', error => faults.push(error.message)); + await desktop.getByRole('button',{name:'新建研究任务'}).waitFor({timeout:20000}); + await desktop.getByRole('button',{name:'新建研究任务'}).click(); + await desktop.locator('.task.active').waitFor(); + await send(desktop,'Electron 实际窗口发送验收'); + await idle(desktop); + await desktop.getByText('完整回答:历史与模型配置均已接入真实应用。',{exact:false}).first().waitFor(); + await screenshot(desktop,'electron-conversation.png'); + await send(desktop,'Electron 慢速取消验收'); + await desktop.getByText('慢速研究已经开始。',{exact:true}).waitFor(); + await desktop.getByRole('button',{name:'停止',exact:false}).click(); + await idle(desktop); + await desktop.reload(); + await desktop.getByText('慢速研究已经开始。',{exact:true}).waitFor(); + } assert.equal(faults.length, 0, faults.join('\n')); - writeFileSync(join(output,'verification.json'),JSON.stringify({passed:true,web:true,electron:true,mobileWidth:390,requests:requests.length,modelSwitch:true,restartHistory:true,stoppedPartialSaved:true,fileRpc:true,staleFileResponseIgnored:true,taskDraftIsolation:true,archiveRestore:true,consoleErrors:faults},null,2)); - console.log(JSON.stringify({passed:true,output,requests:requests.length})); + writeFileSync(join(output,'verification.json'),JSON.stringify({ + passed:true, + web:true, + electron:!isWebOnly, + mobileWidth:390, + requests:requests.length, + modelSwitch:true, + restartHistory:true, + stoppedPartialSaved:true, + fileRpc:true, + staleFileResponseIgnored:true, + taskDraftIsolation:true, + archiveRestore:true, + fileTruncationTested:true, + offlineTerminalTested:true, + consoleErrors:faults + },null,2)); + console.log(JSON.stringify({passed:true,output,requests:requests.length,webOnly:isWebOnly})); })().catch(error => { console.error(error); process.exitCode=1; }).finally(async () => { if (browser) await browser.close().catch(()=>{}); if (electron) await electron.close().catch(()=>{});