启动时自动创建raw-dump目录

This commit is contained in:
zbc 2026-05-27 16:11:14 +08:00
parent 14d8e59d8c
commit 2e6dd38603
6 changed files with 66 additions and 59 deletions

View File

@ -368,9 +368,11 @@ export CSC_RAW_DUMP_DIR=/tmp/raw-dump-debug
### 状态文件 ### 状态文件
``` ```
~/.claude/raw-dump/csc-state.json ~/.claude/raw-dump/csc-state.json
~/.claude/raw-dump/csc-work-queue.jsonl
~/.claude/raw-dump/csc-dead-letter.jsonl
``` ```
内容格式 **csc-state.json** — 去重状态
```json ```json
{ {
"conversation": { "conversation": {
@ -378,7 +380,7 @@ export CSC_RAW_DUMP_DIR=/tmp/raw-dump-debug
"session-id-1:msg-uuid-2": true "session-id-1:msg-uuid-2": true
}, },
"summary": { "summary": {
"session-id-1": 1747123456789 "session-id-1": "2026-05-12T10:30:00.000Z"
}, },
"commits": { "commits": {
"git@github.com:org/repo.git#main#/Users/xxx/code/repo": "abc123def" "git@github.com:org/repo.git#main#/Users/xxx/code/repo": "abc123def"
@ -386,6 +388,12 @@ export CSC_RAW_DUMP_DIR=/tmp/raw-dump-debug
} }
``` ```
**csc-dead-letter.jsonl** — 上报失败记录(超过最大重试次数的任务),每行一个 JSON 对象:
```jsonl
{"sessionID":"...","messageID":"...","directory":"...","attemptCount":4,"error":"...","endpoint":"/raw-store/task-conversation","failedAt":"2026-05-12T10:30:00.000Z"}
{"sessionID":"...","messageID":"__summary__","directory":"...","attemptCount":4,"error":"...","endpoint":"/raw-store/task-summary","failedAt":"2026-05-12T10:31:00.000Z"}
```
### 日志文件 ### 日志文件
``` ```
~/.claude/raw-dump/csc-raw-dump.log ~/.claude/raw-dump/csc-raw-dump.log
@ -461,7 +469,10 @@ tail -f ~/.claude/raw-dump/csc-raw-dump.log
cat ~/.claude/raw-dump/csc-state.json cat ~/.claude/raw-dump/csc-state.json
# 查看队列文件 # 查看队列文件
cat ~/.claude/raw-dump/csc-work-queue.json cat ~/.claude/raw-dump/csc-work-queue.jsonl
# 查看失败记录dead letter
cat ~/.claude/raw-dump/csc-dead-letter.jsonl
# 查看是否有 worker 在运行(锁文件) # 查看是否有 worker 在运行(锁文件)
cat ~/.claude/raw-dump/csc-work-queue.lock cat ~/.claude/raw-dump/csc-work-queue.lock

View File

@ -220,7 +220,7 @@ async function runBatch() {
} }
const uniqueTasks = Array.from(deduped.values()).sort( const uniqueTasks = Array.from(deduped.values()).sort(
(a, b) => a.enqueuedAt - b.enqueuedAt, (a, b) => a.enqueuedAt.localeCompare(b.enqueuedAt),
) )
log.info(`deduped to ${uniqueTasks.length} unique tasks`) log.info(`deduped to ${uniqueTasks.length} unique tasks`)

View File

@ -4,6 +4,7 @@
*/ */
import { import {
ensureRawDumpDirCreated,
getRawDumpMode, getRawDumpMode,
RAW_DUMP_MODE, RAW_DUMP_MODE,
} from './localStorage.js' } from './localStorage.js'
@ -84,13 +85,14 @@ function shouldEnqueue(sessionID: string, messageID: string): boolean {
* *
* batch worker * batch worker
*/ */
export function reportTurn( export async function reportTurn(
sessionID: string, sessionID: string,
messageID: string, messageID: string,
directory: string, directory: string,
): void { ): Promise<void> {
if (!isEnabled()) return if (!isEnabled()) return
if (!shouldEnqueue(sessionID, messageID)) return if (!shouldEnqueue(sessionID, messageID)) return
ensureRawDumpDirCreated()
enqueue({ sessionID, messageID, directory }) enqueue({ sessionID, messageID, directory })
ensureBatchWorker() ensureBatchWorker()
} }
@ -100,9 +102,13 @@ export function reportTurn(
* 使 messageID '__summary__' conversation * 使 messageID '__summary__' conversation
* batch worker * batch worker
*/ */
export function reportSession(sessionID: string, directory: string): void { export async function reportSession(
sessionID: string,
directory: string,
): Promise<void> {
if (!isEnabled()) return if (!isEnabled()) return
if (!shouldEnqueue(sessionID, '__summary__')) return if (!shouldEnqueue(sessionID, '__summary__')) return
ensureRawDumpDirCreated()
enqueue({ sessionID, messageID: '__summary__', directory }) enqueue({ sessionID, messageID: '__summary__', directory })
ensureBatchWorker() ensureBatchWorker()
} }
@ -123,11 +129,11 @@ const STATS_DEBOUNCE_MS = 60_000 // 同一 session 1 分钟内不重复 enqueue
* session数conversation数token数 * session数conversation数token数
* messageID '__statistics__' * messageID '__statistics__'
*/ */
export function reportStatistics( export async function reportStatistics(
sessionID: string, sessionID: string,
directory: string, directory: string,
data: StatisticsData, data: StatisticsData,
): void { ): Promise<void> {
if (!isEnabled()) return if (!isEnabled()) return
const key = `${sessionID}:__statistics__` const key = `${sessionID}:__statistics__`
const now = Date.now() const now = Date.now()
@ -137,6 +143,7 @@ export function reportStatistics(
return return
} }
lastReportStatsMap.set(key, now) lastReportStatsMap.set(key, now)
ensureRawDumpDirCreated()
enqueue({ enqueue({
sessionID, sessionID,
messageID: '__statistics__', messageID: '__statistics__',

View File

@ -24,15 +24,21 @@ export const RAW_DUMP_MODE = {
BOTH: 3, BOTH: 3,
} as const } as const
let rawDumpDirCreated = false
export async function ensureRawDumpDirCreated(): Promise<void> {
if (rawDumpDirCreated) return
const RAW_DUMP_DIR = path.join(os.homedir(), '.claude', 'raw-dump')
await fs.mkdir(RAW_DUMP_DIR, { recursive: true })
rawDumpDirCreated = true
}
/** /**
* *
* CSC_RAW_DUMP_DIR * CSC_RAW_DUMP_DIR
*/ */
export function getLocalDumpDir(): string { export function getLocalDumpDir(): string {
return (process.env.CSC_RAW_DUMP_DIR || DEFAULT_LOCAL_DIR).replace( return (process.env.CSC_RAW_DUMP_DIR || DEFAULT_LOCAL_DIR).replace(/\/$/, '')
/\/$/,
'',
)
} }
/** /**
@ -62,13 +68,17 @@ export async function writeLocalDump(
body: Record<string, unknown>, body: Record<string, unknown>,
): Promise<void> { ): Promise<void> {
const dir = getLocalDumpDir() const dir = getLocalDumpDir()
const taskId = (body.task_id as string) || (body.commit_id as string)|| 'unknown' const taskId =
(body.task_id as string) || (body.commit_id as string) || 'unknown'
const taskDir = path.join(dir, type, taskId) const taskDir = path.join(dir, type, taskId)
await fs.mkdir(taskDir, { recursive: true }) await fs.mkdir(taskDir, { recursive: true })
const timestamp = new Date().toISOString().replace(/[:.]/g, '-').slice(0, -5) const timestamp = new Date().toISOString().replace(/[:.]/g, '-').slice(0, -5)
const requestId = const requestId =
(body.request_id as string) || (body.commit_id as string) || (body.task_id as string) || 'unknown' (body.request_id as string) ||
(body.commit_id as string) ||
(body.task_id as string) ||
'unknown'
const filename = `${timestamp}-${requestId}.json` const filename = `${timestamp}-${requestId}.json`
const filePath = path.join(taskDir, filename) const filePath = path.join(taskDir, filename)

View File

@ -14,12 +14,7 @@ import os from 'node:os'
import path from 'path' import path from 'path'
import type { StatisticsData } from './index.js' import type { StatisticsData } from './index.js'
const QUEUE_FILE = path.join( const QUEUE_FILE = path.join(os.homedir(), '.claude', 'raw-dump', 'csc-work-queue.jsonl')
os.homedir(),
'.claude',
'raw-dump',
'csc-work-queue.jsonl',
)
const LOCK_FILE = path.join(os.homedir(), '.claude', 'raw-dump', 'csc-work-queue.lock') const LOCK_FILE = path.join(os.homedir(), '.claude', 'raw-dump', 'csc-work-queue.lock')
export const MAX_ATTEMPTS = 4 // 最多尝试次数 export const MAX_ATTEMPTS = 4 // 最多尝试次数
@ -28,7 +23,7 @@ export interface QueueTask {
sessionID: string sessionID: string
messageID: string messageID: string
directory: string directory: string
enqueuedAt: number enqueuedAt: string // RFC3339 format, e.g. "2026-05-27T10:00:00.000Z"
attemptCount: number attemptCount: number
statsData?: StatisticsData statsData?: StatisticsData
} }
@ -74,7 +69,7 @@ export async function flushQueue(): Promise<void> {
* *
*/ */
export function enqueue(task: Omit<QueueTask, 'enqueuedAt' | 'attemptCount'>): void { export function enqueue(task: Omit<QueueTask, 'enqueuedAt' | 'attemptCount'>): void {
const item: QueueTask = { ...task, enqueuedAt: Date.now(), attemptCount: 0 } const item: QueueTask = { ...task, enqueuedAt: new Date().toISOString(), attemptCount: 0 }
queue.push(item) queue.push(item)
// 同步写文件,不阻塞主进程 event loop // 同步写文件,不阻塞主进程 event loop
fs.writeFile(QUEUE_FILE, JSON.stringify(item) + '\n', { flag: 'a', encoding: 'utf-8' }).catch(() => {}) fs.writeFile(QUEUE_FILE, JSON.stringify(item) + '\n', { flag: 'a', encoding: 'utf-8' }).catch(() => {})

View File

@ -30,7 +30,11 @@ import {
toCommitComment, toCommitComment,
} from './git.js' } from './git.js'
import { createLogger } from './logger.js' import { createLogger } from './logger.js'
import { getRawDumpMode, RAW_DUMP_MODE, writeLocalDump } from './localStorage.js' import {
getRawDumpMode,
RAW_DUMP_MODE,
writeLocalDump,
} from './localStorage.js'
import { readState, writeState } from './state.js' import { readState, writeState } from './state.js'
import { RAW_DUMP_EVENT_ENV_KEY, type RawDumpEventPayload } from './types.js' import { RAW_DUMP_EVENT_ENV_KEY, type RawDumpEventPayload } from './types.js'
import type { import type {
@ -103,17 +107,15 @@ function getRawDumpUrl(
} }
/** /**
* JSON POST * Raw Dump API
* - JSON * - JSON
* - /429 3 退 * - /429 3 退
* @param baseUrl API * @param authData baseUrl headers
* @param headers HTTP * @param endpoint API /raw-store/task-conversation
* @param endpoint API
* @param body JSON * @param body JSON
*/ */
async function postJson( async function uploadReport(
baseUrl: string, authData: { baseUrl: string; headers: Headers },
headers: Headers,
endpoint: string, endpoint: string,
body: object, body: object,
): Promise<void> { ): Promise<void> {
@ -143,8 +145,8 @@ async function postJson(
// REMOTE / BOTH 模式下继续执行 remote 上报逻辑 // REMOTE / BOTH 模式下继续执行 remote 上报逻辑
const isAnonymous = !headers.get('Authorization') const isAnonymous = !authData.headers.get('Authorization')
const url = getRawDumpUrl(baseUrl, endpoint, isAnonymous) const url = getRawDumpUrl(authData.baseUrl, endpoint, isAnonymous)
log.debug(`POST ${endpoint}`, { url, isAnonymous }) log.debug(`POST ${endpoint}`, { url, isAnonymous })
let lastError: Error | undefined let lastError: Error | undefined
@ -160,7 +162,7 @@ async function postJson(
try { try {
const res = await fetch(url, { const res = await fetch(url, {
method: 'POST', method: 'POST',
headers, headers: authData.headers,
body: JSON.stringify(body), body: JSON.stringify(body),
signal: controller.signal, signal: controller.signal,
}) })
@ -906,12 +908,7 @@ export async function uploadConversation(
request_id: requestID, request_id: requestID,
bodyKeys: Object.keys(body), bodyKeys: Object.keys(body),
}) })
await postJson( await uploadReport(authData, '/raw-store/task-conversation', body)
authData.baseUrl,
authData.headers,
'/raw-store/task-conversation',
body,
)
state.conversation[key] = true state.conversation[key] = true
log.info('conversation uploaded', { log.info('conversation uploaded', {
task_id: payload.sessionID, task_id: payload.sessionID,
@ -947,9 +944,7 @@ export async function uploadSummary(
const lastReported = state.summary[payload.sessionID] const lastReported = state.summary[payload.sessionID]
if ( if (
lastReported && lastReported &&
Date.now() - Date.now() - new Date(lastReported).getTime() < SUMMARY_DEDUP_WINDOW_MS
new Date(lastReported).getTime() <
SUMMARY_DEDUP_WINDOW_MS
) { ) {
log.info('summary skipped: reported recently', { log.info('summary skipped: reported recently', {
task_id: payload.sessionID, task_id: payload.sessionID,
@ -975,12 +970,7 @@ export async function uploadSummary(
caller: process.env.CSC_RAW_DUMP_CALLER || 'chat', caller: process.env.CSC_RAW_DUMP_CALLER || 'chat',
} }
await postJson( await uploadReport(authData, '/raw-store/task-summary', body)
authData.baseUrl,
authData.headers,
'/raw-store/task-summary',
body,
)
state.summary[payload.sessionID] = new Date().toISOString() state.summary[payload.sessionID] = new Date().toISOString()
log.info('summary uploaded', { task_id: payload.sessionID }) log.info('summary uploaded', { task_id: payload.sessionID })
} }
@ -1056,9 +1046,8 @@ export async function uploadCommits(
subject: commit.subject, subject: commit.subject,
parent_ids: commit.parent_ids, parent_ids: commit.parent_ids,
} }
await postJson( await uploadReport(
authData.baseUrl, authData,
authData.headers,
'/raw-store/commit', '/raw-store/commit',
body as unknown as Record<string, unknown>, body as unknown as Record<string, unknown>,
) )
@ -1115,12 +1104,7 @@ export async function uploadStatistics(
downstream_tokens: payload.downstreamTokens, downstream_tokens: payload.downstreamTokens,
} }
await postJson( await uploadReport(authData, '/raw-store/statistics', body)
authData.baseUrl,
authData.headers,
'/raw-store/statistics',
body,
)
state.summary[key] = new Date().toISOString() state.summary[key] = new Date().toISOString()
log.info('statistics uploaded', { log.info('statistics uploaded', {
task_id: payload.sessionID, task_id: payload.sessionID,
@ -1157,7 +1141,7 @@ function normalizeProjectPath(dir: string): string {
// 将路径中的路径分隔符替换为 -,统一处理 / 和 \ (Windows) // 将路径中的路径分隔符替换为 -,统一处理 / 和 \ (Windows)
// 如 /Users/linkai/code/csc → -Users-linkai-code-csc // 如 /Users/linkai/code/csc → -Users-linkai-code-csc
// 如 D:\shenma\zgsm-ai\csc → D--shenma-zgsm-ai-csc (drive letter 后的 \ 也转为 -) // 如 D:\shenma\zgsm-ai\csc → D--shenma-zgsm-ai-csc (drive letter 后的 \ 也转为 -)
return dir.replace(/:/, '-').replace(/[\/\\]/g, '-') return dir.replace(/:/, '-').replace(/[/\\]/g, '-')
} }
/** /**