From e24ed3b63d124429467bc7a3e53a207f2593c73e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9E=97=E5=87=AF90331?= <90331@sangfor.com> Date: Wed, 13 May 2026 23:49:20 +0800 Subject: [PATCH 1/3] feat(rawDump): add raw dump reporting module with local mode and deduplication MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Implement the raw dump data reporting module for uploading conversation, summary, and commit data to the CoStrict endpoint. Key features: - Local dump mode for debugging without network - PID-based file locking to prevent state corruption with multiple csc instances - 5-minute deduplication window for summary uploads - Auto-detect sender (user/agent) and conditional user_input - Repository fields (repo_addr, repo_branch, work_dir) in conversation - Anonymous interface fallback when auth fails - Filter commits to current user and limit fetch range - Skip empty intermediate conversation turns - Fix worker spawning for build and binary modes - Remove git diff HEAD fallback to avoid unrelated working tree changes Signed-off-by: 林凯90331 <90331@sangfor.com> Co-authored-by: CoStrict --- build.ts | 34 +++- docs/config-list.md | 13 ++ docs/user-guide.md | 34 ++++ package.json | 9 +- src/main.tsx | 3 + src/services/rawDump/README.md | 191 +++++++++++++++--- src/services/rawDump/batchWorker.ts | 59 ++++-- src/services/rawDump/git.ts | 8 +- src/services/rawDump/index.ts | 10 +- src/services/rawDump/localStorage.ts | 51 +++++ src/services/rawDump/spawn.ts | 91 +++++++-- src/services/rawDump/state.ts | 79 ++++++-- src/services/rawDump/types.ts | 13 +- src/services/rawDump/worker.ts | 283 ++++++++++++++++++--------- 14 files changed, 685 insertions(+), 193 deletions(-) create mode 100644 src/services/rawDump/localStorage.ts diff --git a/build.ts b/build.ts index cb28fd106..06654c40e 100644 --- a/build.ts +++ b/build.ts @@ -1,4 +1,4 @@ -import { readdir, readFile, writeFile, cp } from 'fs/promises' +import { readdir, readFile, writeFile, cp, unlink, mkdir } from 'fs/promises' import { join } from 'path' import { getMacroDefines, DEFAULT_BUILD_FEATURES } from './scripts/defines.ts' @@ -25,7 +25,7 @@ const envFeatures = Object.keys(process.env) .map(k => k.replace('FEATURE_', '')) const features = [...new Set([...DEFAULT_BUILD_FEATURES, ...envFeatures])] -// Step 2: Bundle with splitting +// Step 2: Bundle main entrypoint const result = await Bun.build({ entrypoints: ['src/entrypoints/cli.tsx'], outdir, @@ -42,9 +42,22 @@ const result = await Bun.build({ features, }) -if (!result.success) { +// Step 2.5: Bundle raw dump worker (standalone, no splitting) +const workerResult = await Bun.build({ + entrypoints: ['src/services/rawDump/batchWorker.ts'], + outdir, + target: 'bun', + sourcemap: 'linked', + define: { + ...getMacroDefines(), + 'process.env.NODE_ENV': JSON.stringify('production'), + }, + features, +}) + +if (!result.success || !workerResult.success) { console.error('Build failed:') - for (const log of result.logs) { + for (const log of [...result.logs, ...workerResult.logs]) { console.error(log) } process.exit(1) @@ -113,6 +126,19 @@ console.log( `Bundled ${result.outputs.length} files to ${outdir}/ (patched ${patched} for import.meta.require, ${bunPatched} for Bun destructure, ${featureReplaced} feature flags)`, ) +// Step 3.9: Move raw dump worker to expected subdir (must happen after patches, before map cleanup) +const workerSrc = join(outdir, 'batchWorker.js') +const workerDstDir = join(outdir, 'services', 'rawDump') +const workerDst = join(workerDstDir, 'batchWorker.js') +try { + await mkdir(workerDstDir, { recursive: true }) + await cp(workerSrc, workerDst) + await unlink(workerSrc) + console.log(`Moved batchWorker.js → ${workerDst}`) +} catch (err) { + console.warn('Warning: could not move batchWorker.js:', (err as Error).message) +} + // Step 4: Copy native .node addon files (audio-capture) and vendored binaries (ripgrep) const audioCaptureDir = join(outdir, 'vendor', 'audio-capture') await cp('vendor/audio-capture', audioCaptureDir, { recursive: true }) diff --git a/docs/config-list.md b/docs/config-list.md index 3071d141d..8df65ad5a 100644 --- a/docs/config-list.md +++ b/docs/config-list.md @@ -184,6 +184,19 @@ | `CLAUDE_CODE_BUBBLEWRAP` | boolean | `false` | 启用 Bubblewrap 沙箱 | | `CLAUDE_CODE_ENABLE_XAA` | boolean | `false` | 启用 XAA(外部认证代理) | +### Raw Dump 数据上报 + +| 变量名 | 类型 | 默认值 | 说明 | +|--------|------|--------|------| +| `CSC_DISABLE_RAW_DUMP` | boolean | `false` | 禁用数据上报(设为 `0` 或 `false` 表示启用) | +| `COSTRICT_DISABLE_RAW_DUMP` | boolean | `false` | 兼容 opencode 的禁用开关 | +| `CSC_RAW_DUMP_DEBUG` | boolean | `false` | 开启 raw dump 调试日志 | +| `CSC_RAW_DUMP_BASE_URL` | string | - | 自定义上报服务端地址 | +| `COSTRICT_RAW_DUMP_BASE_URL` | string | - | 兼容 opencode 的自定义上报地址 | +| `COSTRICT_BASE_URL` | string | `https://zgsm.sangfor.com` | CoStrict 服务地址 | +| `CSC_RAW_DUMP_LOCAL_MODE` | boolean | `false` | **本地留存模式**:数据只写入本地文件,不上报服务端 | +| `CSC_RAW_DUMP_LOCAL_DIR` | string | `~/.claude/raw-dump-local` | 本地留存目录 | + ### Bash/终端 | 变量名 | 类型 | 默认值 | 说明 | diff --git a/docs/user-guide.md b/docs/user-guide.md index 15f438fe2..955cbefca 100644 --- a/docs/user-guide.md +++ b/docs/user-guide.md @@ -428,4 +428,38 @@ bun run build # 构建后运行 dist/cli.js 切换主题或颜色配置。 +#### Q: 排查 raw dump 上报数据 + +开启**本地留存模式**,让 raw dump 数据只写入本地文件而不上报服务端,用于调试和排障: + +```bash +# 在启动 csc 前设置环境变量 +export CSC_RAW_DUMP_LOCAL_MODE=1 +export CSC_RAW_DUMP_LOCAL_DIR=/tmp/raw-dump-debug + +csc +``` + +留存文件按 `sessionID` 分目录存储: + +``` +/tmp/raw-dump-debug/ +└── {sessionID}/ + ├── 2026-05-12T10-30-00-conversation-msg-uuid.json + ├── 2026-05-12T10-30-01-summary-msg-uuid.json + └── 2026-05-12T10-30-02-commit-abc123.json +``` + +每个 JSON 文件包含完整的上报 payload,并在 `_dumpMeta` 字段标注类型和时间戳。本地模式特点: + +- **无需登录** — auth 失败自动降级,不会阻断流程 +- **不触发 HTTP** — 零网络依赖,零 429 风险 +- **与正常流程一致** — 队列 + worker 机制完全保留,只是输出到本地文件 + +关闭本地模式(恢复正常上报): + +```bash +unset CSC_RAW_DUMP_LOCAL_MODE +``` + --- \ No newline at end of file diff --git a/package.json b/package.json index e5d1ec370..2cff8656f 100644 --- a/package.json +++ b/package.json @@ -50,13 +50,20 @@ "build:binary:linux-musl-baseline": "bun run build && bun run compile:linux-musl-baseline", "build:binary:win": "bun run build && bun run compile:win", "build:binary:win-baseline": "bun run build && bun run compile:win-baseline", - "build:binary:all": "bun run build && bun run compile:linux && bun run compile:linux-baseline && bun run compile:linux-musl && bun run compile:linux-musl-baseline && bun run compile:win && bun run compile:win-baseline", + "build:binary:mac": "bun run build && bun run compile:mac-arm64 && bun run compile:mac-x64", + "build:binary:mac-arm64": "bun run build && bun run compile:mac-arm64", + "build:binary:mac-x64": "bun run build && bun run compile:mac-x64", + "build:binary:mac-x64-baseline": "bun run build && bun run compile:mac-x64-baseline", + "build:binary:all": "bun run build && bun run compile:linux && bun run compile:linux-baseline && bun run compile:linux-musl && bun run compile:linux-musl-baseline && bun run compile:win && bun run compile:win-baseline && bun run compile:mac-arm64 && bun run compile:mac-x64 && bun run compile:mac-x64-baseline", "compile:linux": "bun build dist/cli.js --compile --target=bun-linux-x64 --outfile dist/csc-linux-x64", "compile:linux-baseline": "bun build dist/cli.js --compile --target=bun-linux-x64-baseline --outfile dist/csc-linux-x64-baseline", "compile:linux-musl": "bun build dist/cli.js --compile --target=bun-linux-x64-musl --outfile dist/csc-linux-x64-musl", "compile:linux-musl-baseline": "bun build dist/cli.js --compile --target=bun-linux-x64-musl-baseline --outfile dist/csc-linux-x64-musl-baseline", "compile:win": "bun build dist/cli.js --compile --target=bun-windows-x64 --outfile dist/csc-windows-x64.exe", "compile:win-baseline": "bun build dist/cli.js --compile --target=bun-windows-x64-baseline --outfile dist/csc-windows-x64-baseline.exe", + "compile:mac-arm64": "bun build dist/cli.js --compile --target=bun-darwin-arm64 --outfile dist/csc-darwin-arm64", + "compile:mac-x64": "bun build dist/cli.js --compile --target=bun-darwin-x64 --outfile dist/csc-darwin-x64", + "compile:mac-x64-baseline": "bun build dist/cli.js --compile --target=bun-darwin-x64-baseline --outfile dist/csc-darwin-x64-baseline", "dev": "bun run scripts/dev.ts", "dev:inspect": "bun run scripts/dev-debug.ts", "prepublishOnly": "bun run build", diff --git a/src/main.tsx b/src/main.tsx index d8bcf23e5..bf14a30cd 100644 --- a/src/main.tsx +++ b/src/main.tsx @@ -3883,6 +3883,7 @@ async function run(): Promise { if (isBareMode()) { // skip — no-op } else if (isNonInteractiveSession) { + process.env.CSC_RAW_DUMP_CALLER = 'headless' // In headless mode, await to ensure plugin sync completes before CLI exits await initializeVersionedPlugins(); profileCheckpoint("action_after_plugins_init"); @@ -4292,6 +4293,8 @@ async function run(): Promise { return; } + process.env.CSC_RAW_DUMP_CALLER = 'chat' + // Log model config at startup logEvent("tengu_startup_manual_model_config", { cli_flag: diff --git a/src/services/rawDump/README.md b/src/services/rawDump/README.md index 6c30c0098..2f2306992 100644 --- a/src/services/rawDump/README.md +++ b/src/services/rawDump/README.md @@ -31,25 +31,62 @@ src/services/rawDump/ ## 上报流程 -``` -主进程:assistant message 完成 - → reportTurn(sessionID, messageID, directory) - → enqueue({ sessionID, messageID, directory }) 写入队列文件 - → ensureBatchWorker() 启动 detached batch worker(仅一次) +```mermaid +flowchart TD + subgraph 主进程 + A[reportTurn / reportSession] --> B{isEnabled?} + B -->|否| Z([返回]) + B -->|是| C[enqueue 写入 JSONL 队列] + C --> D{batchWorkerSpawned?} + D -->|是| Z + D -->|否| E[ensureBatchWorker] + end -Batch Worker 进程(独立,每 30s + 随机抖动检查一次): - → acquireLock() # 文件锁,防止多 worker 并发 - → readQueue() # 读取队列文件 - → dedup tasks # 同一个 session 的多个 task 只保留最新一个 - → for each task: - → auth() # 加载凭证、刷新 token - → loadSessionMessages() # 从 JSONL 加载会话消息 - → uploadConversation() → POST /raw-store/task-conversation - → uploadSummary() → POST /raw-store/task-summary - → uploadCommits() → POST /raw-store/commit(逐条更新 state) - → writeState(state) # finally 中执行,确保 state 一定写入 - → clearQueue() # 清空队列 - → releaseLock() + subgraph Worker 启动 + E --> F[spawnBatchWorker] + F -->|dev / bun| G1[bun run batchWorker.ts] + F -->|build / node| G2[node batchWorker.js] + F -->|binary 模式失败| H[startBatchWorker 内联] + end + + subgraph BatchWorker 循环 + I[startBatchWorker] --> J[setTimeout 30s] + J --> K[runBatch] + K --> L{isRunning?} + L -->|是| J + L -->|否| M[acquireLock] + M -->|失败| J + M -->|成功| N[readQueue] + N --> O{队列空?} + O -->|是| P[releaseLock] --> J + O -->|否| Q[clearQueue] + Q --> R[去重] + R --> S[processTask] + S --> T[loadSessionMessages] + T --> U[authWithFallback] + U --> V[预加载 git 信息] + V --> W[uploadConversation] + V --> X[uploadSummary] + V --> Y[uploadCommits] + W --> Z1[writeState] + X --> Z1 + Y --> Z1 + Z1 --> P + end + + subgraph 实际上报 + AA[uploadConversation] --> AB{本地模式?} + AB -->|是| AC[writeLocalDump] + AB -->|否| AD[HTTP POST 服务端] + AD --> AE{失败?} + AE -->|是 3次重试| AD + AE -->|否| AF[写入 state 去重] + end + + G1 --> I + G2 --> I + H --> I + S --> AA ``` --- @@ -83,16 +120,28 @@ reportTurn(sessionId, assistantMessage.uuid, cwd) |-----|------|------| | `task_id` | `sessionID` | 会话唯一标识 | | `request_id` | `message.id` 或 `message.uuid` | assistant message ID | +| `prompt_mode` | `user.variant` | 用户消息变体(如 `normal` / `plan`) | +| `mode` | `assistant.mode` / `assistant.agent` | 默认 `"code"` | | `model` | `assistant.message.model` | 使用的模型 | -| `mode` | `assistant.mode` / `assistant.agent` | 默认 "code" | | `start_time` | parent user message `timestamp` | 用户请求时间 | | `end_time` | assistant message `timestamp` | assistant 完成时间 | +| `process_time` | `end_time - start_time` | 本轮处理耗时(毫秒) | +| `process_ttft` | `assistant.ttftMs` | 首 token 延迟(毫秒) | | `upstream_tokens` | `usage.input + cache_read + cache_creation` | 输入 token 总量 | | `downstream_tokens` | `usage.output` | 输出 token 量 | +| `cost` | 固定 `0`(待接入 cost-tracker) | 本轮调用成本(USD) | +| `sender` | `detectSender(assistant, user)` | `"user"` 或 `"agent"`,根据消息来源自动识别 | | `request_content` | user message text content | 用户请求文本 | | `response_content` | assistant text content | assistant 回复文本 | +| `user_input` | `sender === 'user'` 时同 `request_content`,否则为空 | 真实用户输入文本 | | `diff` | tool_use diff → fallback `git diff HEAD` | 本轮代码变更 | -| `error_code` | error name 映射 | 401/413/499/500 | +| `diff_lines` | diff 中 `+` 行数统计 | 新增/变更行数 | +| `files` | diff 中涉及文件列表 | 本轮变更的文件列表 | +| `repo_addr` | `git remote get-url origin` | 仓库远程地址 | +| `repo_branch` | `git branch --show-current` | 当前分支 | +| `work_dir` | 当前工作目录 | 项目路径 | +| `error_code` | error name 映射 | `401` / `413` / `499` / `500` 或 API 原始状态码 | +| `error_reason` | `error.message` | 错误原因描述 | ### Summary(会话汇总) @@ -101,21 +150,36 @@ reportTurn(sessionId, assistantMessage.uuid, cwd) | `task_id` | `sessionID` | 会话唯一标识 | | `start_time` | 第一条消息 `timestamp` | 会话开始时间 | | `end_time` | 最后一条消息 `timestamp` | 会话最后更新时间 | -| `upstream_tokens` | 所有 assistant messages 累计 | 会话总输入 token | -| `downstream_tokens` | 所有 assistant messages 累计 | 会话总输出 token | | `user_id` | refresh_token JWT `universal_id` | 用户唯一标识 | -| `repo_addr` | `git remote get-url origin` | 仓库地址 | -| `repo_branch` | `git branch --show-current` | 当前分支 | -| `diff` | `git diff HEAD` | 工作区完整变更 | +| `user_name` | refresh_token JWT `properties.oauth_GitHub_username` 或 `id` | 当前登录用户名 | +| `client_id` | `creds.machine_id` / `CSC_MACHINE_ID` | 设备唯一标识 | +| `client_version` | `package.json` version | CLI 版本号 | +| `client_ide` | 固定值 `"cli"` | 客户端类型 | +| `client_os` | `os.platform()` | 操作系统 | +| `client_os_version` | `os.release()` | 操作系统版本 | +| `caller` | `"chat"`(REPL 交互模式)/ `"headless"`(`--print`/管道模式) | 调用来源,自动识别 | ### Commits(Git 提交) | 字段 | 来源 | 说明 | |-----|------|------| -| `commit_id` | `git log` | commit hash | +| `commit_id` | `git log %H` | commit hash | | `commit_time` | `git log %aI` | 作者时间(ISO) | +| `repo_addr` | `git remote get-url origin` | 仓库远程地址 | +| `repo_branch` | `git branch --show-current` | 当前分支 | +| `git_user_name` | `git log %an` | commit 作者姓名 | +| `git_user_email` | `git log %ae` | commit 作者邮箱 | +| `user_id` | refresh_token JWT `universal_id` | 当前登录用户 ID | +| `user_name` | refresh_token JWT `properties.oauth_GitHub_username` 或 `id` | 当前登录用户名 | +| `client_id` | `creds.machine_id` / `CSC_MACHINE_ID` | 设备唯一标识 | +| `client_version` | `package.json` version | CLI 版本号 | +| `client_ide` | 固定值 `"cli"` | 客户端类型 | +| `work_dir` | 当前工作目录 | 项目路径 | | `diff` | `git show --diff-filter=ACDMR` | 变更内容 | +| `diff_lines` | diff 中 `+` 行数 | 新增行数统计 | +| `files` | diff 中涉及文件列表 | 变更文件列表 | | `comment` | `subject.slice(0, 150)` | 截断后的提交信息 | +| `subject` | `git log %s` | 原始提交信息(完整) | --- @@ -124,8 +188,8 @@ reportTurn(sessionId, assistantMessage.uuid, cwd) csc 没有 opencode 中的 `step-start`/`step-finish` snapshot 机制,采用以下策略: ### Conversation diff -1. **优先**:从 assistant message 的 `tool_use` blocks 中提取 `input.content` / `new_string` / `diff` / `patch` -2. **Fallback**:执行 `git diff HEAD` 获取当前工作区未提交的变更 +- **唯一来源**:从当前 assistant message 的 `tool_use` blocks 中提取 `input.content` / `new_string` / `diff` / `patch` +- **无 fallback**:不执行 `git diff HEAD`。工作区中历史未提交的改动与当前对话轮次无关,不应混入 conversation diff。 ### Summary diff - 直接执行 `git diff HEAD`,获取整个工作区相对于最新 commit 的变更 @@ -157,7 +221,18 @@ if (!existing || task.enqueuedAt > existing.enqueuedAt) { } ``` -### 3. Commits 去重(磁盘,逐条更新) +### 3. Summary 去重(磁盘,时间窗口) +```typescript +// 以 sessionID 为 key,记录上次上报时间戳 +{ + "summary": { + "session-id-1": 1747123456789 + } +} +``` +- **时间窗口**:同一 session 的 summary 在 **5 分钟内只上报一次**。长会话跨窗口后才会再次上报,既避免批量消费时的重复,又保留会话更新的机会。 + +### 4. Commits 去重(磁盘,逐条更新) ```typescript // 以 repo#branch#workDir 为 key { @@ -168,8 +243,8 @@ if (!existing || task.enqueuedAt > existing.enqueuedAt) { ``` - **逐 commit 更新**:每成功上传一个 commit,立即更新 `state.commits[stateKey]` 为该 commit 的 hash。即使后续失败,已成功的 commits 不会重复上报。 - **获取范围**: - - 有 lastCommit:`git log ${lastCommit}..HEAD --max-count=50` - - 无 lastCommit:`git log --since=7 days ago --max-count=50` + - 有 lastCommit:`git log ${lastCommit}..HEAD --max-count=50 --author $(git config user.email)` + - 无 lastCommit:`git log --since=1 day ago --max-count=50 --author $(git config user.email)` - **批次延迟**:每上传 10 个 commits 后暂停 500ms,避免触发限流 --- @@ -214,6 +289,50 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js' | `CSC_RAW_DUMP_BASE_URL` | 自定义上报 base URL | 从凭证读取 | | `COSTRICT_RAW_DUMP_BASE_URL` | 兼容 opencode 的自定义 URL | 从凭证读取 | | `COSTRICT_BASE_URL` | CoStrict 服务地址 | `https://zgsm.sangfor.com` | +| `CSC_RAW_DUMP_LOCAL_MODE` | 本地留存模式,数据只写入本地文件不上报服务端 | `false` | +| `CSC_RAW_DUMP_LOCAL_DIR` | 本地留存目录 | `~/.claude/raw-dump-local` | + +--- + +## 本地留存模式(调试排障) + +通过环境变量开启,开启后数据**不上报服务端**,仅写入本地 JSON 文件: + +```bash +# 开启本地留存模式 +export CSC_RAW_DUMP_LOCAL_MODE=1 + +# 可选:自定义留存目录(默认 ~/.claude/raw-dump-local) +export CSC_RAW_DUMP_LOCAL_DIR=/tmp/raw-dump-debug +``` + +留存文件结构: +``` +{localDir}/ +└── {sessionID}/ + ├── 2026-05-12T10-30-00-conversation-msg-uuid.json + ├── 2026-05-12T10-30-01-summary-msg-uuid.json + └── 2026-05-12T10-30-02-commit-abc123.json +``` + +每个 JSON 文件包含完整的上报 payload,并在 `_dumpMeta` 字段中标注类型和时间戳: +```json +{ + "_dumpMeta": { + "type": "conversation", + "dumpedAt": "2026-05-12T10:30:00.000Z", + "endpoint": "/raw-store/task-conversation" + }, + "task_id": "...", + "request_id": "..." +} +``` + +**本地模式特点:** +- 无需登录(auth 失败自动降级,使用 `local-mode` 占位值) +- 不依赖网络,不触发任何 HTTP 请求 +- 零 429 / 限流风险 +- 与正常队列机制完全一致,只是最终输出目的地从服务端改为本地文件 --- @@ -231,6 +350,9 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js' "session-id-1:msg-uuid-1": true, "session-id-1:msg-uuid-2": true }, + "summary": { + "session-id-1": 1747123456789 + }, "commits": { "git@github.com:org/repo.git#main#/Users/xxx/code/repo": "abc123def" } @@ -264,7 +386,12 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js' `model` 字段取自 `assistant.message.model`。若该字段不可靠,可从 `bootstrap/state.ts` 的 `getCurrentModel()` 获取。 6. **Sender 识别** - 当前固定为 `"user"`。若 csc 支持 agent/agentic 模式,需根据消息来源判断 `"user"` 或 `"agent"`。 + 已实现自动识别: + - `assistant.mode === 'agent'` / `'auto'` → `"agent"` + - `assistant.agent` 存在 → `"agent"` + - `assistant.isSidechain === true` → `"agent"` + - 父 user 消息 `isMeta === true` → `"agent"` + - 其他情况 → `"user"` --- diff --git a/src/services/rawDump/batchWorker.ts b/src/services/rawDump/batchWorker.ts index 070e8b8e8..561089cae 100644 --- a/src/services/rawDump/batchWorker.ts +++ b/src/services/rawDump/batchWorker.ts @@ -4,11 +4,11 @@ * 独立进程,通过自循环 setTimeout 严格串行执行 */ -import { uploadConversation, uploadSummary, uploadCommits, auth } from './worker.js' +import { uploadConversation, uploadSummary, uploadCommits, authWithFallback } from './worker.js' import { readQueue, clearQueue, acquireLock, releaseLock, type QueueTask } from './queue.js' import { readState, writeState } from './state.js' import { getSessionDirectory, loadSessionMessages } from './worker.js' -import { getRepoInfo, getWorkingTreeDiff } from './git.js' +import { getRepoInfo } from './git.js' import { createLogger } from './logger.js' const log = createLogger('raw-dump-batch') @@ -17,22 +17,34 @@ const BATCH_INTERVAL_MS = 30_000 // 每轮间隔 // 进程内重入保护:文件锁不防同进程重入,必须用内存 flag 兜底 let isRunning = false +const PARENT_PID = process.ppid +const IS_WORKER_PROCESS = process.argv[1]?.includes('batchWorker') || false + +function isParentAlive(): boolean { + if (!IS_WORKER_PROCESS) return true + try { + process.kill(PARENT_PID, 0) + return true + } catch { + return false + } +} + async function processTask(task: QueueTask) { - log('info', 'processing task', { sessionID: task.sessionID, messageID: task.messageID }) + log.info('processing task', { sessionID: task.sessionID, messageID: task.messageID }) const sessionDir = getSessionDirectory(task.directory, task.sessionID) const messages = await loadSessionMessages(sessionDir, task.sessionID, task.messageID) if (messages.length === 0) { - log('warn', 'no messages found', { sessionDir, sessionID: task.sessionID }) + log.warn('no messages found', { sessionDir, sessionID: task.sessionID }) } - const authData = await auth() + const authData = await authWithFallback() const state = await readState() - // 预加载 git 信息,三次上传共享,避免每个 task 重复 spawn 8+ 个 git 进程 + // 预加载 git 信息,commits 和 repo 字段共享,避免每个 task 重复 spawn git 进程 const repoInfo = await getRepoInfo(task.directory) - const workingTreeDiff = await getWorkingTreeDiff(task.directory) try { // conversation @@ -40,20 +52,20 @@ async function processTask(task: QueueTask) { { sessionID: task.sessionID, messageID: task.messageID, directory: task.directory, messages }, authData, state, - { workingTreeDiff }, + { repoInfo }, ) - // summary(每个 turn 都报,但内容会累积) + // summary(5 分钟内同一 session 只上报一次) await uploadSummary( { sessionID: task.sessionID, directory: task.directory, messages }, authData, - { repoInfo, workingTreeDiff }, + state, ) // commits(限制频率,避免重复上报) await uploadCommits({ directory: task.directory }, authData, state, { repoInfo }) - log('info', 'task completed', { sessionID: task.sessionID, conversationUploaded }) + log.info('task completed', { sessionID: task.sessionID, conversationUploaded }) } finally { // 无论成功或失败,都写入 state(commits 已逐条更新) await writeState(state) @@ -63,7 +75,7 @@ async function processTask(task: QueueTask) { async function runBatch() { // 第一道防线:同进程重入保护 if (isRunning) { - log('debug', 'runBatch already running in-process, skip') + log.debug('runBatch already running in-process, skip') return } isRunning = true @@ -71,14 +83,14 @@ async function runBatch() { try { // 第二道防线:跨进程文件锁 if (!acquireLock()) { - log('debug', 'another worker process holds the lock, skip') + log.debug('another worker process holds the lock, skip') return } try { const tasks = readQueue() if (tasks.length === 0) { - log('debug', 'queue empty') + log.debug('queue empty') return } @@ -87,7 +99,7 @@ async function runBatch() { // - 即使有意外的并发 runBatch 拿到锁,也只会看到空队列直接返回 clearQueue() - log('info', `processing ${tasks.length} tasks`) + log.info(`processing ${tasks.length} tasks`) // 去重:同一个 session 的多个 task,只保留最新的一个 const deduped = new Map() @@ -100,20 +112,20 @@ async function runBatch() { } const uniqueTasks = Array.from(deduped.values()).sort((a, b) => a.enqueuedAt - b.enqueuedAt) - log('info', `deduped to ${uniqueTasks.length} unique tasks`) + log.info(`deduped to ${uniqueTasks.length} unique tasks`) for (const task of uniqueTasks) { try { await processTask(task) } catch (err) { - log('error', 'task failed', { + log.error('task failed', { error: err instanceof Error ? err.message : String(err), sessionID: task.sessionID, }) } } - log('info', 'batch completed') + log.info('batch completed') } finally { releaseLock() } @@ -123,16 +135,20 @@ async function runBatch() { } export function startBatchWorker() { - log('info', 'batch worker started', { interval: BATCH_INTERVAL_MS }) + log.info('batch worker started', { interval: BATCH_INTERVAL_MS }) // 自循环 setTimeout:上一轮跑完才安排下一轮,从源头消除并发 // 即便 runBatch 抛错也确保下一轮被排上,避免 worker 卡死 const scheduleNext = (delay: number) => { setTimeout(async () => { + if (!isParentAlive()) { + log.info('parent process exited, stopping batch worker') + process.exit(0) + } try { await runBatch() } catch (err) { - log('error', 'runBatch threw', { error: err instanceof Error ? err.message : String(err) }) + log.error('runBatch threw', { error: err instanceof Error ? err.message : String(err) }) } const jitter = Math.floor(Math.random() * 5_000) scheduleNext(BATCH_INTERVAL_MS + jitter) @@ -144,6 +160,7 @@ export function startBatchWorker() { } // 如果直接运行此文件 -if (process.argv[1]?.includes('batchWorker')) { +const scriptPath = process.argv[1] || '' +if (scriptPath.endsWith('batchWorker.ts') || scriptPath.endsWith('batchWorker.js')) { startBatchWorker() } diff --git a/src/services/rawDump/git.ts b/src/services/rawDump/git.ts index 92c8e4b38..ed351726b 100644 --- a/src/services/rawDump/git.ts +++ b/src/services/rawDump/git.ts @@ -14,6 +14,7 @@ async function gitExec(args: string[], cwd: string): Promise { cwd, encoding: 'utf-8', maxBuffer: 50 * 1024 * 1024, // 50MB + windowsHide: true, }) return stdout.trim() } catch { @@ -90,14 +91,17 @@ export function parseCommitLog(output: string): Array<{ } export async function getCommitLog(cwd: string, lastCommit?: string): Promise { + const authorEmail = await gitExec(['config', 'user.email'], cwd) + const authorFilter = authorEmail ? ['--author', authorEmail] : [] + if (lastCommit) { return gitExec( - ['log', `${lastCommit}..HEAD`, '--max-count=50', '--format=%H|%aI|%an|%ae|%s'], + ['log', `${lastCommit}..HEAD`, '--max-count=50', ...authorFilter, '--format=%H|%aI|%an|%ae|%s'], cwd, ) } return gitExec( - ['log', '--since=7 days ago', '--max-count=50', '--format=%H|%aI|%an|%ae|%s'], + ['log', '--since=1 day ago', '--max-count=50', ...authorFilter, '--format=%H|%aI|%an|%ae|%s'], cwd, ) } diff --git a/src/services/rawDump/index.ts b/src/services/rawDump/index.ts index 71c956fff..85b55f571 100644 --- a/src/services/rawDump/index.ts +++ b/src/services/rawDump/index.ts @@ -3,12 +3,16 @@ * 队列模式:主进程只 enqueue,单 batch worker 顺序消费 */ +import { isLocalDumpMode } from './localStorage.js' import { enqueue } from './queue.js' import { spawnBatchWorker } from './spawn.js' +import { startBatchWorker } from './batchWorker.js' let batchWorkerSpawned = false function isEnabled(): boolean { + // 本地调试模式自动启用 + if (isLocalDumpMode()) return true // 默认禁用 raw dump,除非显式设置为 0/false if (process.env.CSC_DISABLE_RAW_DUMP === '0' || process.env.CSC_DISABLE_RAW_DUMP === 'false') return true if (process.env.COSTRICT_DISABLE_RAW_DUMP === '0' || process.env.COSTRICT_DISABLE_RAW_DUMP === 'false') return true @@ -18,7 +22,11 @@ function isEnabled(): boolean { function ensureBatchWorker() { if (batchWorkerSpawned) return batchWorkerSpawned = true - spawnBatchWorker() + const spawned = spawnBatchWorker() + if (!spawned) { + // Fallback: compiled binary mode where external worker file is not available + startBatchWorker() + } } /** diff --git a/src/services/rawDump/localStorage.ts b/src/services/rawDump/localStorage.ts new file mode 100644 index 000000000..37682d8f3 --- /dev/null +++ b/src/services/rawDump/localStorage.ts @@ -0,0 +1,51 @@ +/** + * Raw Dump 本地存储模式 + * 开启后上报数据不落服务端,仅写入本地 JSON 文件,用于排障和调试 + */ + +import { promises as fs } from 'node:fs' +import os from 'node:os' +import path from 'node:path' + +const DEFAULT_LOCAL_DIR = path.join(os.homedir(), '.claude', 'raw-dump-local') + +export function getLocalDumpDir(): string { + return (process.env.CSC_RAW_DUMP_LOCAL_DIR || DEFAULT_LOCAL_DIR).replace(/\/$/, '') +} + +export function isLocalDumpMode(): boolean { + const mode = process.env.CSC_RAW_DUMP_LOCAL_MODE + return mode === '1' || mode === 'true' +} + +export async function writeLocalDump( + type: 'conversation' | 'summary' | 'commit', + body: Record, +): Promise { + const dir = getLocalDumpDir() + const taskId = (body.task_id as string) || 'unknown' + const taskDir = path.join(dir, taskId) + await fs.mkdir(taskDir, { recursive: true }) + + const timestamp = new Date().toISOString().replace(/[:.]/g, '-').slice(0, -5) + const requestId = + (body.request_id as string) || (body.commit_id as string) || 'unknown' + const filename = `${timestamp}-${type}-${requestId}.json` + const filePath = path.join(taskDir, filename) + + const payload = { + _dumpMeta: { + type, + dumpedAt: new Date().toISOString(), + endpoint: + type === 'conversation' + ? '/raw-store/task-conversation' + : type === 'summary' + ? '/raw-store/task-summary' + : '/raw-store/commit', + }, + ...body, + } + + await fs.writeFile(filePath, JSON.stringify(payload, null, 2) + '\n', 'utf-8') +} diff --git a/src/services/rawDump/spawn.ts b/src/services/rawDump/spawn.ts index 3039bc057..095180e53 100644 --- a/src/services/rawDump/spawn.ts +++ b/src/services/rawDump/spawn.ts @@ -4,29 +4,82 @@ */ import { spawn } from 'node:child_process' +import { existsSync } from 'node:fs' import path from 'node:path' import { fileURLToPath } from 'node:url' -export function spawnBatchWorker(): void { - const entry = process.execPath - const isDev = path.basename(entry).toLowerCase().startsWith('bun') +function resolveWorkerPath(): string { + const entry = process.execPath + const isDev = path.basename(entry).toLowerCase().startsWith('bun') + const __dirname = path.dirname(fileURLToPath(import.meta.url)) - const __dirname = path.dirname(fileURLToPath(import.meta.url)) - const workerPath = path.resolve(__dirname, 'batchWorker.ts') + if (isDev) { + return path.resolve(__dirname, 'batchWorker.ts') + } - const args = isDev - ? ['run', workerPath] - : [workerPath] + // Build mode: locate dist root (same pattern as ripgrep.ts / audio-capture-napi) + // Bun.build strips the 'src/' prefix, so worker lands at dist/services/rawDump/batchWorker.js + const parts = __dirname.split(path.sep) + const distIdx = parts.lastIndexOf('dist') + if (distIdx !== -1) { + const distRoot = parts.slice(0, distIdx + 1).join(path.sep) + return path.resolve(distRoot, 'services', 'rawDump', 'batchWorker.js') + } - const child = spawn(entry, args, { - detached: true, - windowsHide: true, - stdio: 'ignore', - }) - - child.on('error', (err) => { - console.error('[raw-dump] batch worker spawn error:', err.message) - }) - - child.unref() + return path.resolve(__dirname, 'batchWorker.js') +} + +function resolveRuntime(): { entry: string; isBun: boolean } | null { + const execPath = process.execPath + const basename = path.basename(execPath).toLowerCase() + + if (basename.startsWith('bun')) { + return { entry: execPath, isBun: true } + } + if (basename.startsWith('node')) { + return { entry: execPath, isBun: false } + } + + // Compiled binary (e.g. csc-darwin-arm64): find bun or node on PATH + if (typeof Bun !== 'undefined' && Bun.which) { + const bun = Bun.which('bun') + if (bun) return { entry: bun, isBun: true } + } + + return null +} + +/** + * 尝试 spawn 独立的 batch worker 进程 + * @returns 是否成功 spawn(false 表示应 fallback 到内联启动) + */ +export function spawnBatchWorker(): boolean { + const workerPath = resolveWorkerPath() + const runtime = resolveRuntime() + + // Worker file missing or no suitable runtime (compiled binary without external worker file) + if (!existsSync(workerPath) || !runtime) { + return false + } + + const args = runtime.isBun + ? ['run', workerPath] + : [workerPath] + + try { + const child = spawn(runtime.entry, args, { + detached: true, + windowsHide: true, + stdio: 'ignore', + }) + + child.on('error', (err) => { + console.error('[raw-dump] batch worker spawn error:', err.message) + }) + + child.unref() + return true + } catch { + return false + } } diff --git a/src/services/rawDump/state.ts b/src/services/rawDump/state.ts index 1a01d6e09..64ecad61f 100644 --- a/src/services/rawDump/state.ts +++ b/src/services/rawDump/state.ts @@ -1,37 +1,92 @@ /** * Raw Dump 磁盘状态管理 - * 用于 conversation 和 commits 的去重 + * 用于 conversation、summary 和 commits 的去重 + * 通过文件锁保证多进程并发读写安全 */ -import { promises as fs } from 'node:fs' +import { promises as fs, readFileSync, writeFileSync } from 'node:fs' import os from 'node:os' import path from 'node:path' import type { RawDumpState } from './types.js' const STATE_DIR = path.join(os.homedir(), '.claude') const STATE_FILE = path.join(STATE_DIR, 'csc-raw-dump-state.json') +const STATE_LOCK_FILE = path.join(STATE_DIR, 'csc-raw-dump-state.lock') function createEmptyState(): RawDumpState { return { conversation: {}, + summary: {}, commits: {}, } } -export async function readState(): Promise { +function acquireStateLock(): boolean { try { - const text = await fs.readFile(STATE_FILE, 'utf-8') - const parsed = JSON.parse(text) as Partial - return { - conversation: parsed.conversation ?? {}, - commits: parsed.commits ?? {}, + try { + const stat = readFileSync(STATE_LOCK_FILE, 'utf-8') + const pid = parseInt(stat, 10) + if (!isNaN(pid) && pid !== process.pid) { + try { + process.kill(pid, 0) + return false // 其他进程持有锁 + } catch { + // 进程已退出,锁是陈旧的,可以抢占 + } + } + } catch { + // 锁文件不存在 } + writeFileSync(STATE_LOCK_FILE, String(process.pid), 'utf-8') + return true } catch { - return createEmptyState() + return false } } -export async function writeState(state: RawDumpState): Promise { - await fs.mkdir(STATE_DIR, { recursive: true }) - await fs.writeFile(STATE_FILE, JSON.stringify(state, null, 2), 'utf-8') +function releaseStateLock(): void { + try { + writeFileSync(STATE_LOCK_FILE, '', 'utf-8') + } catch { + // ignore + } +} + +async function withStateLock(fn: () => Promise): Promise { + const start = Date.now() + while (!acquireStateLock()) { + if (Date.now() - start > 5_000) { + // 5 秒超时:降级为无锁执行,避免永久挂起 + break + } + await new Promise((r) => setTimeout(r, 10)) + } + try { + return await fn() + } finally { + releaseStateLock() + } +} + +export async function readState(): Promise { + return withStateLock(async () => { + try { + const text = await fs.readFile(STATE_FILE, 'utf-8') + const parsed = JSON.parse(text) as Partial + return { + conversation: parsed.conversation ?? {}, + summary: parsed.summary ?? {}, + commits: parsed.commits ?? {}, + } + } catch { + return createEmptyState() + } + }) +} + +export async function writeState(state: RawDumpState): Promise { + return withStateLock(async () => { + await fs.mkdir(STATE_DIR, { recursive: true }) + await fs.writeFile(STATE_FILE, JSON.stringify(state, null, 2), 'utf-8') + }) } diff --git a/src/services/rawDump/types.ts b/src/services/rawDump/types.ts index 5dafc4133..cf1c22f82 100644 --- a/src/services/rawDump/types.ts +++ b/src/services/rawDump/types.ts @@ -13,6 +13,7 @@ export interface RawDumpEventPayload { export interface RawDumpState { conversation: Record + summary: Record commits: Record } @@ -47,6 +48,9 @@ export interface ConversationPayload { diff: string diff_lines: number files: string[] + repo_addr: string + repo_branch: string + work_dir: string error_code?: number error_reason?: string } @@ -63,15 +67,6 @@ export interface SummaryPayload { client_os: string client_os_version: string caller: string - repo_addr: string - repo_branch: string - work_dir: string - upstream_tokens: number - downstream_tokens: number - cost: number - diff: string - diff_lines: number - files: string[] } export interface CommitPayload { diff --git a/src/services/rawDump/worker.ts b/src/services/rawDump/worker.ts index d98f824ad..aec15f59e 100644 --- a/src/services/rawDump/worker.ts +++ b/src/services/rawDump/worker.ts @@ -30,6 +30,7 @@ import { toCommitComment, } from './git.js' import { createLogger } from './logger.js' +import { isLocalDumpMode, writeLocalDump } from './localStorage.js' import { readState, writeState } from './state.js' import { RAW_DUMP_EVENT_ENV_KEY, type RawDumpEventPayload } from './types.js' import type { @@ -68,25 +69,44 @@ function resolveRawDumpBaseUrl(baseUrl?: string): string { return raw.replace(/\/cloud-api$/, '') } -function getRawDumpUrl(baseUrl: string, endpoint: string): string { +function getRawDumpUrl(baseUrl: string, endpoint: string, isAnonymous: boolean = false): string { const suffix = endpoint.startsWith('/') ? endpoint : `/${endpoint}` - return `${baseUrl}/user-indicator/api/v1${suffix}` + const prefix = isAnonymous ? '/user-indicator/public' : '/user-indicator/api/v1' + return `${baseUrl}${prefix}${suffix}` } async function postJson( baseUrl: string, headers: Headers, endpoint: string, - body: Record, + body: object, ): Promise { - const url = getRawDumpUrl(baseUrl, endpoint) - log('debug', `POST ${endpoint}`, { url }) + if (isLocalDumpMode()) { + const type = + endpoint === '/raw-store/task-conversation' + ? 'conversation' + : endpoint === '/raw-store/task-summary' + ? 'summary' + : 'commit' + await writeLocalDump(type, body as Record) + const b = body as Record + log.info(`local dump: ${type} saved`, { + task_id: b.task_id, + request_id: b.request_id, + commit_id: b.commit_id, + }) + return + } + + const isAnonymous = !headers.get('Authorization') + const url = getRawDumpUrl(baseUrl, endpoint, isAnonymous) + log.debug(`POST ${endpoint}`, { url, isAnonymous }) let lastError: Error | undefined for (let attempt = 0; attempt < 3; attempt++) { if (attempt > 0) { const delay = 5000 * Math.pow(2, attempt - 1) // 5s, 10s - log('debug', `retrying ${endpoint} after ${delay}ms`, { attempt }) + log.debug(`retrying ${endpoint} after ${delay}ms`, { attempt }) await new Promise((r) => setTimeout(r, delay)) } @@ -101,14 +121,14 @@ async function postJson( }) if (res.ok) { - log('debug', `POST ${endpoint} ok`, { status: res.status }) + log.debug(`POST ${endpoint} ok`, { status: res.status }) return } const text = await res.text().catch(() => '') // 429 限流时重试,其他错误直接抛 if (res.status === 429) { - log('warn', `${endpoint} got 429, will retry`, { attempt, text: text.slice(0, 200) }) + log.warn(`${endpoint} got 429, will retry`, { attempt, text: text.slice(0, 200) }) lastError = new Error(`${endpoint} failed: ${res.status} ${text}`) continue } @@ -117,7 +137,7 @@ async function postJson( lastError = err instanceof Error ? err : new Error(String(err)) const isAbort = lastError.name === 'AbortError' // 网络错误 / 超时也重试 - log('warn', `${endpoint} ${isAbort ? 'timeout' : 'network error'}, will retry`, { + log.warn(`${endpoint} ${isAbort ? 'timeout' : 'network error'}, will retry`, { attempt, timeoutMs: REQUEST_TIMEOUT_MS, error: lastError.message, @@ -149,14 +169,14 @@ function detectOs(): string { } export async function auth() { - log('debug', 'auth start') + log.debug('auth start') let creds = await loadCoStrictCredentials() if (!creds?.access_token) throw new Error('Not authenticated') - log('debug', 'credentials loaded', { hasRefreshToken: !!creds.refresh_token, baseUrl: creds.base_url }) + log.debug('credentials loaded', { hasRefreshToken: !!creds.refresh_token, baseUrl: creds.base_url }) // Token 刷新 if (creds.refresh_token && !isCoStrictTokenValid(creds)) { - log('debug', 'token expired, refreshing...') + log.debug('token expired, refreshing...') const next = await refreshCoStrictToken({ baseUrl: creds.base_url, refreshToken: creds.refresh_token, @@ -171,7 +191,7 @@ export async function auth() { expired_at: new Date(extractExpiryFromJWT(next.access_token)).toISOString(), }) creds = { ...creds, access_token: next.access_token, refresh_token: next.refresh_token } - log('debug', 'token refreshed') + log.debug('token refreshed') } const headers = new Headers() @@ -209,7 +229,7 @@ export async function auth() { const user = parseUser(accessPayload, refreshPayload) const baseUrl = resolveRawDumpBaseUrl(creds.base_url) - log('debug', 'auth success', { baseUrl, user_id: user.user_id, clientId, version }) + log.debug('auth success', { baseUrl, user_id: user.user_id, clientId, version }) return { baseUrl, @@ -226,7 +246,7 @@ export async function loadSessionMessages(sessionDir: string, sessionId: string, try { const entries = await fs.readdir(sessionDir) const jsonlFiles = entries.filter((f) => f.endsWith('.jsonl')) - log('debug', 'found jsonl files', { sessionDir, count: jsonlFiles.length, files: jsonlFiles.slice(0, 5) }) + log.debug('found jsonl files', { sessionDir, count: jsonlFiles.length, files: jsonlFiles.slice(0, 5) }) for (const file of jsonlFiles) { const filePath = path.join(sessionDir, file) @@ -250,7 +270,7 @@ export async function loadSessionMessages(sessionDir: string, sessionId: string, ) const hasMessage = messageId ? lines.some((m) => m.uuid === messageId || (m.message as Record)?.id === messageId) : false if (hasSession || hasMessage) { - log('debug', 'loaded messages from file', { file, count: lines.length, hasSession, hasMessage }) + log.debug('loaded messages from file', { file, count: lines.length, hasSession, hasMessage }) return lines } } catch { @@ -283,6 +303,26 @@ function findParentUserMessage( return undefined } +function detectSender( + assistant: Record, + user: Record | undefined, +): string { + // 1. assistant 消息自身标记了 agent 模式 + const mode = String(assistant.mode ?? '') + if (mode === 'agent' || mode === 'auto') return 'agent' + + // 2. assistant.agent 字段存在且非空 + if (assistant.agent) return 'agent' + + // 3. 子 agent 会话(isSidechain 为 true) + if (assistant.isSidechain === true) return 'agent' + + // 4. 父 user 消息是 meta/system 生成的(非真实用户输入) + if (user?.isMeta === true) return 'agent' + + return 'user' +} + function extractTextContent(msg: Record): string { const content = (msg.message as Record)?.content if (!Array.isArray(content)) return String(content ?? '') @@ -357,51 +397,63 @@ export async function uploadConversation( }, authData: Awaited>, state: Awaited>, - options?: { workingTreeDiff?: string }, + options?: { repoInfo?: RepoInfo }, ): Promise { - log('debug', 'uploadConversation start', { messageID: payload.messageID, messageCount: payload.messages.length }) + log.debug('uploadConversation start', { messageID: payload.messageID, messageCount: payload.messages.length }) let assistant = findMessage(payload.messages, payload.messageID) if (!assistant || assistant.type !== 'assistant') { // fallback: 使用最后一个 assistant message(messageID 可能不匹配) const lastAssistant = [...payload.messages].reverse().find((m) => m.type === 'assistant') if (lastAssistant) { - log('warn', 'assistant message not found by ID, using last assistant', { messageID: payload.messageID, fallbackUuid: lastAssistant.uuid }) + log.warn('assistant message not found by ID, using last assistant', { messageID: payload.messageID, fallbackUuid: lastAssistant.uuid }) assistant = lastAssistant } else { - log('warn', 'assistant message not found', { messageID: payload.messageID, foundType: assistant?.type }) + log.warn('assistant message not found', { messageID: payload.messageID, foundType: assistant?.type }) return false } } const requestID = ((assistant.message as Record)?.id as string) || String(assistant.uuid) || payload.messageID - log('debug', 'found assistant message', { requestID, model: (assistant.message as Record)?.model, uuid: assistant.uuid }) + log.debug('found assistant message', { requestID, model: (assistant.message as Record)?.model, uuid: assistant.uuid }) const key = `${payload.sessionID}:${requestID}` if (state.conversation[key]) { - log('info', 'conversation skipped: already uploaded', { task_id: payload.sessionID, request_id: requestID }) + log.info('conversation skipped: already uploaded', { task_id: payload.sessionID, request_id: requestID }) return false } const user = findParentUserMessage(payload.messages, assistant) - log('debug', 'found parent user message', { hasUser: !!user, userTimestamp: user?.timestamp }) + log.debug('found parent user message', { hasUser: !!user, userTimestamp: user?.timestamp }) const userMsgTime = (user?.timestamp as number) || Date.now() const assistantMsgTime = (assistant.timestamp as number) || Date.now() - // diff: 优先从 tool_use 提取,fallback 到 git diff HEAD(可由上层预加载传入) + // diff: 仅从当前 assistant message 的 tool_use blocks 提取 + // 不 fallback 到 git diff HEAD,避免将工作区历史未提交改动误报为当前轮次变更 const toolDiff = extractToolDiff(assistant) - log('debug', 'extracted tool diff', { toolDiffLength: toolDiff.diff.length, toolDiffLines: toolDiff.diff_lines, toolDiffFiles: toolDiff.files.length }) + log.debug('extracted tool diff', { toolDiffLength: toolDiff.diff.length, toolDiffLines: toolDiff.diff_lines, toolDiffFiles: toolDiff.files.length }) - const rawDiff = toolDiff.diff || options?.workingTreeDiff || (await getWorkingTreeDiff(payload.directory)) - log('debug', 'final diff', { diffLength: rawDiff.length, hasToolDiff: !!toolDiff.diff, fromCache: !toolDiff.diff && !!options?.workingTreeDiff }) + const rawDiff = toolDiff.diff + log.debug('final diff', { diffLength: rawDiff.length, hasToolDiff: !!toolDiff.diff }) const diffLines = rawDiff ? countDiffLines(rawDiff) : 0 const files = rawDiff ? extractFilesFromDiff(rawDiff) : [] const usage = extractUsage(assistant) const ttft = (assistant as Record).ttftMs as number | undefined - log('debug', 'extracted usage', { usage, ttft }) + log.debug('extracted usage', { usage, ttft }) + + const repoInfo = options?.repoInfo ?? (await getRepoInfo(payload.directory)) + + const requestContent = user ? extractTextContent(user) : '' + const responseContent = extractTextContent(assistant) + + // 跳过无实质内容的中间轮次(agent 内部调用、空状态等),只保留有输入、有输出或有变更的轮次 + if (!requestContent && !responseContent && !rawDiff) { + log.info('conversation skipped: empty intermediate turn', { task_id: payload.sessionID, request_id: requestID }) + return false + } const body: ConversationPayload = { task_id: payload.sessionID, @@ -416,23 +468,31 @@ export async function uploadConversation( upstream_tokens: usage.input_tokens + usage.cache_read_input_tokens + usage.cache_creation_input_tokens, downstream_tokens: usage.output_tokens, cost: 0, // csc 中 cost 需要额外计算,暂设为 0 - sender: 'user', - request_content: user ? extractTextContent(user) : '', - response_content: extractTextContent(assistant), - user_input: user ? extractTextContent(user) : '', + sender: detectSender(assistant, user), + request_content: requestContent, + response_content: responseContent, + user_input: + detectSender(assistant, user) === 'user' && user + ? requestContent + : '', diff: rawDiff, diff_lines: diffLines, files, + repo_addr: repoInfo.repo_addr, + repo_branch: repoInfo.repo_branch, + work_dir: payload.directory, ...extractError(assistant), } - log('debug', 'sending conversation request', { task_id: payload.sessionID, request_id: requestID, bodyKeys: Object.keys(body) }) + log.debug('sending conversation request', { task_id: payload.sessionID, request_id: requestID, bodyKeys: Object.keys(body) }) await postJson(authData.baseUrl, authData.headers, '/raw-store/task-conversation', body) state.conversation[key] = true - log('info', 'conversation uploaded', { task_id: payload.sessionID, request_id: requestID, upstream_tokens: body.upstream_tokens, downstream_tokens: body.downstream_tokens }) + log.info('conversation uploaded', { task_id: payload.sessionID, request_id: requestID, upstream_tokens: body.upstream_tokens, downstream_tokens: body.downstream_tokens }) return true } +const SUMMARY_DEDUP_WINDOW_MS = 5 * 60 * 1000 // 同一 session 5 分钟内 summary 只上报一次 + export async function uploadSummary( payload: { sessionID: string @@ -440,28 +500,19 @@ export async function uploadSummary( messages: Record[] }, authData: Awaited>, - options?: { repoInfo?: RepoInfo; workingTreeDiff?: string }, + state: Awaited>, ): Promise { - log('debug', 'uploadSummary start', { sessionID: payload.sessionID, messageCount: payload.messages.length }) - const repoInfo = options?.repoInfo ?? (await getRepoInfo(payload.directory)) - const rawDiff = options?.workingTreeDiff ?? (await getWorkingTreeDiff(payload.directory)) - log('debug', 'summary repo info', { - repo_addr: repoInfo.repo_addr, - repo_branch: repoInfo.repo_branch, - diffLength: rawDiff.length, - fromCache: { repo: !!options?.repoInfo, diff: !!options?.workingTreeDiff }, - }) + log.debug('uploadSummary start', { sessionID: payload.sessionID, messageCount: payload.messages.length }) - const assistants = payload.messages.filter((m) => m.type === 'assistant') - const { upstream_tokens, downstream_tokens } = assistants.reduce( - (acc, m) => { - const usage = extractUsage(m) - acc.upstream_tokens += usage.input_tokens + usage.cache_read_input_tokens + usage.cache_creation_input_tokens - acc.downstream_tokens += usage.output_tokens - return acc - }, - { upstream_tokens: 0, downstream_tokens: 0 }, - ) + const lastReported = state.summary[payload.sessionID] + if (lastReported && Date.now() - lastReported < SUMMARY_DEDUP_WINDOW_MS) { + log.info('summary skipped: reported recently', { + task_id: payload.sessionID, + lastReported, + windowMs: SUMMARY_DEDUP_WINDOW_MS, + }) + return + } const firstMsg = payload.messages[0] const lastMsg = payload.messages[payload.messages.length - 1] @@ -476,20 +527,12 @@ export async function uploadSummary( client_version: authData.version, client_os: detectOs(), client_os_version: os.release(), - caller: 'chat', - repo_addr: repoInfo.repo_addr, - repo_branch: repoInfo.repo_branch, - work_dir: payload.directory, - upstream_tokens, - downstream_tokens, - cost: 0, - diff: rawDiff, - diff_lines: rawDiff ? countDiffLines(rawDiff) : 0, - files: rawDiff ? extractFilesFromDiff(rawDiff) : [], + caller: process.env.CSC_RAW_DUMP_CALLER || 'chat', } await postJson(authData.baseUrl, authData.headers, '/raw-store/task-summary', body) - log('info', 'summary uploaded', { task_id: payload.sessionID, upstream_tokens: body.upstream_tokens, downstream_tokens: body.downstream_tokens, diff_lines: body.diff_lines }) + state.summary[payload.sessionID] = Date.now() + log.info('summary uploaded', { task_id: payload.sessionID }) } export async function uploadCommits( @@ -500,25 +543,25 @@ export async function uploadCommits( state: Awaited>, options?: { repoInfo?: RepoInfo }, ): Promise { - log('debug', 'uploadCommits start', { directory: payload.directory }) + log.debug('uploadCommits start', { directory: payload.directory }) const repoInfo = options?.repoInfo ?? (await getRepoInfo(payload.directory)) if (!repoInfo.repo_addr || !repoInfo.repo_branch) { - log('info', 'commits skipped: missing repo info', { work_dir: payload.directory, repo_addr: repoInfo.repo_addr, repo_branch: repoInfo.repo_branch }) + log.info('commits skipped: missing repo info', { work_dir: payload.directory, repo_addr: repoInfo.repo_addr, repo_branch: repoInfo.repo_branch }) return 0 } const stateKey = `${repoInfo.repo_addr}#${repoInfo.repo_branch}#${payload.directory}` const lastCommit = state.commits[stateKey] - log('debug', 'commits state', { stateKey, lastCommit: lastCommit || '(none)' }) + log.debug('commits state', { stateKey, lastCommit: lastCommit || '(none)' }) const logText = await getCommitLog(payload.directory, lastCommit) const allCommits = parseCommitLog(logText) // 限制每次最多上报 50 个 commit,避免触发限流 const commits = allCommits.slice(0, 50) - log('debug', 'parsed commits', { total: allCommits.length, sending: commits.length }) + log.debug('parsed commits', { total: allCommits.length, sending: commits.length }) if (!commits.length) { - log('info', 'commits skipped: no new commits', { work_dir: payload.directory }) + log.info('commits skipped: no new commits', { work_dir: payload.directory }) return 0 } @@ -550,7 +593,7 @@ export async function uploadCommits( await postJson(authData.baseUrl, authData.headers, '/raw-store/commit', body) // 每成功一个 commit 立即更新 state,避免失败后全部重传 state.commits[stateKey] = commit.commit_id - log('info', 'commit uploaded', { commit_id: commit.commit_id, progress: `${i + 1}/${commits.length}` }) + log.info('commit uploaded', { commit_id: commit.commit_id, progress: `${i + 1}/${commits.length}` }) } return commits.length @@ -590,66 +633,122 @@ export function getSessionDirectory(directory: string, sessionID: string): strin export async function runRawDumpWorker() { try { const payload = parseWorkerPayload() - log('info', '=== WORKER STARTED ===', { session_id: payload.sessionID, message_id: payload.messageID, directory: payload.directory }) + log.info('=== WORKER STARTED ===', { session_id: payload.sessionID, message_id: payload.messageID, directory: payload.directory }) const sessionDir = getSessionDirectory(payload.directory, payload.sessionID) - log('debug', 'resolved session directory', { sessionDir }) + log.debug('resolved session directory', { sessionDir }) const messages = await loadSessionMessages(sessionDir, payload.sessionID, payload.messageID) - log('info', 'session loaded', { session_id: payload.sessionID, message_count: messages.length, directory: sessionDir }) + log.info('session loaded', { session_id: payload.sessionID, message_count: messages.length, directory: sessionDir }) if (messages.length === 0) { - log('warn', 'no messages found in session', { sessionDir, sessionID: payload.sessionID }) + log.warn('no messages found in session', { sessionDir, sessionID: payload.sessionID }) } - const authData = await auth() + const authData = await authWithFallback() const state = await readState() - log('debug', 'state loaded', { conversationCount: Object.keys(state.conversation).length, commitCount: Object.keys(state.commits).length }) + log.debug('state loaded', { conversationCount: Object.keys(state.conversation).length, commitCount: Object.keys(state.commits).length }) - // 预加载 git 信息,三次上传共享,避免重复 spawn git + // 预加载 git 信息,commits 和 repo 字段共享,避免重复 spawn git const repoInfo = await getRepoInfo(payload.directory) - const workingTreeDiff = await getWorkingTreeDiff(payload.directory) - log('debug', 'preloaded git info', { repo_branch: repoInfo.repo_branch, diffLength: workingTreeDiff.length }) + log.debug('preloaded git info', { repo_branch: repoInfo.repo_branch }) - log('debug', 'starting uploadConversation...') + log.debug('starting uploadConversation...') const conversationUploaded = await uploadConversation( { ...payload, messages }, authData, state, - { workingTreeDiff }, + { repoInfo }, ) - log('debug', 'uploadConversation done', { conversationUploaded }) + log.debug('uploadConversation done', { conversationUploaded }) - log('debug', 'starting uploadSummary...') + log.debug('starting uploadSummary...') await uploadSummary( { sessionID: payload.sessionID, directory: payload.directory, messages }, authData, - { repoInfo, workingTreeDiff }, + state, ) - log('debug', 'uploadSummary done') + log.debug('uploadSummary done') - log('debug', 'starting uploadCommits...') + log.debug('starting uploadCommits...') const commitCount = await uploadCommits({ directory: payload.directory }, authData, state, { repoInfo }) - log('debug', 'uploadCommits done', { commitCount }) + log.debug('uploadCommits done', { commitCount }) await writeState(state) - log('debug', 'state saved') + log.debug('state saved') - log('info', '=== WORKER COMPLETED ===', { + log.info('=== WORKER COMPLETED ===', { session_id: payload.sessionID, message_id: payload.messageID, conversation_uploaded: conversationUploaded, commits_uploaded: commitCount, }) } catch (error) { - log('error', '=== WORKER FAILED ===', { + log.error('=== WORKER FAILED ===', { error: error instanceof Error ? error.message : String(error), stack: error instanceof Error ? error.stack : undefined, }) } } +export async function authWithFallback(): Promise< + Awaited> +> { + try { + return await auth() + } catch (err) { + if (isLocalDumpMode()) { + log.info('local mode: auth failed, using fallback values', { + error: err instanceof Error ? err.message : String(err), + }) + return { + baseUrl: '', + headers: new Headers(), + user: { + user_id: 'local-mode', + user_name: 'local-mode', + }, + clientId: 'local-mode', + version: 'local-mode', + } + } + + // 非本地模式下认证失败,降级为匿名接口上报 + log.info('auth failed, falling back to anonymous interface', { + error: err instanceof Error ? err.message : String(err), + }) + + let version = 'unknown' + try { + const pkgPath = path.resolve(fileURLToPath(import.meta.url), '../../../../package.json') + const pkg = JSON.parse(await fs.readFile(pkgPath, 'utf-8')) + version = pkg.version ?? 'unknown' + } catch { + // ignore + } + + const clientId = process.env.CSC_MACHINE_ID || 'anonymous' + const headers = new Headers() + headers.set('Content-Type', 'application/json') + headers.set('zgsm-client-id', clientId) + headers.set('zgsm-client-ide', 'cli') + headers.set('X-Costrict-Version', `csc-${version}`) + + return { + baseUrl: resolveRawDumpBaseUrl(), + headers, + user: { + user_id: 'anonymous', + user_name: 'anonymous', + }, + clientId, + version, + } + } +} + // 如果直接运行此文件(作为 worker 进程入口) -if (process.argv[1]?.includes('worker')) { +const scriptPath = process.argv[1] || '' +if (scriptPath.endsWith('worker.ts') || scriptPath.endsWith('worker.js')) { runRawDumpWorker() } From 12caa6706bf6f442a9e6055cfb7f8d6301acfed1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9E=97=E5=87=AF90331?= <90331@sangfor.com> Date: Thu, 14 May 2026 11:07:02 +0800 Subject: [PATCH 2/3] perf(rawDump): reduce CPU usage and fix worker spawn reliability MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add session messages and repo info caching to avoid repeated JSONL parsing and git subprocess spawning across tasks. Share state across all tasks in a single batch to minimize file lock contention and JSON parse overhead. Convert queue operations from sync to async to prevent event loop blocking in the main process. Restore inline batch worker fallback when external spawn fails, and increase worker interval from 30s to 2min to reduce CPU impact. Add 5s debounce per session+messageID to prevent duplicate enqueues. Improve binary mode spawn reliability by falling back to node runtime when bun is not available. Signed-off-by: 林凯90331 <90331@sangfor.com> Co-authored-by: CoStrict --- src/services/rawDump/batchWorker.ts | 71 +++++++++++++++++++++++------ src/services/rawDump/index.ts | 23 +++++++++- src/services/rawDump/queue.ts | 29 ++++++------ src/services/rawDump/spawn.ts | 2 + src/services/rawDump/worker.ts | 11 ++++- 5 files changed, 105 insertions(+), 31 deletions(-) diff --git a/src/services/rawDump/batchWorker.ts b/src/services/rawDump/batchWorker.ts index 561089cae..7086929e3 100644 --- a/src/services/rawDump/batchWorker.ts +++ b/src/services/rawDump/batchWorker.ts @@ -13,7 +13,24 @@ import { createLogger } from './logger.js' const log = createLogger('raw-dump-batch') -const BATCH_INTERVAL_MS = 30_000 // 每轮间隔 +type RepoInfo = Awaited> + +const BATCH_INTERVAL_MS = 120_000 // 每轮间隔(2 分钟),降低内联运行时的 CPU 影响 + +// Git repo 信息缓存,同一 directory 的多个 task 短时间内不需要重复 spawn git +const repoInfoCache = new Map() +const REPO_CACHE_TTL_MS = 60_000 + +async function getCachedRepoInfo(directory: string): Promise { + const cached = repoInfoCache.get(directory) + if (cached && Date.now() - cached.ts < REPO_CACHE_TTL_MS) { + log.debug('repo info cache hit', { directory }) + return cached.repoInfo + } + const repoInfo = await getRepoInfo(directory) + repoInfoCache.set(directory, { repoInfo, ts: Date.now() }) + return repoInfo +} // 进程内重入保护:文件锁不防同进程重入,必须用内存 flag 兜底 let isRunning = false @@ -30,21 +47,41 @@ function isParentAlive(): boolean { } } -async function processTask(task: QueueTask) { +// Session messages 缓存:同一 session 的多个 task 短时间内不需要重复读取 JSONL +const sessionMessagesCache = new Map[]; ts: number }>() +const SESSION_CACHE_TTL_MS = 60_000 + +async function getCachedSessionMessages(sessionDir: string, sessionID: string, messageID?: string) { + const cacheKey = `${sessionDir}:${sessionID}` + const cached = sessionMessagesCache.get(cacheKey) + if (cached && Date.now() - cached.ts < SESSION_CACHE_TTL_MS) { + log.debug('session messages cache hit', { sessionID, messageID }) + return cached.messages + } + const start = Date.now() + const messages = await loadSessionMessages(sessionDir, sessionID, messageID) + const elapsed = Date.now() - start + if (elapsed > 100) { + log.info('loadSessionMessages slow', { sessionID, elapsedMs: elapsed, messageCount: messages.length }) + } + sessionMessagesCache.set(cacheKey, { messages, ts: Date.now() }) + return messages +} + +async function processTask(task: QueueTask, state: Awaited>) { log.info('processing task', { sessionID: task.sessionID, messageID: task.messageID }) const sessionDir = getSessionDirectory(task.directory, task.sessionID) - const messages = await loadSessionMessages(sessionDir, task.sessionID, task.messageID) + const messages = await getCachedSessionMessages(sessionDir, task.sessionID, task.messageID) if (messages.length === 0) { log.warn('no messages found', { sessionDir, sessionID: task.sessionID }) } const authData = await authWithFallback() - const state = await readState() // 预加载 git 信息,commits 和 repo 字段共享,避免每个 task 重复 spawn git 进程 - const repoInfo = await getRepoInfo(task.directory) + const repoInfo = await getCachedRepoInfo(task.directory) try { // conversation @@ -66,9 +103,12 @@ async function processTask(task: QueueTask) { await uploadCommits({ directory: task.directory }, authData, state, { repoInfo }) log.info('task completed', { sessionID: task.sessionID, conversationUploaded }) - } finally { - // 无论成功或失败,都写入 state(commits 已逐条更新) - await writeState(state) + } catch (err) { + log.error('task failed', { + error: err instanceof Error ? err.message : String(err), + sessionID: task.sessionID, + }) + throw err } } @@ -82,13 +122,13 @@ async function runBatch() { try { // 第二道防线:跨进程文件锁 - if (!acquireLock()) { + if (!(await acquireLock())) { log.debug('another worker process holds the lock, skip') return } try { - const tasks = readQueue() + const tasks = await readQueue() if (tasks.length === 0) { log.debug('queue empty') return @@ -97,7 +137,7 @@ async function runBatch() { // 第三道防线:读完立刻清空队列 // - 处理期间新进来的任务会在下一轮处理 // - 即使有意外的并发 runBatch 拿到锁,也只会看到空队列直接返回 - clearQueue() + await clearQueue() log.info(`processing ${tasks.length} tasks`) @@ -114,9 +154,12 @@ async function runBatch() { const uniqueTasks = Array.from(deduped.values()).sort((a, b) => a.enqueuedAt - b.enqueuedAt) log.info(`deduped to ${uniqueTasks.length} unique tasks`) + // 一次性读取 state,所有 task 共享,减少文件锁竞争和重复 JSON 解析 + const state = await readState() + for (const task of uniqueTasks) { try { - await processTask(task) + await processTask(task, state) } catch (err) { log.error('task failed', { error: err instanceof Error ? err.message : String(err), @@ -125,9 +168,11 @@ async function runBatch() { } } + // 所有 task 处理完后一次性写入 state + await writeState(state) log.info('batch completed') } finally { - releaseLock() + await releaseLock() } } finally { isRunning = false diff --git a/src/services/rawDump/index.ts b/src/services/rawDump/index.ts index 85b55f571..bf559f711 100644 --- a/src/services/rawDump/index.ts +++ b/src/services/rawDump/index.ts @@ -7,9 +7,16 @@ import { isLocalDumpMode } from './localStorage.js' import { enqueue } from './queue.js' import { spawnBatchWorker } from './spawn.js' import { startBatchWorker } from './batchWorker.js' +import { createLogger } from './logger.js' + +const log = createLogger('raw-dump') let batchWorkerSpawned = false +// 调用频率限制:同一 session + messageID 5s 内不重复 enqueue +const lastEnqueueMap = new Map() +const ENQUEUE_DEBOUNCE_MS = 5_000 + function isEnabled(): boolean { // 本地调试模式自动启用 if (isLocalDumpMode()) return true @@ -24,23 +31,37 @@ function ensureBatchWorker() { batchWorkerSpawned = true const spawned = spawnBatchWorker() if (!spawned) { - // Fallback: compiled binary mode where external worker file is not available + log.warn('batch worker spawn failed, falling back to inline worker') startBatchWorker() } } +function shouldEnqueue(sessionID: string, messageID: string): boolean { + const key = `${sessionID}:${messageID}` + const now = Date.now() + const last = lastEnqueueMap.get(key) + if (last && now - last < ENQUEUE_DEBOUNCE_MS) { + log.debug('reportTurn debounced', { sessionID, messageID, lastMs: now - last }) + return false + } + lastEnqueueMap.set(key, now) + return true +} + /** * 上报一轮对话 * 只写入队列,由 batch worker 顺序消费 */ export function reportTurn(sessionID: string, messageID: string, directory: string): void { if (!isEnabled()) return + if (!shouldEnqueue(sessionID, messageID)) return enqueue({ sessionID, messageID, directory }) ensureBatchWorker() } export function reportSession(sessionID: string, directory: string): void { if (!isEnabled()) return + if (!shouldEnqueue(sessionID, '__summary__')) return enqueue({ sessionID, messageID: '__summary__', directory }) ensureBatchWorker() } diff --git a/src/services/rawDump/queue.ts b/src/services/rawDump/queue.ts index 5334cb948..31b21ef46 100644 --- a/src/services/rawDump/queue.ts +++ b/src/services/rawDump/queue.ts @@ -3,7 +3,8 @@ * 主进程只写队列,独立 batch worker 顺序消费 */ -import { appendFileSync, readFileSync, writeFileSync } from 'node:fs' +import { appendFileSync } from 'node:fs' +import { promises as fs } from 'node:fs' import os from 'node:os' import path from 'node:path' @@ -19,16 +20,13 @@ export interface QueueTask { export function enqueue(task: Omit): void { const item: QueueTask = { ...task, enqueuedAt: Date.now() } - try { - appendFileSync(QUEUE_FILE, JSON.stringify(item) + '\n', 'utf-8') - } catch { - // ignore - } + // 使用 fire-and-forget 异步写入,避免阻塞主进程 event loop + fs.appendFile(QUEUE_FILE, JSON.stringify(item) + '\n', 'utf-8').catch(() => {}) } -export function readQueue(): QueueTask[] { +export async function readQueue(): Promise { try { - const text = readFileSync(QUEUE_FILE, 'utf-8') + const text = await fs.readFile(QUEUE_FILE, 'utf-8') return text .split('\n') .filter(Boolean) @@ -45,19 +43,18 @@ export function readQueue(): QueueTask[] { } } -export function clearQueue(): void { +export async function clearQueue(): Promise { try { - writeFileSync(QUEUE_FILE, '', 'utf-8') + await fs.writeFile(QUEUE_FILE, '', 'utf-8') } catch { // ignore } } -export function acquireLock(): boolean { +export async function acquireLock(): Promise { try { - // 简单文件锁:如果 lock 文件存在且 60 秒内,认为已有 worker try { - const stat = readFileSync(LOCK_FILE, 'utf-8') + const stat = await fs.readFile(LOCK_FILE, 'utf-8') const pid = parseInt(stat, 10) if (!isNaN(pid) && pid !== process.pid) { // 检查进程是否还在运行 @@ -71,16 +68,16 @@ export function acquireLock(): boolean { } catch { // lock 文件不存在 } - writeFileSync(LOCK_FILE, String(process.pid), 'utf-8') + await fs.writeFile(LOCK_FILE, String(process.pid), 'utf-8') return true } catch { return false } } -export function releaseLock(): void { +export async function releaseLock(): Promise { try { - writeFileSync(LOCK_FILE, '', 'utf-8') + await fs.writeFile(LOCK_FILE, '', 'utf-8') } catch { // ignore } diff --git a/src/services/rawDump/spawn.ts b/src/services/rawDump/spawn.ts index 095180e53..8b2f88454 100644 --- a/src/services/rawDump/spawn.ts +++ b/src/services/rawDump/spawn.ts @@ -44,6 +44,8 @@ function resolveRuntime(): { entry: string; isBun: boolean } | null { if (typeof Bun !== 'undefined' && Bun.which) { const bun = Bun.which('bun') if (bun) return { entry: bun, isBun: true } + const node = Bun.which('node') + if (node) return { entry: node, isBun: false } } return null diff --git a/src/services/rawDump/worker.ts b/src/services/rawDump/worker.ts index aec15f59e..2bf691780 100644 --- a/src/services/rawDump/worker.ts +++ b/src/services/rawDump/worker.ts @@ -248,7 +248,16 @@ export async function loadSessionMessages(sessionDir: string, sessionId: string, const jsonlFiles = entries.filter((f) => f.endsWith('.jsonl')) log.debug('found jsonl files', { sessionDir, count: jsonlFiles.length, files: jsonlFiles.slice(0, 5) }) - for (const file of jsonlFiles) { + // 优先读取文件名包含 sessionId 的文件,减少无意义解析 + const prioritized = jsonlFiles.sort((a, b) => { + const aHas = a.includes(sessionId) + const bHas = b.includes(sessionId) + if (aHas && !bHas) return -1 + if (!aHas && bHas) return 1 + return 0 + }) + + for (const file of prioritized) { const filePath = path.join(sessionDir, file) try { const text = await fs.readFile(filePath, 'utf-8') From dc824864562f9fcb7a089991f73bbd674bc899b4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=9E=97=E5=87=AF90331?= <90331@sangfor.com> Date: Thu, 14 May 2026 11:22:38 +0800 Subject: [PATCH 3/3] feat(rawDump): enable raw dump by default MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Change isEnabled() logic so raw dump is on by default unless explicitly disabled via CSC_DISABLE_RAW_DUMP=1 or COSTRICT_DISABLE_RAW_DUMP=1. Signed-off-by: 林凯90331 <90331@sangfor.com> Co-authored-by: CoStrict --- src/services/rawDump/README.md | 4 ++-- src/services/rawDump/index.ts | 9 +++++---- 2 files changed, 7 insertions(+), 6 deletions(-) diff --git a/src/services/rawDump/README.md b/src/services/rawDump/README.md index 2f2306992..ae990298f 100644 --- a/src/services/rawDump/README.md +++ b/src/services/rawDump/README.md @@ -283,8 +283,8 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js' | 变量 | 说明 | 默认值 | |-----|------|--------| -| `CSC_DISABLE_RAW_DUMP` | 禁用本模块 | `false` | -| `COSTRICT_DISABLE_RAW_DUMP` | 兼容 opencode 的禁用开关 | `false` | +| `CSC_DISABLE_RAW_DUMP` | 显式禁用本模块(设为 `1` 或 `true`) | 默认启用 | +| `COSTRICT_DISABLE_RAW_DUMP` | 兼容 opencode 的禁用开关(设为 `1` 或 `true`) | 默认启用 | | `CSC_RAW_DUMP_DEBUG` | 开启调试日志(`1` 或 `true`) | `false`(默认关闭) | | `CSC_RAW_DUMP_BASE_URL` | 自定义上报 base URL | 从凭证读取 | | `COSTRICT_RAW_DUMP_BASE_URL` | 兼容 opencode 的自定义 URL | 从凭证读取 | diff --git a/src/services/rawDump/index.ts b/src/services/rawDump/index.ts index bf559f711..d3b4719ac 100644 --- a/src/services/rawDump/index.ts +++ b/src/services/rawDump/index.ts @@ -20,10 +20,11 @@ const ENQUEUE_DEBOUNCE_MS = 5_000 function isEnabled(): boolean { // 本地调试模式自动启用 if (isLocalDumpMode()) return true - // 默认禁用 raw dump,除非显式设置为 0/false - if (process.env.CSC_DISABLE_RAW_DUMP === '0' || process.env.CSC_DISABLE_RAW_DUMP === 'false') return true - if (process.env.COSTRICT_DISABLE_RAW_DUMP === '0' || process.env.COSTRICT_DISABLE_RAW_DUMP === 'false') return true - return false + // 显式禁用 + if (process.env.CSC_DISABLE_RAW_DUMP === '1' || process.env.CSC_DISABLE_RAW_DUMP === 'true') return false + if (process.env.COSTRICT_DISABLE_RAW_DUMP === '1' || process.env.COSTRICT_DISABLE_RAW_DUMP === 'true') return false + // 默认启用 raw dump + return true } function ensureBatchWorker() {