From 14d8e59d8c85e8468347dc2b9dcb359967a8f32a Mon Sep 17 00:00:00 2001 From: zbc Date: Wed, 27 May 2026 10:49:21 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E5=A2=9E=E5=8A=A0=E7=BB=9F=E8=AE=A1?= =?UTF-8?q?=E4=BF=A1=E6=81=AF=E5=B9=B6=E4=B8=8A=E6=8A=A5=E9=80=BB=E8=BE=91?= =?UTF-8?q?=EF=BC=8C=E7=94=A8=E4=BA=8E=E5=AF=B9=E8=B4=A6=EF=BC=9B2.=20?= =?UTF-8?q?=E4=B8=8A=E6=8A=A5=E5=A4=B1=E8=B4=A5=E8=B6=85=E8=BF=87=E6=9C=80?= =?UTF-8?q?=E5=A4=A7=E9=87=8D=E8=AF=95=E6=AC=A1=E6=95=B0=E7=9A=84=E6=B6=88?= =?UTF-8?q?=E6=81=AF=E8=AE=B0=E5=BD=95=E5=88=B0dead=20letter=20jsonl?= =?UTF-8?q?=E4=B8=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/config-list.md | 4 +- src/services/rawDump/batchWorker.ts | 69 +++++++++++++++++++---------- src/services/rawDump/index.ts | 44 +++++++++++++++++- src/services/rawDump/queue.ts | 2 + src/services/rawDump/state.ts | 34 ++++++++++---- src/services/rawDump/types.ts | 33 ++++++++++---- src/services/rawDump/worker.ts | 59 ++++++++++++++++++++++++ 7 files changed, 201 insertions(+), 44 deletions(-) diff --git a/docs/config-list.md b/docs/config-list.md index 8df65ad5a..c824a61ba 100644 --- a/docs/config-list.md +++ b/docs/config-list.md @@ -194,8 +194,8 @@ | `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` | 本地留存目录 | +| `CSC_RAW_DUMP_MODE` | boolean | `false` | **本地留存模式**:数据只写入本地文件,不上报服务端 | +| `CSC_RAW_DUMP_DIR` | string | `~/.claude/raw-dump` | 本地留存目录 | ### Bash/终端 diff --git a/src/services/rawDump/batchWorker.ts b/src/services/rawDump/batchWorker.ts index 9ff99b660..7b83d65e5 100644 --- a/src/services/rawDump/batchWorker.ts +++ b/src/services/rawDump/batchWorker.ts @@ -16,6 +16,7 @@ import { uploadConversation, uploadSummary, uploadCommits, + uploadStatistics, authWithFallback, } from './worker.js' import { @@ -29,8 +30,7 @@ import { MAX_ATTEMPTS, type QueueTask, } from './queue.js' -import { readState, writeState } from './state.js' -import type { RawDumpError } from './types.js' +import { readState, writeState, appendDeadLetter } from './state.js' import { getSessionDirectory, loadSessionMessages } from './worker.js' import { getRepoInfo } from './git.js' import { createLogger } from './logger.js' @@ -134,6 +134,25 @@ async function processTask( const authData = await authWithFallback() const repoInfo = await getCachedRepoInfo(task.directory) + if (task.messageID === '__statistics__' && task.statsData) { + await uploadStatistics( + { + sessionID: task.sessionID, + directory: task.directory, + sessionCount: task.statsData.sessionCount, + conversationCount: task.statsData.conversationCount, + upstreamTokens: task.statsData.upstreamTokens, + downstreamTokens: task.statsData.downstreamTokens, + startTime: task.statsData.startTime, + endTime: task.statsData.endTime, + }, + authData, + state, + ) + log.info('statistics task completed', { sessionID: task.sessionID }) + return + } + const conversationUploaded = await uploadConversation( { sessionID: task.sessionID, @@ -156,7 +175,10 @@ async function processTask( repoInfo, }) - log.info('task completed', { sessionID: task.sessionID, conversationUploaded }) + log.info('task completed', { + sessionID: task.sessionID, + conversationUploaded, + }) } /** @@ -218,13 +240,16 @@ async function runBatch() { }) if (task.attemptCount >= MAX_ATTEMPTS) { - // 彻底失败:记录错误,从队列移除 - state.errors[key] = { - message: errorMsg, - count: task.attemptCount, - lastAt: new Date().toISOString(), + // 彻底失败:写入 dead letter,从队列移除 + await appendDeadLetter({ + sessionID: task.sessionID, + messageID: task.messageID, + directory: task.directory, + attemptCount: task.attemptCount, + error: errorMsg, endpoint: extractEndpointFromError(errorMsg), - } satisfies RawDumpError + failedAt: new Date().toISOString(), + }) removeTask(key) log.warn('task permanently failed, removed from queue', { sessionID: task.sessionID, @@ -232,12 +257,6 @@ async function runBatch() { }) } else { // 失败但未超限:attemptCount++ 写回文件,下次重试 - state.errors[key] = { - message: errorMsg, - count: task.attemptCount, - lastAt: new Date().toISOString(), - endpoint: extractEndpointFromError(errorMsg), - } satisfies RawDumpError await flushQueue() } } @@ -279,15 +298,17 @@ export function startBatchWorker() { } // 启动时加载队列 + 随机抖动 - loadQueue().then(() => { - log.info('queue loaded', { count: getQueue().length }) - scheduleNext(Math.floor(Math.random() * 10_000)) - }).catch(err => { - log.error('failed to load queue', { - error: err instanceof Error ? err.message : String(err), + loadQueue() + .then(() => { + log.info('queue loaded', { count: getQueue().length }) + scheduleNext(Math.floor(Math.random() * 10_000)) + }) + .catch(err => { + log.error('failed to load queue', { + error: err instanceof Error ? err.message : String(err), + }) + process.exit(1) }) - process.exit(1) - }) } const scriptPath = process.argv[1] || '' @@ -296,4 +317,4 @@ if ( scriptPath.endsWith('batchWorker.js') ) { startBatchWorker() -} \ No newline at end of file +} diff --git a/src/services/rawDump/index.ts b/src/services/rawDump/index.ts index b1008e71c..0e1a9ec16 100644 --- a/src/services/rawDump/index.ts +++ b/src/services/rawDump/index.ts @@ -3,7 +3,10 @@ * 队列模式:主进程只 enqueue,单 batch worker 顺序消费 */ -import { getRawDumpMode, RAW_DUMP_MODE } from './localStorage.js' +import { + getRawDumpMode, + RAW_DUMP_MODE, +} from './localStorage.js' import { enqueue } from './queue.js' import { spawnBatchWorker } from './spawn.js' import { startBatchWorker } from './batchWorker.js' @@ -103,3 +106,42 @@ export function reportSession(sessionID: string, directory: string): void { enqueue({ sessionID, messageID: '__summary__', directory }) ensureBatchWorker() } + +export interface StatisticsData { + sessionCount: number + conversationCount: number + upstreamTokens: number + downstreamTokens: number + startTime: number + endTime: number +} + +const lastReportStatsMap = new Map() +const STATS_DEBOUNCE_MS = 60_000 // 同一 session 1 分钟内不重复 enqueue + +/** + * 上报对账统计数据(session数、conversation数、token数) + * 通过特殊 messageID '__statistics__' 标识入队 + */ +export function reportStatistics( + sessionID: string, + directory: string, + data: StatisticsData, +): void { + if (!isEnabled()) return + const key = `${sessionID}:__statistics__` + const now = Date.now() + const last = lastReportStatsMap.get(key) + if (last && now - last < STATS_DEBOUNCE_MS) { + log.debug('reportStatistics debounced', { sessionID, lastMs: now - last }) + return + } + lastReportStatsMap.set(key, now) + enqueue({ + sessionID, + messageID: '__statistics__', + directory, + statsData: data, + }) + ensureBatchWorker() +} diff --git a/src/services/rawDump/queue.ts b/src/services/rawDump/queue.ts index 1b49341b4..42108a768 100644 --- a/src/services/rawDump/queue.ts +++ b/src/services/rawDump/queue.ts @@ -12,6 +12,7 @@ import { promises as fs } from 'node:fs' import os from 'node:os' import path from 'path' +import type { StatisticsData } from './index.js' const QUEUE_FILE = path.join( os.homedir(), @@ -29,6 +30,7 @@ export interface QueueTask { directory: string enqueuedAt: number attemptCount: number + statsData?: StatisticsData } // ----------- 内存中的队列 ----------- diff --git a/src/services/rawDump/state.ts b/src/services/rawDump/state.ts index 2e2217da2..4721098e7 100644 --- a/src/services/rawDump/state.ts +++ b/src/services/rawDump/state.ts @@ -1,21 +1,22 @@ /** * Raw Dump 磁盘状态管理 - * 用于 conversation、summary、commits 的去重,以及错误聚合 + * 用于 conversation、summary、commits 的去重 * 通过文件锁保证多进程并发读写安全 * - * state.errors:按 sessionID:messageID 键控的错误聚合表 - * - 每批次累加 count,记录最新错误消息和时间戳 - * - 超过最大重试次数时由 batchWorker 写入 dead letter + * 错误追踪(dead letter): + * - 超过最大重试次数的任务追加写入 dead letter jsonl 文件 + * - 路径: ~/.claude/raw-dump/csc-dead-letter.jsonl */ 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' +import type { RawDumpState, DeadLetterEntry } from './types.js' const STATE_DIR = path.join(os.homedir(), '.claude', 'raw-dump') const STATE_FILE = path.join(STATE_DIR, 'csc-state.json') const STATE_LOCK_FILE = path.join(STATE_DIR, 'csc-state.lock') +const DEAD_LETTER_FILE = path.join(STATE_DIR, 'csc-dead-letter.jsonl') /** * 创建空的 RawDumpState 对象 @@ -26,8 +27,6 @@ function createEmptyState(): RawDumpState { conversation: {}, summary: {}, commits: {}, - errors: {}, - // summary 值为 RFC3339 格式时间戳字符串,解析时无需转换 } } @@ -106,7 +105,6 @@ export async function readState(): Promise { conversation: parsed.conversation ?? {}, summary: parsed.summary ?? {}, commits: parsed.commits ?? {}, - errors: parsed.errors ?? {}, } } catch { return createEmptyState() @@ -126,3 +124,23 @@ export async function writeState(state: RawDumpState): Promise { await fs.writeFile(STATE_FILE, JSON.stringify(state, null, 2), 'utf-8') }) } + +export async function readDeadLetter(): Promise { + try { + const text = await fs.readFile(DEAD_LETTER_FILE, 'utf-8') + return text + .split('\n') + .filter(Boolean) + .map(line => JSON.parse(line) as DeadLetterEntry) + } catch { + return [] + } +} + +export async function appendDeadLetter(entry: DeadLetterEntry): Promise { + await fs.mkdir(STATE_DIR, { recursive: true }) + await fs.writeFile(DEAD_LETTER_FILE, JSON.stringify(entry) + '\n', { + flag: 'a', + encoding: 'utf-8', + }) +} diff --git a/src/services/rawDump/types.ts b/src/services/rawDump/types.ts index c6ca33479..9e6774481 100644 --- a/src/services/rawDump/types.ts +++ b/src/services/rawDump/types.ts @@ -7,10 +7,9 @@ * - summary:会话统计(按 session 去重,5 分钟内只上报一次) * - commits:提交记录(按 commit ID 去重) * - * 错误追踪: - * - 上报失败的任务写入 failed queue,最多重试 MAX_RETRIES 次 - * - 超过最大重试次数的任务移入 dead letter 文件 - * - 所有错误聚合到 state.errors,按 sessionID:messageID 键控 + * 错误追踪(dead letter): + * - 超过最大重试次数的任务追加写入 dead letter jsonl 文件 + * - 路径: ~/.claude/raw-dump/csc-dead-letter.jsonl */ export const RAW_DUMP_EVENT_ENV_KEY = '__CSC_RAW_DUMP_EVENT__' @@ -25,14 +24,16 @@ export interface RawDumpState { conversation: Record summary: Record // RFC3339 时间戳,如 "2024-01-01T12:00:00.000Z" commits: Record - errors: Record } -export interface RawDumpError { - message: string - count: number - lastAt: string +export interface DeadLetterEntry { + sessionID: string + messageID: string + directory: string + attemptCount: number + error: string endpoint?: string + failedAt: string } export interface JwtPayload { @@ -107,3 +108,17 @@ export interface CommitPayload { subject: string parent_ids: string[] } + +export interface StatisticsPayload { + task_id: string + start_time: string + end_time: string + user_id: string + user_name: string + client_id: string + client_version: string + session_count: number + conversation_count: number + upstream_tokens: number + downstream_tokens: number +} diff --git a/src/services/rawDump/worker.ts b/src/services/rawDump/worker.ts index d14d40a9e..ac395c039 100644 --- a/src/services/rawDump/worker.ts +++ b/src/services/rawDump/worker.ts @@ -37,6 +37,7 @@ import type { CommitPayload, ConversationPayload, JwtPayload, + StatisticsPayload, SummaryPayload, } from './types.js' @@ -1072,6 +1073,64 @@ export async function uploadCommits( return commits.length } +const STATS_DEDUP_WINDOW_MS = 60 * 60 * 1000 // 同一 session 1 小时内只上报一次 statistics + +export async function uploadStatistics( + payload: { + sessionID: string + directory: string + sessionCount: number + conversationCount: number + upstreamTokens: number + downstreamTokens: number + startTime: number + endTime: number + }, + authData: Awaited>, + state: Awaited>, +): Promise { + const key = `stats:${payload.sessionID}` + const lastReported = state.summary[key] + if ( + lastReported && + Date.now() - new Date(lastReported).getTime() < STATS_DEDUP_WINDOW_MS + ) { + log.debug('statistics skipped: reported recently', { + task_id: payload.sessionID, + lastReported, + }) + return + } + + const body: StatisticsPayload = { + task_id: payload.sessionID, + start_time: formatIso(payload.startTime), + end_time: formatIso(payload.endTime), + ...authData.user, + client_id: authData.clientId, + client_version: authData.version, + session_count: payload.sessionCount, + conversation_count: payload.conversationCount, + upstream_tokens: payload.upstreamTokens, + downstream_tokens: payload.downstreamTokens, + } + + await postJson( + authData.baseUrl, + authData.headers, + '/raw-store/statistics', + body, + ) + state.summary[key] = new Date().toISOString() + log.info('statistics uploaded', { + task_id: payload.sessionID, + session_count: payload.sessionCount, + conversation_count: payload.conversationCount, + upstream_tokens: payload.upstreamTokens, + downstream_tokens: payload.downstreamTokens, + }) +} + /** * 从环境变量中解析 worker 入口传入的 payload * 抛出错误如果环境变量不存在或 JSON 解析失败