refactor(rawDump): prevent 429 with queue + batch worker
Replace per-turn detached workers with a file-backed queue consumed by
a single long-running batch worker (30s interval, 0-10s jitter).
- queue.ts, batchWorker.ts: file queue with pid-based lock for worker
singleton; tasks deduped by sessionID:messageID before processing
- worker.ts: retry with backoff (5s, 10s) for 429 and network errors;
update commit state per upload so partial failures resume cleanly;
export auth, loadSessionMessages, getSessionDirectory and upload*
helpers for batchWorker reuse
- git.ts: cap commit log at 50 entries within last 7 days; pause 500ms
every 10 commits to spread load
- worker.ts: resolve session jsonl via ~/.claude/projects/{normalized}
with fallbacks, scanning for the file containing the target session
- logger.ts: file + stderr logger gated by CSC_RAW_DUMP_DEBUG, default
silent
- sessionDataUploader.ts: implement createSessionTurnUploader to pick
the last assistant message; query.ts fires uploadSessionTurn after
query_api_streaming_end (non-blocking)
Signed-off-by: 林凯90331 <90331@sangfor.com>
Co-authored-by: CoStrict <zgsm@sangfor.com.cn>
This commit is contained in:
parent
28d4f1f4e5
commit
f3ece2f469
|
|
@ -113,6 +113,7 @@ import { createBudgetTracker, checkTokenBudget } from './query/tokenBudget.js'
|
|||
import { count } from './utils/array.js'
|
||||
import { createTrace, endTrace, isLangfuseEnabled } from './services/langfuse/index.js'
|
||||
import { getAPIProvider } from './utils/model/providers.js'
|
||||
import { uploadSessionTurn } from './utils/sessionDataUploader.js'
|
||||
|
||||
/* eslint-disable @typescript-eslint/no-require-imports */
|
||||
const snipModule = feature('HISTORY_SNIP')
|
||||
|
|
@ -909,6 +910,13 @@ async function* queryLoop(
|
|||
}
|
||||
queryCheckpoint('query_api_streaming_end')
|
||||
|
||||
// Report conversation/summary/commits to CoStrict raw-dump endpoint.
|
||||
// Fire-and-forget; non-blocking.
|
||||
const lastAssistant = assistantMessages.at(-1)
|
||||
if (lastAssistant) {
|
||||
uploadSessionTurn(getSessionId(), String(lastAssistant.uuid))
|
||||
}
|
||||
|
||||
// Yield deferred microcompact boundary message using actual API-reported
|
||||
// token deletion count instead of client-side estimates.
|
||||
// Entire block gated behind feature() so the excluded string
|
||||
|
|
|
|||
|
|
@ -6,7 +6,8 @@
|
|||
|
||||
**设计原则:**
|
||||
- **与框架解耦**:不依赖 React、Effect-TS、Ink 等任何 UI 框架
|
||||
- **非阻塞**:使用 detached 子进程执行上报,主进程立即返回
|
||||
- **非阻塞**:主进程只写入队列,由独立 batch worker 顺序消费,不阻塞主流程
|
||||
- **防限流**:队列 + 单 worker 顺序执行 + 请求间延迟 + 随机抖动,避免并发 429
|
||||
- **协议兼容**:与 opencode 的 raw-dump 插件保持接口对齐
|
||||
|
||||
---
|
||||
|
|
@ -15,13 +16,15 @@
|
|||
|
||||
```
|
||||
src/services/rawDump/
|
||||
├── README.md # 本文档
|
||||
├── types.ts # 类型定义 + 环境变量常量
|
||||
├── state.ts # 磁盘状态管理(去重)
|
||||
├── git.ts # Git 辅助函数封装
|
||||
├── worker.ts # 独立 worker 进程(实际上报逻辑)
|
||||
├── spawn.ts # 子进程启动器
|
||||
└── index.ts # 主入口 API
|
||||
├── README.md # 本文档
|
||||
├── types.ts # 类型定义 + 环境变量常量
|
||||
├── state.ts # 磁盘状态管理(去重)
|
||||
├── git.ts # Git 辅助函数封装
|
||||
├── worker.ts # 实际上报逻辑(被 batch worker 调用)
|
||||
├── queue.ts # 文件队列(主进程写入,worker 消费)
|
||||
├── batchWorker.ts # 独立 batch worker 进程(顺序消费队列)
|
||||
├── spawn.ts # 子进程启动器
|
||||
└── index.ts # 主入口 API
|
||||
```
|
||||
|
||||
---
|
||||
|
|
@ -31,14 +34,22 @@ src/services/rawDump/
|
|||
```
|
||||
主进程:assistant message 完成
|
||||
→ reportTurn(sessionID, messageID, directory)
|
||||
→ 内存去重(Set)
|
||||
→ spawn detached worker 进程
|
||||
→ worker 加载会话消息(JSONL)
|
||||
→ auth():加载凭证、刷新 token
|
||||
→ uploadConversation() → POST /raw-store/task-conversation
|
||||
→ uploadSummary() → POST /raw-store/task-summary
|
||||
→ uploadCommits() → POST /raw-store/commit
|
||||
→ writeState() → 更新去重状态
|
||||
→ enqueue({ sessionID, messageID, directory }) 写入队列文件
|
||||
→ ensureBatchWorker() 启动 detached batch worker(仅一次)
|
||||
|
||||
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()
|
||||
```
|
||||
|
||||
---
|
||||
|
|
@ -60,8 +71,7 @@ reportTurn(sessionId, assistantMessage.uuid, cwd)
|
|||
**推荐集成点:**
|
||||
|
||||
1. `src/query.ts` 中 streaming 结束后(`query_api_streaming_end` 之后)
|
||||
2. `src/costrict/provider/index.ts` 中 `message_stop` 事件后
|
||||
3. `src/utils/sessionDataUploader.ts` 已提供 `uploadSessionTurn()` 封装
|
||||
2. `src/utils/sessionDataUploader.ts` 已提供 `uploadSessionTurn()` 封装
|
||||
|
||||
---
|
||||
|
||||
|
|
@ -127,13 +137,15 @@ csc 没有 opencode 中的 `step-start`/`step-finish` snapshot 机制,采用
|
|||
|
||||
## 去重机制
|
||||
|
||||
### 1. 内存去重(进程内)
|
||||
### 1. 队列去重(进程内)
|
||||
同一个 session + messageID 的多个 task,batch worker 消费时只保留最新一个:
|
||||
```typescript
|
||||
const spawned = new Set<string>()
|
||||
const key = `${sessionID}:${messageID}`
|
||||
if (spawned.has(key)) return
|
||||
const key = `${task.sessionID}:${task.messageID}`
|
||||
const existing = deduped.get(key)
|
||||
if (!existing || task.enqueuedAt > existing.enqueuedAt) {
|
||||
deduped.set(key, task)
|
||||
}
|
||||
```
|
||||
限制 1024 条,超过后清理一半。
|
||||
|
||||
### 2. Conversation 去重(磁盘)
|
||||
```typescript
|
||||
|
|
@ -145,7 +157,7 @@ if (spawned.has(key)) return
|
|||
}
|
||||
```
|
||||
|
||||
### 3. Commits 去重(磁盘)
|
||||
### 3. Commits 去重(磁盘,逐条更新)
|
||||
```typescript
|
||||
// 以 repo#branch#workDir 为 key
|
||||
{
|
||||
|
|
@ -154,8 +166,22 @@ if (spawned.has(key)) return
|
|||
}
|
||||
}
|
||||
```
|
||||
- 有 lastCommit:取 `lastCommit..HEAD`
|
||||
- 无 lastCommit:取 30 天内 commits
|
||||
- **逐 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`
|
||||
- **批次延迟**:每上传 10 个 commits 后暂停 500ms,避免触发限流
|
||||
|
||||
---
|
||||
|
||||
## 429 防护机制
|
||||
|
||||
1. **队列 + 单 worker**:主进程只 enqueue,只有一个 batch worker 顺序消费,天然避免并发
|
||||
2. **文件锁**:`acquireLock()` / `releaseLock()` 确保同一时刻只有一个 worker 在运行
|
||||
3. **请求重试**:`postJson()` 对 429 和网络错误自动重试 3 次,退避间隔 5s、10s
|
||||
4. **commit 批次延迟**:每 10 个 commits 暂停 500ms
|
||||
5. **随机抖动**:batch worker 启动后首次执行有 0-10s 随机延迟,避免规律性请求
|
||||
6. **commit 数量限制**:单次最多获取 50 个 commits,时间范围限制为 7 天
|
||||
|
||||
---
|
||||
|
||||
|
|
@ -184,14 +210,16 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js'
|
|||
|-----|------|--------|
|
||||
| `CSC_DISABLE_RAW_DUMP` | 禁用本模块 | `false` |
|
||||
| `COSTRICT_DISABLE_RAW_DUMP` | 兼容 opencode 的禁用开关 | `false` |
|
||||
| `CSC_RAW_DUMP_DEBUG` | 开启调试日志(`1` 或 `true`) | `false`(默认关闭) |
|
||||
| `CSC_RAW_DUMP_BASE_URL` | 自定义上报 base URL | 从凭证读取 |
|
||||
| `COSTRICT_RAW_DUMP_BASE_URL` | 兼容 opencode 的自定义 URL | 从凭证读取 |
|
||||
| `COSTRICT_BASE_URL` | CoStrict 服务地址 | `https://zgsm.sangfor.com` |
|
||||
|
||||
---
|
||||
|
||||
## 状态文件
|
||||
## 状态文件与日志文件
|
||||
|
||||
### 状态文件
|
||||
```
|
||||
~/.claude/csc-raw-dump-state.json
|
||||
```
|
||||
|
|
@ -209,26 +237,33 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js'
|
|||
}
|
||||
```
|
||||
|
||||
### 日志文件
|
||||
```
|
||||
~/.claude/csc-raw-dump.log
|
||||
```
|
||||
|
||||
主进程和 batch worker 的日志都追加写入该文件。由于 worker 是 detached 进程(`stdio: 'ignore'`),日志只能通过文件查看。
|
||||
|
||||
---
|
||||
|
||||
## 注意事项与待完善项
|
||||
|
||||
1. **Cost 计算**
|
||||
1. **Cost 计算**
|
||||
当前 `cost` 字段设为 0。需接入 `src/cost-tracker.ts` 的 `calculateUSDCost()` 或从 `bootstrap/state.ts` 获取每轮/累计 cost。
|
||||
|
||||
2. **TTFT 获取**
|
||||
2. **TTFT 获取**
|
||||
当前从 assistant message 的 `ttftMs` 字段读取。需确认 csc 是否在 message 对象上保存了该值,否则需要在 streaming 开始时手动计时。
|
||||
|
||||
3. **会话目录**
|
||||
`getSessionDirectory()` 使用启发式查找(`directory/.claude/sessions`、`directory/.claude`、directory 本身)。需根据 csc 实际会话 JSONL 存放路径校准。
|
||||
3. **会话目录**
|
||||
`getSessionDirectory()` 使用启发式查找(`~/.claude/projects/{normalizedPath}` 等)。csc 实际会话 JSONL 存放路径为 `~/.claude/projects/{sanitizePath(cwd)}/{sessionId}.jsonl`。
|
||||
|
||||
4. **User 消息关联**
|
||||
4. **User 消息关联**
|
||||
当前按消息列表顺序查找前一个 `type === 'user'` 的消息。若 csc 存在明确的 parent-child 关系,应改用 `parentID` 或类似字段。
|
||||
|
||||
5. **Model 信息**
|
||||
5. **Model 信息**
|
||||
`model` 字段取自 `assistant.message.model`。若该字段不可靠,可从 `bootstrap/state.ts` 的 `getCurrentModel()` 获取。
|
||||
|
||||
6. **Sender 识别**
|
||||
6. **Sender 识别**
|
||||
当前固定为 `"user"`。若 csc 支持 agent/agentic 模式,需根据消息来源判断 `"user"` 或 `"agent"`。
|
||||
|
||||
---
|
||||
|
|
@ -242,18 +277,38 @@ import { refreshCoStrictToken } from '../../costrict/provider/token.js'
|
|||
| 会话加载 | 内存 Session 对象 | JSONL 文件解析 |
|
||||
| Cost 来源 | `assistant.info.cost` | 待接入 cost-tracker |
|
||||
| 运行时 | Effect-TS | Bun + 纯 Node.js API |
|
||||
| Worker 启动 | `bun run index.ts raw-dump _worker` | `bun run worker.ts` |
|
||||
| 上报模式 | 单条即时上报 | 队列 + batch worker 顺序消费 |
|
||||
| 限流防护 | 无 | 队列 + 单 worker + 重试 + 批次延迟 + 抖动 |
|
||||
| 凭证路径 | `~/.costrict/credentials.json` | `~/.claude/csc-auth.json` |
|
||||
|
||||
---
|
||||
|
||||
## 调试
|
||||
|
||||
Worker 进程的错误和日志通过 `console.error` 输出到 stderr,可在启动时重定向:
|
||||
调试日志默认关闭,通过环境变量开启:
|
||||
|
||||
```bash
|
||||
# 查看 worker 日志
|
||||
CSC_RAW_DUMP_DEBUG=1 csc
|
||||
# 开启调试日志
|
||||
export CSC_RAW_DUMP_DEBUG=1
|
||||
|
||||
# 查看日志
|
||||
tail -f ~/.claude/csc-raw-dump.log
|
||||
```
|
||||
|
||||
或查看系统日志(worker 为 detached 进程,日志不输出到主进程终端)。
|
||||
关键日志标识:
|
||||
- `[raw-dump:info]` / `[raw-dump:debug]` — worker.ts 中的日志
|
||||
- `[raw-dump-batch:info]` / `[raw-dump-batch:debug]` — batchWorker.ts 中的日志
|
||||
|
||||
日志模块完全独立(`logger.ts`),默认不产生任何输出,不创建日志文件。
|
||||
|
||||
常用排查命令:
|
||||
```bash
|
||||
# 查看 state 文件
|
||||
cat ~/.claude/csc-raw-dump-state.json
|
||||
|
||||
# 查看队列文件
|
||||
cat ~/.claude/csc-raw-dump-queue.jsonl
|
||||
|
||||
# 查看是否有 worker 在运行(锁文件)
|
||||
cat ~/.claude/csc-raw-dump.lock
|
||||
```
|
||||
|
|
|
|||
113
src/services/rawDump/batchWorker.ts
Normal file
113
src/services/rawDump/batchWorker.ts
Normal file
|
|
@ -0,0 +1,113 @@
|
|||
/**
|
||||
* Raw Dump Batch Worker
|
||||
* 顺序消费队列,避免并发 429
|
||||
* 独立进程,通过 setInterval 定期执行
|
||||
*/
|
||||
|
||||
import { uploadConversation, uploadSummary, uploadCommits, auth } 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 { createLogger } from './logger.js'
|
||||
|
||||
const log = createLogger('raw-dump-batch')
|
||||
|
||||
const BATCH_INTERVAL_MS = 30_000 // 30 秒检查一次队列
|
||||
|
||||
async function processTask(task: QueueTask) {
|
||||
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 })
|
||||
}
|
||||
|
||||
const authData = await auth()
|
||||
const state = await readState()
|
||||
|
||||
try {
|
||||
// conversation
|
||||
const conversationUploaded = await uploadConversation(
|
||||
{ sessionID: task.sessionID, messageID: task.messageID, directory: task.directory, messages },
|
||||
authData,
|
||||
state,
|
||||
)
|
||||
|
||||
// summary(每个 turn 都报,但内容会累积)
|
||||
await uploadSummary({ sessionID: task.sessionID, directory: task.directory, messages }, authData)
|
||||
|
||||
// commits(限制频率,避免重复上报)
|
||||
await uploadCommits({ directory: task.directory }, authData, state)
|
||||
|
||||
log('info', 'task completed', { sessionID: task.sessionID, conversationUploaded })
|
||||
} finally {
|
||||
// 无论成功或失败,都写入 state(commits 已逐条更新)
|
||||
await writeState(state)
|
||||
}
|
||||
}
|
||||
|
||||
async function runBatch() {
|
||||
if (!acquireLock()) {
|
||||
log('debug', 'another worker is running, skip')
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
const tasks = readQueue()
|
||||
if (tasks.length === 0) {
|
||||
log('debug', 'queue empty')
|
||||
return
|
||||
}
|
||||
|
||||
log('info', `processing ${tasks.length} tasks`)
|
||||
|
||||
// 去重:同一个 session 的多个 task,只保留最新的一个
|
||||
const deduped = new Map<string, QueueTask>()
|
||||
for (const task of tasks) {
|
||||
const key = `${task.sessionID}:${task.messageID}`
|
||||
const existing = deduped.get(key)
|
||||
if (!existing || task.enqueuedAt > existing.enqueuedAt) {
|
||||
deduped.set(key, task)
|
||||
}
|
||||
}
|
||||
|
||||
const uniqueTasks = Array.from(deduped.values()).sort((a, b) => a.enqueuedAt - b.enqueuedAt)
|
||||
log('info', `deduped to ${uniqueTasks.length} unique tasks`)
|
||||
|
||||
for (const task of uniqueTasks) {
|
||||
try {
|
||||
await processTask(task)
|
||||
} catch (err) {
|
||||
log('error', 'task failed', { error: err instanceof Error ? err.message : String(err), sessionID: task.sessionID })
|
||||
}
|
||||
}
|
||||
|
||||
clearQueue()
|
||||
log('info', 'batch completed')
|
||||
} finally {
|
||||
releaseLock()
|
||||
}
|
||||
}
|
||||
|
||||
export function startBatchWorker() {
|
||||
log('info', 'batch worker started', { interval: BATCH_INTERVAL_MS })
|
||||
|
||||
// 立即执行一次
|
||||
void runBatch()
|
||||
|
||||
// 定期执行,添加随机抖动避免规律性 429
|
||||
const jitter = Math.floor(Math.random() * 10_000)
|
||||
setTimeout(() => {
|
||||
void runBatch()
|
||||
setInterval(() => {
|
||||
void runBatch()
|
||||
}, BATCH_INTERVAL_MS)
|
||||
}, jitter)
|
||||
}
|
||||
|
||||
// 如果直接运行此文件
|
||||
if (process.argv[1]?.includes('batchWorker')) {
|
||||
startBatchWorker()
|
||||
}
|
||||
|
|
@ -90,10 +90,16 @@ export function parseCommitLog(output: string): Array<{
|
|||
}
|
||||
|
||||
export async function getCommitLog(cwd: string, lastCommit?: string): Promise<string> {
|
||||
const args = lastCommit
|
||||
? ['log', `${lastCommit}..HEAD`, '--format=%H|%aI|%an|%ae|%s']
|
||||
: ['log', '--since=30 days ago', '--format=%H|%aI|%an|%ae|%s']
|
||||
return gitExec(args, cwd)
|
||||
if (lastCommit) {
|
||||
return gitExec(
|
||||
['log', `${lastCommit}..HEAD`, '--max-count=50', '--format=%H|%aI|%an|%ae|%s'],
|
||||
cwd,
|
||||
)
|
||||
}
|
||||
return gitExec(
|
||||
['log', '--since=7 days ago', '--max-count=50', '--format=%H|%aI|%an|%ae|%s'],
|
||||
cwd,
|
||||
)
|
||||
}
|
||||
|
||||
export async function getCommitDiff(cwd: string, commitId: string): Promise<string> {
|
||||
|
|
|
|||
|
|
@ -1,68 +1,37 @@
|
|||
/**
|
||||
* Raw Dump 主入口
|
||||
* 提供非阻塞的上报 API,与框架解耦
|
||||
* 队列模式:主进程只 enqueue,单 batch worker 顺序消费
|
||||
*/
|
||||
|
||||
import { spawnRawDumpWorker } from './spawn.js'
|
||||
import { enqueue } from './queue.js'
|
||||
import { spawnBatchWorker } from './spawn.js'
|
||||
|
||||
const SPAWNED_LIMIT = 1024
|
||||
const spawned = new Set<string>()
|
||||
|
||||
function rememberSpawned(key: string) {
|
||||
if (spawned.size >= SPAWNED_LIMIT) {
|
||||
const target = Math.floor(SPAWNED_LIMIT / 2)
|
||||
let dropped = 0
|
||||
for (const k of spawned) {
|
||||
if (dropped >= target) break
|
||||
spawned.delete(k)
|
||||
dropped++
|
||||
}
|
||||
}
|
||||
spawned.add(key)
|
||||
}
|
||||
let batchWorkerSpawned = false
|
||||
|
||||
function isEnabled(): boolean {
|
||||
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
|
||||
}
|
||||
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
|
||||
return true
|
||||
}
|
||||
|
||||
function ensureBatchWorker() {
|
||||
if (batchWorkerSpawned) return
|
||||
batchWorkerSpawned = true
|
||||
spawnBatchWorker()
|
||||
}
|
||||
|
||||
/**
|
||||
* 上报一轮对话的 Conversation + Summary + Commits
|
||||
* 非阻塞:通过 spawn detached 子进程执行,主进程立即返回
|
||||
*
|
||||
* @param sessionID 会话 ID
|
||||
* @param messageID 当前 assistant message 的 UUID
|
||||
* @param directory 工作目录(用于 git diff 和 repo 信息)
|
||||
* 上报一轮对话
|
||||
* 只写入队列,由 batch worker 顺序消费
|
||||
*/
|
||||
export function reportTurn(sessionID: string, messageID: string, directory: string): void {
|
||||
if (!isEnabled()) return
|
||||
|
||||
const key = `${sessionID}:${messageID}`
|
||||
if (spawned.has(key)) return
|
||||
rememberSpawned(key)
|
||||
|
||||
spawnRawDumpWorker({
|
||||
sessionID,
|
||||
messageID,
|
||||
directory,
|
||||
})
|
||||
enqueue({ sessionID, messageID, directory })
|
||||
ensureBatchWorker()
|
||||
}
|
||||
|
||||
/**
|
||||
* 批量上报(用于会话结束时补报)
|
||||
* 非阻塞
|
||||
*/
|
||||
export function reportSession(sessionID: string, directory: string): void {
|
||||
if (!isEnabled()) return
|
||||
// 使用一个特殊 messageID 表示 summary + commits 上报
|
||||
spawnRawDumpWorker({
|
||||
sessionID,
|
||||
messageID: '__summary__',
|
||||
directory,
|
||||
})
|
||||
enqueue({ sessionID, messageID: '__summary__', directory })
|
||||
ensureBatchWorker()
|
||||
}
|
||||
|
|
|
|||
39
src/services/rawDump/logger.ts
Normal file
39
src/services/rawDump/logger.ts
Normal file
|
|
@ -0,0 +1,39 @@
|
|||
/**
|
||||
* Raw Dump 日志模块
|
||||
* 通过环境变量开关控制,默认关闭,与业务逻辑完全解耦
|
||||
*/
|
||||
|
||||
import { appendFileSync } from 'node:fs'
|
||||
import os from 'node:os'
|
||||
import path from 'node:path'
|
||||
|
||||
const LOG_FILE = path.join(os.homedir(), '.claude', 'csc-raw-dump.log')
|
||||
|
||||
function isDebugEnabled(): boolean {
|
||||
const v = process.env.CSC_RAW_DUMP_DEBUG
|
||||
return v === '1' || v === 'true'
|
||||
}
|
||||
|
||||
export function createLogger(prefix: string) {
|
||||
const enabled = isDebugEnabled()
|
||||
|
||||
function write(level: string, msg: string, meta?: Record<string, unknown>) {
|
||||
if (!enabled) return
|
||||
const timestamp = new Date().toISOString()
|
||||
const metaStr = meta ? ` ${JSON.stringify(meta)}` : ''
|
||||
const line = `[${timestamp}] [${prefix}:${level}] ${msg}${metaStr}\n`
|
||||
console.error(line.trimEnd())
|
||||
try {
|
||||
appendFileSync(LOG_FILE, line)
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
debug: (msg: string, meta?: Record<string, unknown>) => write('debug', msg, meta),
|
||||
info: (msg: string, meta?: Record<string, unknown>) => write('info', msg, meta),
|
||||
warn: (msg: string, meta?: Record<string, unknown>) => write('warn', msg, meta),
|
||||
error: (msg: string, meta?: Record<string, unknown>) => write('error', msg, meta),
|
||||
}
|
||||
}
|
||||
87
src/services/rawDump/queue.ts
Normal file
87
src/services/rawDump/queue.ts
Normal file
|
|
@ -0,0 +1,87 @@
|
|||
/**
|
||||
* Raw Dump 任务队列
|
||||
* 主进程只写队列,独立 batch worker 顺序消费
|
||||
*/
|
||||
|
||||
import { appendFileSync, readFileSync, writeFileSync } from 'node:fs'
|
||||
import os from 'node:os'
|
||||
import path from 'node:path'
|
||||
|
||||
const QUEUE_FILE = path.join(os.homedir(), '.claude', 'csc-raw-dump-queue.jsonl')
|
||||
const LOCK_FILE = path.join(os.homedir(), '.claude', 'csc-raw-dump.lock')
|
||||
|
||||
export interface QueueTask {
|
||||
sessionID: string
|
||||
messageID: string
|
||||
directory: string
|
||||
enqueuedAt: number
|
||||
}
|
||||
|
||||
export function enqueue(task: Omit<QueueTask, 'enqueuedAt'>): void {
|
||||
const item: QueueTask = { ...task, enqueuedAt: Date.now() }
|
||||
try {
|
||||
appendFileSync(QUEUE_FILE, JSON.stringify(item) + '\n', 'utf-8')
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
|
||||
export function readQueue(): QueueTask[] {
|
||||
try {
|
||||
const text = readFileSync(QUEUE_FILE, 'utf-8')
|
||||
return text
|
||||
.split('\n')
|
||||
.filter(Boolean)
|
||||
.map((line) => {
|
||||
try {
|
||||
return JSON.parse(line) as QueueTask
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
})
|
||||
.filter((t): t is QueueTask => t !== null)
|
||||
} catch {
|
||||
return []
|
||||
}
|
||||
}
|
||||
|
||||
export function clearQueue(): void {
|
||||
try {
|
||||
writeFileSync(QUEUE_FILE, '', 'utf-8')
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
|
||||
export function acquireLock(): boolean {
|
||||
try {
|
||||
// 简单文件锁:如果 lock 文件存在且 60 秒内,认为已有 worker
|
||||
try {
|
||||
const stat = readFileSync(LOCK_FILE, 'utf-8')
|
||||
const pid = parseInt(stat, 10)
|
||||
if (!isNaN(pid) && pid !== process.pid) {
|
||||
// 检查进程是否还在运行
|
||||
try {
|
||||
process.kill(pid, 0)
|
||||
return false // 已有 worker 在运行
|
||||
} catch {
|
||||
// 进程已退出,可以抢占锁
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
// lock 文件不存在
|
||||
}
|
||||
writeFileSync(LOCK_FILE, String(process.pid), 'utf-8')
|
||||
return true
|
||||
} catch {
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
export function releaseLock(): void {
|
||||
try {
|
||||
writeFileSync(LOCK_FILE, '', 'utf-8')
|
||||
} catch {
|
||||
// ignore
|
||||
}
|
||||
}
|
||||
|
|
@ -1,24 +1,18 @@
|
|||
/**
|
||||
* Raw Dump Worker 进程启动器
|
||||
* 使用 detached 子进程,确保不阻塞主业务流程
|
||||
* 启动独立的 batch worker 顺序消费队列
|
||||
*/
|
||||
|
||||
import { spawn } from 'node:child_process'
|
||||
import path from 'node:path'
|
||||
import { fileURLToPath } from 'node:url'
|
||||
import { RAW_DUMP_EVENT_ENV_KEY, type RawDumpEventPayload } from './types.js'
|
||||
|
||||
export function getRawDumpEventEnvKey(): string {
|
||||
return RAW_DUMP_EVENT_ENV_KEY
|
||||
}
|
||||
|
||||
export function spawnRawDumpWorker(payload: RawDumpEventPayload): void {
|
||||
export function spawnBatchWorker(): void {
|
||||
const entry = process.execPath
|
||||
const isDev = path.basename(entry).toLowerCase().startsWith('bun')
|
||||
|
||||
// 计算 worker.ts 的绝对路径
|
||||
const __dirname = path.dirname(fileURLToPath(import.meta.url))
|
||||
const workerPath = path.resolve(__dirname, 'worker.ts')
|
||||
const workerPath = path.resolve(__dirname, 'batchWorker.ts')
|
||||
|
||||
const args = isDev
|
||||
? ['run', workerPath]
|
||||
|
|
@ -28,15 +22,10 @@ export function spawnRawDumpWorker(payload: RawDumpEventPayload): void {
|
|||
detached: true,
|
||||
windowsHide: true,
|
||||
stdio: 'ignore',
|
||||
env: {
|
||||
...process.env,
|
||||
[RAW_DUMP_EVENT_ENV_KEY]: JSON.stringify(payload),
|
||||
},
|
||||
})
|
||||
|
||||
child.on('error', (err) => {
|
||||
// 静默处理,不影响主进程
|
||||
console.error('[raw-dump] spawn error:', err.message)
|
||||
console.error('[raw-dump] batch worker spawn error:', err.message)
|
||||
})
|
||||
|
||||
child.unref()
|
||||
|
|
|
|||
|
|
@ -29,6 +29,7 @@ import {
|
|||
parseCommitLog,
|
||||
toCommitComment,
|
||||
} from './git.js'
|
||||
import { createLogger } from './logger.js'
|
||||
import { readState, writeState } from './state.js'
|
||||
import { RAW_DUMP_EVENT_ENV_KEY, type RawDumpEventPayload } from './types.js'
|
||||
import type {
|
||||
|
|
@ -38,12 +39,7 @@ import type {
|
|||
SummaryPayload,
|
||||
} from './types.js'
|
||||
|
||||
// 简单的日志输出到 stderr,不依赖主进程日志系统
|
||||
function log(level: string, msg: string, meta?: Record<string, unknown>) {
|
||||
const timestamp = new Date().toISOString()
|
||||
const metaStr = meta ? ` ${JSON.stringify(meta)}` : ''
|
||||
console.error(`[${timestamp}] [raw-dump:${level}] ${msg}${metaStr}`)
|
||||
}
|
||||
const log = createLogger('raw-dump')
|
||||
|
||||
function formatIso(ms: number | undefined): string {
|
||||
if (!ms) return ''
|
||||
|
|
@ -82,18 +78,42 @@ async function postJson(
|
|||
const url = getRawDumpUrl(baseUrl, endpoint)
|
||||
log('debug', `POST ${endpoint}`, { url })
|
||||
|
||||
const res = await fetch(url, {
|
||||
method: 'POST',
|
||||
headers,
|
||||
body: JSON.stringify(body),
|
||||
})
|
||||
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 })
|
||||
await new Promise((r) => setTimeout(r, delay))
|
||||
}
|
||||
|
||||
if (!res.ok) {
|
||||
const text = await res.text().catch(() => '')
|
||||
throw new Error(`${endpoint} failed: ${res.status} ${text}`)
|
||||
try {
|
||||
const res = await fetch(url, {
|
||||
method: 'POST',
|
||||
headers,
|
||||
body: JSON.stringify(body),
|
||||
})
|
||||
|
||||
if (res.ok) {
|
||||
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) })
|
||||
lastError = new Error(`${endpoint} failed: ${res.status} ${text}`)
|
||||
continue
|
||||
}
|
||||
throw new Error(`${endpoint} failed: ${res.status} ${text}`)
|
||||
} catch (err) {
|
||||
lastError = err instanceof Error ? err : new Error(String(err))
|
||||
// 网络错误也重试
|
||||
log('warn', `${endpoint} network error, will retry`, { attempt, error: lastError.message })
|
||||
}
|
||||
}
|
||||
|
||||
log('debug', `POST ${endpoint} ok`, { status: res.status })
|
||||
throw lastError || new Error(`${endpoint} failed after retries`)
|
||||
}
|
||||
|
||||
function parseUser(accessPayload: JwtPayload, refreshPayload?: JwtPayload | null) {
|
||||
|
|
@ -114,12 +134,15 @@ function detectOs(): string {
|
|||
return map[process.platform] ?? process.platform
|
||||
}
|
||||
|
||||
async function auth() {
|
||||
export async function auth() {
|
||||
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 })
|
||||
|
||||
// Token 刷新
|
||||
if (creds.refresh_token && !isCoStrictTokenValid(creds)) {
|
||||
log('debug', 'token expired, refreshing...')
|
||||
const next = await refreshCoStrictToken({
|
||||
baseUrl: creds.base_url,
|
||||
refreshToken: creds.refresh_token,
|
||||
|
|
@ -134,6 +157,7 @@ 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')
|
||||
}
|
||||
|
||||
const headers = new Headers()
|
||||
|
|
@ -169,34 +193,60 @@ 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 })
|
||||
|
||||
return {
|
||||
baseUrl: resolveRawDumpBaseUrl(creds.base_url),
|
||||
baseUrl,
|
||||
headers,
|
||||
user: parseUser(accessPayload, refreshPayload),
|
||||
user,
|
||||
clientId,
|
||||
version,
|
||||
}
|
||||
}
|
||||
|
||||
// 从 JSONL 文件加载会话消息
|
||||
async function loadSessionMessages(sessionDir: string, sessionId: string) {
|
||||
const filePath = path.join(sessionDir, `${sessionId}.jsonl`)
|
||||
// csc 的会话文件名可能是 ses_{hash}.jsonl 或 {uuid}.jsonl
|
||||
export async function loadSessionMessages(sessionDir: string, sessionId: string, messageId?: string) {
|
||||
try {
|
||||
const text = await fs.readFile(filePath, 'utf-8')
|
||||
return text
|
||||
.split('\n')
|
||||
.filter(Boolean)
|
||||
.map((line) => {
|
||||
try {
|
||||
return JSON.parse(line)
|
||||
} catch {
|
||||
return null
|
||||
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) })
|
||||
|
||||
for (const file of jsonlFiles) {
|
||||
const filePath = path.join(sessionDir, file)
|
||||
try {
|
||||
const text = await fs.readFile(filePath, 'utf-8')
|
||||
const lines = text
|
||||
.split('\n')
|
||||
.filter(Boolean)
|
||||
.map((line) => {
|
||||
try {
|
||||
return JSON.parse(line)
|
||||
} catch {
|
||||
return null
|
||||
}
|
||||
})
|
||||
.filter((m): m is Record<string, unknown> => m !== null)
|
||||
|
||||
// 检查是否包含目标 sessionId 或 messageId
|
||||
const hasSession = lines.some(
|
||||
(m) => m.sessionId === sessionId || m.session_id === sessionId || m.uuid === sessionId,
|
||||
)
|
||||
const hasMessage = messageId ? lines.some((m) => m.uuid === messageId || (m.message as Record<string, unknown>)?.id === messageId) : false
|
||||
if (hasSession || hasMessage) {
|
||||
log('debug', 'loaded messages from file', { file, count: lines.length, hasSession, hasMessage })
|
||||
return lines
|
||||
}
|
||||
})
|
||||
.filter((m): m is Record<string, unknown> => m !== null)
|
||||
} catch {
|
||||
// ignore per-file errors
|
||||
}
|
||||
}
|
||||
} catch {
|
||||
return []
|
||||
// ignore dir read errors
|
||||
}
|
||||
return []
|
||||
}
|
||||
|
||||
function findMessage(
|
||||
|
|
@ -284,7 +334,7 @@ function extractError(msg: Record<string, unknown>) {
|
|||
return { error_code: errorCode, error_reason: message }
|
||||
}
|
||||
|
||||
async function uploadConversation(
|
||||
export async function uploadConversation(
|
||||
payload: {
|
||||
sessionID: string
|
||||
messageID: string
|
||||
|
|
@ -294,13 +344,24 @@ async function uploadConversation(
|
|||
authData: Awaited<ReturnType<typeof auth>>,
|
||||
state: Awaited<ReturnType<typeof readState>>,
|
||||
): Promise<boolean> {
|
||||
const assistant = findMessage(payload.messages, payload.messageID)
|
||||
log('debug', 'uploadConversation start', { messageID: payload.messageID, messageCount: payload.messages.length })
|
||||
|
||||
let assistant = findMessage(payload.messages, payload.messageID)
|
||||
if (!assistant || assistant.type !== 'assistant') {
|
||||
log('warn', 'assistant message not found', { messageID: payload.messageID })
|
||||
return false
|
||||
// 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 })
|
||||
assistant = lastAssistant
|
||||
} else {
|
||||
log('warn', 'assistant message not found', { messageID: payload.messageID, foundType: assistant?.type })
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
const requestID = ((assistant.message as Record<string, unknown>)?.id as string) || payload.messageID
|
||||
const requestID = ((assistant.message as Record<string, unknown>)?.id as string) || String(assistant.uuid) || payload.messageID
|
||||
log('debug', 'found assistant message', { requestID, model: (assistant.message as Record<string, unknown>)?.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 })
|
||||
|
|
@ -308,17 +369,24 @@ async function uploadConversation(
|
|||
}
|
||||
|
||||
const user = findParentUserMessage(payload.messages, assistant)
|
||||
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
|
||||
const toolDiff = extractToolDiff(assistant)
|
||||
log('debug', 'extracted tool diff', { toolDiffLength: toolDiff.diff.length, toolDiffLines: toolDiff.diff_lines, toolDiffFiles: toolDiff.files.length })
|
||||
|
||||
const rawDiff = toolDiff.diff || (await getWorkingTreeDiff(payload.directory))
|
||||
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<string, unknown>).ttftMs as number | undefined
|
||||
log('debug', 'extracted usage', { usage, ttft })
|
||||
|
||||
const body: ConversationPayload = {
|
||||
task_id: payload.sessionID,
|
||||
|
|
@ -343,13 +411,14 @@ async function uploadConversation(
|
|||
...extractError(assistant),
|
||||
}
|
||||
|
||||
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 })
|
||||
log('info', 'conversation uploaded', { task_id: payload.sessionID, request_id: requestID, upstream_tokens: body.upstream_tokens, downstream_tokens: body.downstream_tokens })
|
||||
return true
|
||||
}
|
||||
|
||||
async function uploadSummary(
|
||||
export async function uploadSummary(
|
||||
payload: {
|
||||
sessionID: string
|
||||
directory: string
|
||||
|
|
@ -357,8 +426,10 @@ async function uploadSummary(
|
|||
},
|
||||
authData: Awaited<ReturnType<typeof auth>>,
|
||||
): Promise<void> {
|
||||
log('debug', 'uploadSummary start', { sessionID: payload.sessionID, messageCount: payload.messages.length })
|
||||
const repoInfo = await getRepoInfo(payload.directory)
|
||||
const rawDiff = await getWorkingTreeDiff(payload.directory)
|
||||
log('debug', 'summary repo info', { repo_addr: repoInfo.repo_addr, repo_branch: repoInfo.repo_branch, diffLength: rawDiff.length })
|
||||
|
||||
const assistants = payload.messages.filter((m) => m.type === 'assistant')
|
||||
const { upstream_tokens, downstream_tokens } = assistants.reduce(
|
||||
|
|
@ -397,33 +468,44 @@ async function uploadSummary(
|
|||
}
|
||||
|
||||
await postJson(authData.baseUrl, authData.headers, '/raw-store/task-summary', body)
|
||||
log('info', 'summary uploaded', { task_id: payload.sessionID })
|
||||
log('info', 'summary uploaded', { task_id: payload.sessionID, upstream_tokens: body.upstream_tokens, downstream_tokens: body.downstream_tokens, diff_lines: body.diff_lines })
|
||||
}
|
||||
|
||||
async function uploadCommits(
|
||||
export async function uploadCommits(
|
||||
payload: {
|
||||
directory: string
|
||||
},
|
||||
authData: Awaited<ReturnType<typeof auth>>,
|
||||
state: Awaited<ReturnType<typeof readState>>,
|
||||
): Promise<number> {
|
||||
log('debug', 'uploadCommits start', { directory: payload.directory })
|
||||
const repoInfo = await getRepoInfo(payload.directory)
|
||||
if (!repoInfo.repo_addr || !repoInfo.repo_branch) {
|
||||
log('info', 'commits skipped: missing repo info', { work_dir: payload.directory })
|
||||
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)' })
|
||||
|
||||
const logText = await getCommitLog(payload.directory, lastCommit)
|
||||
const commits = parseCommitLog(logText)
|
||||
const allCommits = parseCommitLog(logText)
|
||||
// 限制每次最多上报 50 个 commit,避免触发限流
|
||||
const commits = allCommits.slice(0, 50)
|
||||
log('debug', 'parsed commits', { total: allCommits.length, sending: commits.length })
|
||||
|
||||
if (!commits.length) {
|
||||
log('info', 'commits skipped: no new commits', { work_dir: payload.directory })
|
||||
return 0
|
||||
}
|
||||
|
||||
for (const commit of commits) {
|
||||
for (let i = 0; i < commits.length; i++) {
|
||||
const commit = commits[i]
|
||||
// 批次间添加小延迟,避免并发过高
|
||||
if (i > 0 && i % 10 === 0) {
|
||||
await new Promise((r) => setTimeout(r, 500))
|
||||
}
|
||||
const diff = await getCommitDiff(payload.directory, commit.commit_id)
|
||||
const body: CommitPayload = {
|
||||
commit_id: commit.commit_id,
|
||||
|
|
@ -444,10 +526,11 @@ async function uploadCommits(
|
|||
subject: commit.subject,
|
||||
}
|
||||
await postJson(authData.baseUrl, authData.headers, '/raw-store/commit', body)
|
||||
log('info', 'commit uploaded', { commit_id: commit.commit_id })
|
||||
// 每成功一个 commit 立即更新 state,避免失败后全部重传
|
||||
state.commits[stateKey] = commit.commit_id
|
||||
log('info', 'commit uploaded', { commit_id: commit.commit_id, progress: `${i + 1}/${commits.length}` })
|
||||
}
|
||||
|
||||
state.commits[stateKey] = commits[0]!.commit_id
|
||||
return commits.length
|
||||
}
|
||||
|
||||
|
|
@ -457,10 +540,23 @@ function parseWorkerPayload(): RawDumpEventPayload {
|
|||
return JSON.parse(raw) as RawDumpEventPayload
|
||||
}
|
||||
|
||||
function getSessionDirectory(directory: string, sessionID: string): string {
|
||||
// csc 的会话文件通常在项目的 .claude/sessions/ 目录下
|
||||
// 尝试从传入的 directory 或环境变量推断
|
||||
export function getClaudeConfigHomeDir(): string {
|
||||
return process.env.CLAUDE_CONFIG_HOME || path.join(os.homedir(), '.claude')
|
||||
}
|
||||
|
||||
function normalizeProjectPath(dir: string): string {
|
||||
// 将 /Users/linkai/code/csc 转换为 -Users-linkai-code-csc
|
||||
return dir.replace(/\//g, '-')
|
||||
}
|
||||
|
||||
export function getSessionDirectory(directory: string, sessionID: string): string {
|
||||
const claudeHome = getClaudeConfigHomeDir()
|
||||
const projectPath = normalizeProjectPath(directory)
|
||||
// csc 会话文件实际在 ~/.claude/projects/{project-path}/
|
||||
const candidates = [
|
||||
path.join(claudeHome, 'projects', projectPath),
|
||||
path.join(claudeHome, 'transcripts'),
|
||||
path.join(claudeHome, 'sessions'),
|
||||
path.join(directory, '.claude', 'sessions'),
|
||||
path.join(directory, '.claude'),
|
||||
directory,
|
||||
|
|
@ -472,34 +568,51 @@ function getSessionDirectory(directory: string, sessionID: string): string {
|
|||
export async function runRawDumpWorker() {
|
||||
try {
|
||||
const payload = parseWorkerPayload()
|
||||
log('info', 'worker started', { session_id: payload.sessionID, message_id: payload.messageID })
|
||||
log('info', '=== WORKER STARTED ===', { session_id: payload.sessionID, message_id: payload.messageID, directory: payload.directory })
|
||||
|
||||
const sessionDir = getSessionDirectory(payload.directory, payload.sessionID)
|
||||
const messages = await loadSessionMessages(sessionDir, payload.sessionID)
|
||||
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 })
|
||||
|
||||
if (messages.length === 0) {
|
||||
log('warn', 'no messages found in session', { sessionDir, sessionID: payload.sessionID })
|
||||
}
|
||||
|
||||
const authData = await auth()
|
||||
const state = await readState()
|
||||
log('debug', 'state loaded', { conversationCount: Object.keys(state.conversation).length, commitCount: Object.keys(state.commits).length })
|
||||
|
||||
log('debug', 'starting uploadConversation...')
|
||||
const conversationUploaded = await uploadConversation(
|
||||
{ ...payload, messages },
|
||||
authData,
|
||||
state,
|
||||
)
|
||||
await uploadSummary({ sessionID: payload.sessionID, directory: payload.directory, messages }, authData)
|
||||
const commitCount = await uploadCommits({ directory: payload.directory }, authData, state)
|
||||
await writeState(state)
|
||||
log('debug', 'uploadConversation done', { conversationUploaded })
|
||||
|
||||
log('info', 'worker completed', {
|
||||
log('debug', 'starting uploadSummary...')
|
||||
await uploadSummary({ sessionID: payload.sessionID, directory: payload.directory, messages }, authData)
|
||||
log('debug', 'uploadSummary done')
|
||||
|
||||
log('debug', 'starting uploadCommits...')
|
||||
const commitCount = await uploadCommits({ directory: payload.directory }, authData, state)
|
||||
log('debug', 'uploadCommits done', { commitCount })
|
||||
|
||||
await writeState(state)
|
||||
log('debug', 'state saved')
|
||||
|
||||
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,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,23 +5,42 @@
|
|||
*/
|
||||
|
||||
import { reportTurn } from '../services/rawDump/index.js'
|
||||
import { getSessionProjectDir } from '../bootstrap/state.js'
|
||||
import { getSessionProjectDir, getSessionId } from '../bootstrap/state.js'
|
||||
import type { Message } from '../types/message.js'
|
||||
|
||||
/**
|
||||
* 创建 session turn 上报器
|
||||
* 在 assistant message 完成时调用,触发异步 raw-dump 上报
|
||||
* main.tsx 在 onTurnComplete 回调中调用返回的函数
|
||||
*/
|
||||
export function createSessionTurnUploader(): void {
|
||||
// Stub: 实际调用方应直接调用 reportTurn()
|
||||
// 此处保留空实现以兼容现有代码
|
||||
export function createSessionTurnUploader(): (messages: Message[]) => void {
|
||||
return (messages: Message[]) => {
|
||||
const sessionId = getSessionId()
|
||||
if (!sessionId) {
|
||||
console.error('[raw-dump] skip: no sessionId')
|
||||
return
|
||||
}
|
||||
|
||||
// 找到最后一个 assistant message
|
||||
const lastAssistant = [...messages].reverse().find((m) => m.type === 'assistant')
|
||||
if (!lastAssistant) {
|
||||
console.error('[raw-dump] skip: no assistant message in turn')
|
||||
return
|
||||
}
|
||||
|
||||
const messageId = String(lastAssistant.uuid || '')
|
||||
if (!messageId) {
|
||||
console.error('[raw-dump] skip: assistant message has no uuid')
|
||||
return
|
||||
}
|
||||
|
||||
const directory = getSessionProjectDir() || process.cwd()
|
||||
console.error('[raw-dump] trigger reportTurn', { sessionId, messageId, directory })
|
||||
reportTurn(sessionId, messageId, directory)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 上报单个 turn 的数据
|
||||
* 由 query.ts 或 costrict/provider/index.ts 在 streaming 结束后调用
|
||||
*
|
||||
* @param sessionId 会话 ID
|
||||
* @param assistantMessageUuid 刚完成的 assistant message UUID
|
||||
* 手动上报单个 turn(供外部直接调用)
|
||||
*/
|
||||
export function uploadSessionTurn(sessionId: string, assistantMessageUuid: string): void {
|
||||
const directory = getSessionProjectDir() || process.cwd()
|
||||
|
|
|
|||
Loading…
Reference in New Issue
Block a user