1. 修复匿名上报url错误; 2.增加注释; 3.增加异常场景处理提升获取对话diff的健壮性“

This commit is contained in:
zbc 2026-05-26 16:39:19 +08:00
parent 7ae2c9ba08
commit c8586ee6ee
9 changed files with 620 additions and 138 deletions

View File

@ -54,6 +54,22 @@ function isRecord(value: unknown): value is Record<string, unknown> {
return value !== null && typeof value === 'object' && !Array.isArray(value) return value !== null && typeof value === 'object' && !Array.isArray(value)
} }
function mapErrorToSDKError(
error: unknown,
): SDKAssistantMessageError | undefined {
if (!(error instanceof Error)) return undefined
// OpenAI SDK errors expose .status
const status = (error as Record<string, unknown>).status
if (typeof status === 'number') {
if (status === 401 || status === 403) return 'authentication_failed'
if (status === 429) return 'rate_limit'
if (status === 402) return 'billing_error'
if (status === 400) return 'invalid_request'
return 'server_error'
}
return 'unknown'
}
function isPlainRecord(value: unknown): value is Record<string, unknown> { function isPlainRecord(value: unknown): value is Record<string, unknown> {
return isRecord(value) return isRecord(value)
} }
@ -453,10 +469,7 @@ export async function* queryModelCoStrict(
? 'CoStrict API Error: The current model does not support image input. Switch to a multimodal or vision-capable model and try again.' ? 'CoStrict API Error: The current model does not support image input. Switch to a multimodal or vision-capable model and try again.'
: `CoStrict API Error: ${errorMsg}`, : `CoStrict API Error: ${errorMsg}`,
apiError: 'api_error', apiError: 'api_error',
error: error: mapErrorToSDKError(error),
error instanceof Error
? (error as unknown as SDKAssistantMessageError)
: undefined,
}) })
} }
} }

View File

@ -4,8 +4,19 @@
* setTimeout * setTimeout
*/ */
import { uploadConversation, uploadSummary, uploadCommits, authWithFallback } from './worker.js' import {
import { readQueue, clearQueue, acquireLock, releaseLock, type QueueTask } from './queue.js' uploadConversation,
uploadSummary,
uploadCommits,
authWithFallback,
} from './worker.js'
import {
readQueue,
clearQueue,
acquireLock,
releaseLock,
type QueueTask,
} from './queue.js'
import { readState, writeState } from './state.js' import { readState, writeState } from './state.js'
import { getSessionDirectory, loadSessionMessages } from './worker.js' import { getSessionDirectory, loadSessionMessages } from './worker.js'
import { getRepoInfo } from './git.js' import { getRepoInfo } from './git.js'
@ -21,6 +32,11 @@ const BATCH_INTERVAL_MS = 120_000 // 每轮间隔2 分钟),降低内联
const repoInfoCache = new Map<string, { repoInfo: RepoInfo; ts: number }>() const repoInfoCache = new Map<string, { repoInfo: RepoInfo; ts: number }>()
const REPO_CACHE_TTL_MS = 60_000 const REPO_CACHE_TTL_MS = 60_000
/**
* Git
* directory 60 git
* @returns RepoInfo
*/
async function getCachedRepoInfo(directory: string): Promise<RepoInfo> { async function getCachedRepoInfo(directory: string): Promise<RepoInfo> {
const cached = repoInfoCache.get(directory) const cached = repoInfoCache.get(directory)
if (cached && Date.now() - cached.ts < REPO_CACHE_TTL_MS) { if (cached && Date.now() - cached.ts < REPO_CACHE_TTL_MS) {
@ -37,6 +53,12 @@ let isRunning = false
const PARENT_PID = process.ppid const PARENT_PID = process.ppid
const IS_WORKER_PROCESS = process.argv[1]?.includes('batchWorker') || false const IS_WORKER_PROCESS = process.argv[1]?.includes('batchWorker') || false
/**
*
* - worker IS_WORKER_PROCESS = true
* - 使 kill(pid, 0)
* - 退 worker
*/
function isParentAlive(): boolean { function isParentAlive(): boolean {
if (!IS_WORKER_PROCESS) return true if (!IS_WORKER_PROCESS) return true
try { try {
@ -48,10 +70,22 @@ function isParentAlive(): boolean {
} }
// Session messages 缓存:同一 session 的多个 task 短时间内不需要重复读取 JSONL // Session messages 缓存:同一 session 的多个 task 短时间内不需要重复读取 JSONL
const sessionMessagesCache = new Map<string, { messages: Record<string, unknown>[]; ts: number }>() const sessionMessagesCache = new Map<
string,
{ messages: Record<string, unknown>[]; ts: number }
>()
const SESSION_CACHE_TTL_MS = 60_000 const SESSION_CACHE_TTL_MS = 60_000
async function getCachedSessionMessages(sessionDir: string, sessionID: string, messageID?: string) { /**
* session
* sessionDir + sessionID 60 JSONL
* 100ms info 便
*/
async function getCachedSessionMessages(
sessionDir: string,
sessionID: string,
messageID?: string,
) {
const cacheKey = `${sessionDir}:${sessionID}` const cacheKey = `${sessionDir}:${sessionID}`
const cached = sessionMessagesCache.get(cacheKey) const cached = sessionMessagesCache.get(cacheKey)
if (cached && Date.now() - cached.ts < SESSION_CACHE_TTL_MS) { if (cached && Date.now() - cached.ts < SESSION_CACHE_TTL_MS) {
@ -62,17 +96,39 @@ async function getCachedSessionMessages(sessionDir: string, sessionID: string, m
const messages = await loadSessionMessages(sessionDir, sessionID, messageID) const messages = await loadSessionMessages(sessionDir, sessionID, messageID)
const elapsed = Date.now() - start const elapsed = Date.now() - start
if (elapsed > 100) { if (elapsed > 100) {
log.info('loadSessionMessages slow', { sessionID, elapsedMs: elapsed, messageCount: messages.length }) log.info('loadSessionMessages slow', {
sessionID,
elapsedMs: elapsed,
messageCount: messages.length,
})
} }
sessionMessagesCache.set(cacheKey, { messages, ts: Date.now() }) sessionMessagesCache.set(cacheKey, { messages, ts: Date.now() })
return messages return messages
} }
async function processTask(task: QueueTask, state: Awaited<ReturnType<typeof readState>>) { /**
log.info('processing task', { sessionID: task.sessionID, messageID: task.messageID }) *
* - session
* - uploadConversationuploadSummaryuploadCommits
* - runBatch
* @param task
* @param state batch state
*/
async function processTask(
task: QueueTask,
state: Awaited<ReturnType<typeof readState>>,
) {
log.info('processing task', {
sessionID: task.sessionID,
messageID: task.messageID,
})
const sessionDir = getSessionDirectory(task.directory, task.sessionID) const sessionDir = getSessionDirectory(task.directory, task.sessionID)
const messages = await getCachedSessionMessages(sessionDir, task.sessionID, task.messageID) const messages = await getCachedSessionMessages(
sessionDir,
task.sessionID,
task.messageID,
)
if (messages.length === 0) { if (messages.length === 0) {
log.warn('no messages found', { sessionDir, sessionID: task.sessionID }) log.warn('no messages found', { sessionDir, sessionID: task.sessionID })
@ -86,7 +142,12 @@ async function processTask(task: QueueTask, state: Awaited<ReturnType<typeof rea
try { try {
// conversation // conversation
const conversationUploaded = await uploadConversation( const conversationUploaded = await uploadConversation(
{ sessionID: task.sessionID, messageID: task.messageID, directory: task.directory, messages }, {
sessionID: task.sessionID,
messageID: task.messageID,
directory: task.directory,
messages,
},
authData, authData,
state, state,
{ repoInfo }, { repoInfo },
@ -100,9 +161,14 @@ async function processTask(task: QueueTask, state: Awaited<ReturnType<typeof rea
) )
// commits限制频率避免重复上报 // commits限制频率避免重复上报
await uploadCommits({ directory: task.directory }, authData, state, { repoInfo }) await uploadCommits({ directory: task.directory }, authData, state, {
repoInfo,
})
log.info('task completed', { sessionID: task.sessionID, conversationUploaded }) log.info('task completed', {
sessionID: task.sessionID,
conversationUploaded,
})
} catch (err) { } catch (err) {
log.error('task failed', { log.error('task failed', {
error: err instanceof Error ? err.message : String(err), error: err instanceof Error ? err.message : String(err),
@ -112,6 +178,11 @@ async function processTask(task: QueueTask, state: Awaited<ReturnType<typeof rea
} }
} }
/**
* batch
* 线
* state
*/
async function runBatch() { async function runBatch() {
// 第一道防线:同进程重入保护 // 第一道防线:同进程重入保护
if (isRunning) { if (isRunning) {
@ -151,7 +222,9 @@ async function runBatch() {
} }
} }
const uniqueTasks = Array.from(deduped.values()).sort((a, b) => a.enqueuedAt - b.enqueuedAt) 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`)
// 一次性读取 state所有 task 共享,减少文件锁竞争和重复 JSON 解析 // 一次性读取 state所有 task 共享,减少文件锁竞争和重复 JSON 解析
@ -179,6 +252,13 @@ async function runBatch() {
} }
} }
/**
* Batch Worker
* setTimeout API
* 0~10 csc API
* batch 2 0~5
* 退
*/
export function startBatchWorker() { export function startBatchWorker() {
log.info('batch worker started', { interval: BATCH_INTERVAL_MS }) log.info('batch worker started', { interval: BATCH_INTERVAL_MS })
@ -193,7 +273,9 @@ export function startBatchWorker() {
try { try {
await runBatch() await runBatch()
} catch (err) { } 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) const jitter = Math.floor(Math.random() * 5_000)
scheduleNext(BATCH_INTERVAL_MS + jitter) scheduleNext(BATCH_INTERVAL_MS + jitter)
@ -206,6 +288,9 @@ export function startBatchWorker() {
// 如果直接运行此文件 // 如果直接运行此文件
const scriptPath = process.argv[1] || '' const scriptPath = process.argv[1] || ''
if (scriptPath.endsWith('batchWorker.ts') || scriptPath.endsWith('batchWorker.js')) { if (
scriptPath.endsWith('batchWorker.ts') ||
scriptPath.endsWith('batchWorker.js')
) {
startBatchWorker() startBatchWorker()
} }

View File

@ -8,6 +8,11 @@ import { promisify } from 'node:util'
const execFileAsync = promisify(execFile) const execFileAsync = promisify(execFile)
/**
* git
* - git 使cwdencodingmaxBuffer
* -
*/
async function gitExec(args: string[], cwd: string): Promise<string> { async function gitExec(args: string[], cwd: string): Promise<string> {
try { try {
const { stdout } = await execFileAsync('git', args, { const { stdout } = await execFileAsync('git', args, {
@ -22,6 +27,12 @@ async function gitExec(args: string[], cwd: string): Promise<string> {
} }
} }
/**
* Git
* - remote origin URL
* -
* - git
*/
export async function getRepoInfo(cwd: string) { export async function getRepoInfo(cwd: string) {
const [repoAddr, repoBranch, gitUserName, gitUserEmail] = await Promise.all([ const [repoAddr, repoBranch, gitUserName, gitUserEmail] = await Promise.all([
gitExec(['remote', 'get-url', 'origin'], cwd), gitExec(['remote', 'get-url', 'origin'], cwd),
@ -38,7 +49,11 @@ export async function getRepoInfo(cwd: string) {
} }
} }
export async function getRawDiff(cwd: string, from?: string, to?: string): Promise<string> { export async function getRawDiff(
cwd: string,
from?: string,
to?: string,
): Promise<string> {
if (from && to && from !== to) { if (from && to && from !== to) {
return gitExec(['diff', '--no-ext-diff', from, to], cwd) return gitExec(['diff', '--no-ext-diff', from, to], cwd)
} }
@ -46,10 +61,17 @@ export async function getRawDiff(cwd: string, from?: string, to?: string): Promi
return gitExec(['diff', 'HEAD'], cwd) return gitExec(['diff', 'HEAD'], cwd)
} }
/**
* HEAD diff
*/
export async function getWorkingTreeDiff(cwd: string): Promise<string> { export async function getWorkingTreeDiff(cwd: string): Promise<string> {
return gitExec(['diff', 'HEAD'], cwd) return gitExec(['diff', 'HEAD'], cwd)
} }
/**
* diff diff header
* + +++
*/
export function countDiffLines(diff: string): number { export function countDiffLines(diff: string): number {
let count = 0 let count = 0
for (const line of diff.split('\n')) { for (const line of diff.split('\n')) {
@ -59,6 +81,10 @@ export function countDiffLines(diff: string): number {
return count return count
} }
/**
* unified diff
* +++ b/--- a/diff --git a/ b/
*/
export function extractFilesFromDiff(diff: string): string[] { export function extractFilesFromDiff(diff: string): string[] {
const files = new Set<string>() const files = new Set<string>()
for (const line of diff.split('\n')) { for (const line of diff.split('\n')) {
@ -72,6 +98,10 @@ export function extractFilesFromDiff(diff: string): string[] {
return Array.from(files) return Array.from(files)
} }
/**
* git log format: %H|%aI|%an|%ae|%P|%s
* | commit
*/
export function parseCommitLog(output: string): Array<{ export function parseCommitLog(output: string): Array<{
commit_id: string commit_id: string
commit_time: string commit_time: string
@ -83,35 +113,86 @@ export function parseCommitLog(output: string): Array<{
if (!output.trim()) return [] if (!output.trim()) return []
return output return output
.split('\n') .split('\n')
.map((line) => { .map(line => {
const [commit_id, commit_time, git_user_name, git_user_email, parent_ids_str, ...rest] = line.split('|') const [
commit_id,
commit_time,
git_user_name,
git_user_email,
parent_ids_str,
...rest
] = line.split('|')
if (!commit_id || !git_user_email) return null if (!commit_id || !git_user_email) return null
const parent_ids = parent_ids_str ? parent_ids_str.trim().split(' ').filter(Boolean) : [] const parent_ids = parent_ids_str
return { commit_id, commit_time, git_user_name, git_user_email, parent_ids, subject: rest.join('|') } ? parent_ids_str.trim().split(' ').filter(Boolean)
: []
return {
commit_id,
commit_time,
git_user_name,
git_user_email,
parent_ids,
subject: rest.join('|'),
}
}) })
.filter((item): item is NonNullable<typeof item> => !!item) .filter((item): item is NonNullable<typeof item> => !!item)
} }
export async function getCommitLog(cwd: string, lastCommit?: string): Promise<string> { /**
* git commit
* - lastCommit commit commit
* - commit
* - author git
* - 50
*/
export async function getCommitLog(
cwd: string,
lastCommit?: string,
): Promise<string> {
const authorEmail = await gitExec(['config', 'user.email'], cwd) const authorEmail = await gitExec(['config', 'user.email'], cwd)
const authorFilter = authorEmail ? ['--author', authorEmail] : [] const authorFilter = authorEmail ? ['--author', authorEmail] : []
if (lastCommit) { if (lastCommit) {
return gitExec( return gitExec(
['log', `${lastCommit}..HEAD`, '--reverse', '--max-count=50', ...authorFilter, '--format=%H|%aI|%an|%ae|%P|%s'], [
'log',
`${lastCommit}..HEAD`,
'--reverse',
'--max-count=50',
...authorFilter,
'--format=%H|%aI|%an|%ae|%P|%s',
],
cwd, cwd,
) )
} }
return gitExec( return gitExec(
['log', '--since=1 day ago', '--reverse', '--max-count=50', ...authorFilter, '--format=%H|%aI|%an|%ae|%P|%s'], [
'log',
'--since=1 day ago',
'--reverse',
'--max-count=50',
...authorFilter,
'--format=%H|%aI|%an|%ae|%P|%s',
],
cwd, cwd,
) )
} }
export async function getCommitDiff(cwd: string, commitId: string): Promise<string> { /**
* commit diff
* 使 --diff-filter=ACDMR (Add)(Copy)(Delete)(Modify)(Rename)
*/
export async function getCommitDiff(
cwd: string,
commitId: string,
): Promise<string> {
return gitExec(['show', '--format=', '--diff-filter=ACDMR', commitId], cwd) return gitExec(['show', '--format=', '--diff-filter=ACDMR', commitId], cwd)
} }
/**
* commit subject 150
* comment
*/
export function toCommitComment(subject: string): string { export function toCommitComment(subject: string): string {
return Array.from(subject).slice(0, 150).join('') return Array.from(subject).slice(0, 150).join('')
} }

View File

@ -17,16 +17,35 @@ let batchWorkerSpawned = false
const lastEnqueueMap = new Map<string, number>() const lastEnqueueMap = new Map<string, number>()
const ENQUEUE_DEBOUNCE_MS = 5_000 const ENQUEUE_DEBOUNCE_MS = 5_000
/**
* Raw Dump
* - isLocalDumpMode
* - CSC_DISABLE_RAW_DUMP COSTRICT_DISABLE_RAW_DUMP '1'/'true'
* -
*/
function isEnabled(): boolean { function isEnabled(): boolean {
// 本地调试模式自动启用 // 本地调试模式自动启用
if (isLocalDumpMode()) return true if (isLocalDumpMode()) return true
// 显式禁用 // 显式禁用
if (process.env.CSC_DISABLE_RAW_DUMP === '1' || process.env.CSC_DISABLE_RAW_DUMP === 'true') return false if (
if (process.env.COSTRICT_DISABLE_RAW_DUMP === '1' || process.env.COSTRICT_DISABLE_RAW_DUMP === 'true') return false 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 // 默认启用 raw dump
return true return true
} }
/**
* Batch Worker
* - spawn worker
* - spawn worker runtime
*/
function ensureBatchWorker() { function ensureBatchWorker() {
if (batchWorkerSpawned) return if (batchWorkerSpawned) return
batchWorkerSpawned = true batchWorkerSpawned = true
@ -37,12 +56,21 @@ function ensureBatchWorker() {
} }
} }
/**
* session + message enqueue
* key 5 enqueue
* @returns true false
*/
function shouldEnqueue(sessionID: string, messageID: string): boolean { function shouldEnqueue(sessionID: string, messageID: string): boolean {
const key = `${sessionID}:${messageID}` const key = `${sessionID}:${messageID}`
const now = Date.now() const now = Date.now()
const last = lastEnqueueMap.get(key) const last = lastEnqueueMap.get(key)
if (last && now - last < ENQUEUE_DEBOUNCE_MS) { if (last && now - last < ENQUEUE_DEBOUNCE_MS) {
log.debug('reportTurn debounced', { sessionID, messageID, lastMs: now - last }) log.debug('reportTurn debounced', {
sessionID,
messageID,
lastMs: now - last,
})
return false return false
} }
lastEnqueueMap.set(key, now) lastEnqueueMap.set(key, now)
@ -53,13 +81,22 @@ function shouldEnqueue(sessionID: string, messageID: string): boolean {
* *
* batch worker * batch worker
*/ */
export function reportTurn(sessionID: string, messageID: string, directory: string): void { export function reportTurn(
sessionID: string,
messageID: string,
directory: string,
): void {
if (!isEnabled()) return if (!isEnabled()) return
if (!shouldEnqueue(sessionID, messageID)) return if (!shouldEnqueue(sessionID, messageID)) return
enqueue({ sessionID, messageID, directory }) enqueue({ sessionID, messageID, directory })
ensureBatchWorker() ensureBatchWorker()
} }
/**
* session
* 使 messageID '__summary__' conversation
* batch worker
*/
export function reportSession(sessionID: string, directory: string): void { export function reportSession(sessionID: string, directory: string): void {
if (!isEnabled()) return if (!isEnabled()) return
if (!shouldEnqueue(sessionID, '__summary__')) return if (!shouldEnqueue(sessionID, '__summary__')) return

View File

@ -9,15 +9,32 @@ import path from 'node:path'
const DEFAULT_LOCAL_DIR = path.join(os.homedir(), '.claude', 'raw-dump-local') const DEFAULT_LOCAL_DIR = path.join(os.homedir(), '.claude', 'raw-dump-local')
/**
*
* CSC_RAW_DUMP_LOCAL_DIR
*/
export function getLocalDumpDir(): string { export function getLocalDumpDir(): string {
return (process.env.CSC_RAW_DUMP_LOCAL_DIR || DEFAULT_LOCAL_DIR).replace(/\/$/, '') return (process.env.CSC_RAW_DUMP_LOCAL_DIR || DEFAULT_LOCAL_DIR).replace(
/\/$/,
'',
)
} }
/**
*
* CSC_RAW_DUMP_LOCAL_MODE '1' 'true'
*/
export function isLocalDumpMode(): boolean { export function isLocalDumpMode(): boolean {
const mode = process.env.CSC_RAW_DUMP_LOCAL_MODE const mode = process.env.CSC_RAW_DUMP_LOCAL_MODE
return mode === '1' || mode === 'true' return mode === '1' || mode === 'true'
} }
/**
* JSON
* - task_id
* - request_id/commit_id
* - body _dumpMeta API endpoint
*/
export async function writeLocalDump( export async function writeLocalDump(
type: 'conversation' | 'summary' | 'commit', type: 'conversation' | 'summary' | 'commit',
body: Record<string, unknown>, body: Record<string, unknown>,

View File

@ -8,7 +8,11 @@ import { promises as fs } from 'node:fs'
import os from 'node:os' import os from 'node:os'
import path from 'node:path' import path from 'node:path'
const QUEUE_FILE = path.join(os.homedir(), '.claude', 'csc-raw-dump-queue.jsonl') 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') const LOCK_FILE = path.join(os.homedir(), '.claude', 'csc-raw-dump.lock')
export interface QueueTask { export interface QueueTask {
@ -18,19 +22,31 @@ export interface QueueTask {
enqueuedAt: number enqueuedAt: number
} }
/**
* JSONL
* 使 fire-and-forget event loop
*
*/
export function enqueue(task: Omit<QueueTask, 'enqueuedAt'>): void { export function enqueue(task: Omit<QueueTask, 'enqueuedAt'>): void {
const item: QueueTask = { ...task, enqueuedAt: Date.now() } const item: QueueTask = { ...task, enqueuedAt: Date.now() }
// 使用 fire-and-forget 异步写入,避免阻塞主进程 event loop // 使用 fire-and-forget 异步写入,避免阻塞主进程 event loop
fs.appendFile(QUEUE_FILE, JSON.stringify(item) + '\n', 'utf-8').catch(() => {}) fs.appendFile(QUEUE_FILE, JSON.stringify(item) + '\n', 'utf-8').catch(
() => {},
)
} }
/**
* task
* - JSONL task
* -
*/
export async function readQueue(): Promise<QueueTask[]> { export async function readQueue(): Promise<QueueTask[]> {
try { try {
const text = await fs.readFile(QUEUE_FILE, 'utf-8') const text = await fs.readFile(QUEUE_FILE, 'utf-8')
return text return text
.split('\n') .split('\n')
.filter(Boolean) .filter(Boolean)
.map((line) => { .map(line => {
try { try {
return JSON.parse(line) as QueueTask return JSON.parse(line) as QueueTask
} catch { } catch {
@ -43,6 +59,10 @@ export async function readQueue(): Promise<QueueTask[]> {
} }
} }
/**
* truncate
* batch worker
*/
export async function clearQueue(): Promise<void> { export async function clearQueue(): Promise<void> {
try { try {
await fs.writeFile(QUEUE_FILE, '', 'utf-8') await fs.writeFile(QUEUE_FILE, '', 'utf-8')
@ -51,6 +71,11 @@ export async function clearQueue(): Promise<void> {
} }
} }
/**
* batch worker
* - PID false worker
* - 退
*/
export async function acquireLock(): Promise<boolean> { export async function acquireLock(): Promise<boolean> {
try { try {
try { try {
@ -75,6 +100,9 @@ export async function acquireLock(): Promise<boolean> {
} }
} }
/**
* batch worker
*/
export async function releaseLock(): Promise<void> { export async function releaseLock(): Promise<void> {
try { try {
await fs.writeFile(LOCK_FILE, '', 'utf-8') await fs.writeFile(LOCK_FILE, '', 'utf-8')

View File

@ -8,6 +8,11 @@ import { existsSync } from 'node:fs'
import path from 'node:path' import path from 'node:path'
import { fileURLToPath } from 'node:url' import { fileURLToPath } from 'node:url'
/**
* batch worker
* - Dev 使 src/services/rawDump/batchWorker.ts
* - Build dist/services/rawDump/batchWorker.js dist
*/
function resolveWorkerPath(): string { function resolveWorkerPath(): string {
const entry = process.execPath const entry = process.execPath
const isDev = path.basename(entry).toLowerCase().startsWith('bun') const isDev = path.basename(entry).toLowerCase().startsWith('bun')
@ -29,6 +34,12 @@ function resolveWorkerPath(): string {
return path.resolve(__dirname, 'batchWorker.js') return path.resolve(__dirname, 'batchWorker.js')
} }
/**
*
* bun node
* PATH bun node
* @returns { entry: 运行时路径, isBun: 是否为 bun } null
*/
function resolveRuntime(): { entry: string; isBun: boolean } | null { function resolveRuntime(): { entry: string; isBun: boolean } | null {
const execPath = process.execPath const execPath = process.execPath
const basename = path.basename(execPath).toLowerCase() const basename = path.basename(execPath).toLowerCase()
@ -64,9 +75,7 @@ export function spawnBatchWorker(): boolean {
return false return false
} }
const args = runtime.isBun const args = runtime.isBun ? ['run', workerPath] : [workerPath]
? ['run', workerPath]
: [workerPath]
try { try {
const child = spawn(runtime.entry, args, { const child = spawn(runtime.entry, args, {
@ -75,7 +84,7 @@ export function spawnBatchWorker(): boolean {
stdio: 'ignore', stdio: 'ignore',
}) })
child.on('error', (err) => { child.on('error', err => {
console.error('[raw-dump] batch worker spawn error:', err.message) console.error('[raw-dump] batch worker spawn error:', err.message)
}) })

View File

@ -13,6 +13,10 @@ const STATE_DIR = path.join(os.homedir(), '.claude')
const STATE_FILE = path.join(STATE_DIR, 'csc-raw-dump-state.json') 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') const STATE_LOCK_FILE = path.join(STATE_DIR, 'csc-raw-dump-state.lock')
/**
* RawDumpState
* state
*/
function createEmptyState(): RawDumpState { function createEmptyState(): RawDumpState {
return { return {
conversation: {}, conversation: {},
@ -21,6 +25,12 @@ function createEmptyState(): RawDumpState {
} }
} }
/**
* state
* - PID false
* - 退zombie
* - 使 kill(pid, 0)
*/
function acquireStateLock(): boolean { function acquireStateLock(): boolean {
try { try {
try { try {
@ -44,6 +54,9 @@ function acquireStateLock(): boolean {
} }
} }
/**
* state
*/
function releaseStateLock(): void { function releaseStateLock(): void {
try { try {
writeFileSync(STATE_LOCK_FILE, '', 'utf-8') writeFileSync(STATE_LOCK_FILE, '', 'utf-8')
@ -52,6 +65,11 @@ function releaseStateLock(): void {
} }
} }
/**
*
* - 5
* - finally
*/
async function withStateLock<T>(fn: () => Promise<T>): Promise<T> { async function withStateLock<T>(fn: () => Promise<T>): Promise<T> {
const start = Date.now() const start = Date.now()
while (!acquireStateLock()) { while (!acquireStateLock()) {
@ -59,7 +77,7 @@ async function withStateLock<T>(fn: () => Promise<T>): Promise<T> {
// 5 秒超时:降级为无锁执行,避免永久挂起 // 5 秒超时:降级为无锁执行,避免永久挂起
break break
} }
await new Promise((r) => setTimeout(r, 10)) await new Promise(r => setTimeout(r, 10))
} }
try { try {
return await fn() return await fn()
@ -68,6 +86,11 @@ async function withStateLock<T>(fn: () => Promise<T>): Promise<T> {
} }
} }
/**
* state
* -
* - state
*/
export async function readState(): Promise<RawDumpState> { export async function readState(): Promise<RawDumpState> {
return withStateLock(async () => { return withStateLock(async () => {
try { try {
@ -84,6 +107,12 @@ export async function readState(): Promise<RawDumpState> {
}) })
} }
/**
* state
* -
* - recursive: true
* - JSON 便
*/
export async function writeState(state: RawDumpState): Promise<void> { export async function writeState(state: RawDumpState): Promise<void> {
return withStateLock(async () => { return withStateLock(async () => {
await fs.mkdir(STATE_DIR, { recursive: true }) await fs.mkdir(STATE_DIR, { recursive: true })

View File

@ -46,11 +46,20 @@ const REQUEST_TIMEOUT_MS = 30_000 // 单次 HTTP 请求超时,防止 fetch 永
type RepoInfo = Awaited<ReturnType<typeof getRepoInfo>> type RepoInfo = Awaited<ReturnType<typeof getRepoInfo>>
/**
* ISO
* @param ms
* @returns ISO "2024-01-01T12:00:00Z"
*/
function formatIso(ms: number | undefined): string { function formatIso(ms: number | undefined): string {
if (!ms) return '' if (!ms) return ''
return new Date(ms).toISOString().replace(/\.\d{3}Z$/, 'Z') return new Date(ms).toISOString().replace(/\.\d{3}Z$/, 'Z')
} }
/**
* Raw Dump API Base URL
* COSTRICT_RAW_DUMP_BASE_URL CSC_RAW_DUMP_BASE_URL
*/
function resolveRawDumpBaseUrl(baseUrl?: string): string { function resolveRawDumpBaseUrl(baseUrl?: string): string {
const explicit = const explicit =
process.env.COSTRICT_RAW_DUMP_BASE_URL || process.env.CSC_RAW_DUMP_BASE_URL process.env.COSTRICT_RAW_DUMP_BASE_URL || process.env.CSC_RAW_DUMP_BASE_URL
@ -74,6 +83,12 @@ function resolveRawDumpBaseUrl(baseUrl?: string): string {
return raw.replace(/\/cloud-api$/, '') return raw.replace(/\/cloud-api$/, '')
} }
/**
* Raw Dump API URL
* @param baseUrl API
* @param endpoint API /raw-store/task-conversation
* @param isAnonymous 使 Authorization header
*/
function getRawDumpUrl( function getRawDumpUrl(
baseUrl: string, baseUrl: string,
endpoint: string, endpoint: string,
@ -81,11 +96,20 @@ function getRawDumpUrl(
): string { ): string {
const suffix = endpoint.startsWith('/') ? endpoint : `/${endpoint}` const suffix = endpoint.startsWith('/') ? endpoint : `/${endpoint}`
const prefix = isAnonymous const prefix = isAnonymous
? '/user-indicator/public' ? '/user-indicator/public/api/v1'
: '/user-indicator/api/v1' : '/user-indicator/api/v1'
return `${baseUrl}${prefix}${suffix}` return `${baseUrl}${prefix}${suffix}`
} }
/**
* JSON POST
* - JSON
* - /429 3 退
* @param baseUrl API
* @param headers HTTP
* @param endpoint API
* @param body JSON
*/
async function postJson( async function postJson(
baseUrl: string, baseUrl: string,
headers: Headers, headers: Headers,
@ -167,6 +191,10 @@ async function postJson(
throw lastError || new Error(`${endpoint} failed after retries`) throw lastError || new Error(`${endpoint} failed after retries`)
} }
/**
* JWT payload
* 使 refresh token fallback access token
*/
function parseUser( function parseUser(
accessPayload: JwtPayload, accessPayload: JwtPayload,
refreshPayload?: JwtPayload | null, refreshPayload?: JwtPayload | null,
@ -191,6 +219,10 @@ function parseUser(
} }
} }
/**
* OS
* darwin MacOS, win32 Windows, linux Linux,
*/
function detectOs(): string { function detectOs(): string {
const map: Record<string, string> = { const map: Record<string, string> = {
darwin: 'MacOS', darwin: 'MacOS',
@ -200,6 +232,12 @@ function detectOs(): string {
return map[process.platform] ?? process.platform return map[process.platform] ?? process.platform
} }
/**
* CoStrict token API headers
* - credentials
* - token
* - AuthorizationUser-Agent Headers
*/
export async function auth() { export async function auth() {
log.debug('auth start') log.debug('auth start')
let creds = await loadCoStrictCredentials() let creds = await loadCoStrictCredentials()
@ -290,8 +328,14 @@ export async function auth() {
} }
} }
// 从 JSONL 文件加载会话消息 /**
// csc 的会话文件名可能是 ses_{hash}.jsonl 或 {uuid}.jsonl * JSONL
* - sessionDir .jsonl
* - sessionId
* - sessionIdmessageId
* - csc ses_{hash}.jsonl {uuid}.jsonl
* @returns JSONL
*/
export async function loadSessionMessages( export async function loadSessionMessages(
sessionDir: string, sessionDir: string,
sessionId: string, sessionId: string,
@ -364,6 +408,10 @@ export async function loadSessionMessages(
return [] return []
} }
/**
* ID
* message.uuid message.id
*/
function findMessage( function findMessage(
messages: Record<string, unknown>[], messages: Record<string, unknown>[],
messageID: string, messageID: string,
@ -375,6 +423,11 @@ function findMessage(
) )
} }
/**
* assistant message user message
* csc user message assistant message
* assistant type === 'user'
*/
function findParentUserMessage( function findParentUserMessage(
messages: Record<string, unknown>[], messages: Record<string, unknown>[],
assistantMsg: Record<string, unknown>, assistantMsg: Record<string, unknown>,
@ -388,6 +441,10 @@ function findParentUserMessage(
return undefined return undefined
} }
/**
* agent user
* assistant modeagent isSidechain user isMeta
*/
function detectSender( function detectSender(
assistant: Record<string, unknown>, assistant: Record<string, unknown>,
user: Record<string, unknown> | undefined, user: Record<string, unknown> | undefined,
@ -408,6 +465,10 @@ function detectSender(
return 'user' return 'user'
} }
/**
*
* content content block type === 'text'
*/
function extractTextContent(msg: Record<string, unknown>): string { function extractTextContent(msg: Record<string, unknown>): string {
const content = (msg.message as Record<string, unknown>)?.content const content = (msg.message as Record<string, unknown>)?.content
if (!Array.isArray(content)) return String(content ?? '') if (!Array.isArray(content)) return String(content ?? '')
@ -417,6 +478,10 @@ function extractTextContent(msg: Record<string, unknown>): string {
.join('\n') .join('\n')
} }
/**
* structured patch unified diff
* git diff --no-ext-diff patch unified diff
*/
function structuredPatchToUnifiedDiff( function structuredPatchToUnifiedDiff(
filePath: string, filePath: string,
patches: Array<Record<string, unknown>>, patches: Array<Record<string, unknown>>,
@ -437,6 +502,10 @@ function structuredPatchToUnifiedDiff(
return header + '\n' + hunks.join('\n') + '\n' return header + '\n' + hunks.join('\n') + '\n'
} }
/**
* unified diff
* EditNotebookEdit diff
*/
function generateStringDiff( function generateStringDiff(
filePath: string, filePath: string,
oldStr: string, oldStr: string,
@ -453,6 +522,12 @@ function generateStringDiff(
return header + '\n' + hunk + '\n' + body + '\n' return header + '\n' + hunk + '\n' + body + '\n'
} }
/**
* assistant message unified diff
* user message toolUseResult gitDiff.patch unified diff
* fallback structuredPatch tool_use input diff
* @returns diff
*/
function extractToolDiff( function extractToolDiff(
msg: Record<string, unknown>, msg: Record<string, unknown>,
allMessages?: Record<string, unknown>[], allMessages?: Record<string, unknown>[],
@ -594,6 +669,10 @@ function extractToolDiff(
return { diff, diff_lines: countDiffLines(diff), files: Array.from(files) } return { diff, diff_lines: countDiffLines(diff), files: Array.from(files) }
} }
/**
* assistant message token 使
* input_tokensoutput_tokenscache token
*/
function extractUsage(msg: Record<string, unknown>) { function extractUsage(msg: Record<string, unknown>) {
const usage = (msg.message as Record<string, unknown>)?.usage as const usage = (msg.message as Record<string, unknown>)?.usage as
| Record<string, number> | Record<string, number>
@ -606,12 +685,49 @@ function extractUsage(msg: Record<string, unknown>) {
} }
} }
function extractError(msg: Record<string, unknown>) { const SDK_ERROR_CODE_MAP: Record<string, number> = {
const error = msg.error as Record<string, unknown> | undefined authentication_failed: 401,
if (!error) return {} billing_error: 402,
rate_limit: 429,
invalid_request: 400,
server_error: 500,
max_output_tokens: 413,
unknown: 500,
}
const name = String(error.name ?? 'UnknownError') /**
const message = typeof error.message === 'string' ? error.message : name * assistant message
* errorSDK Error CoStrict provider isApiErrorMessage
* HTTP
*/
function extractError(msg: Record<string, unknown>): {
error_code?: number
error_reason?: string
} {
const error = msg.error
// Case 1: msg.error is a string (SDKAssistantMessageError — normal path)
if (typeof error === 'string') {
const errorCode = SDK_ERROR_CODE_MAP[error] ?? 500
// Prefer errorDetails (diagnostic info) over the raw string
const errorDetails = msg.errorDetails
const reason =
typeof errorDetails === 'string' && errorDetails ? errorDetails : error
return { error_code: errorCode, error_reason: reason }
}
// Case 2: msg.error is an Error object (CoStrict provider —不规范存储)
if (typeof error === 'object' && error !== null) {
const err = error as Record<string, unknown>
// Try apiError field first for coarse classification
const apiError = msg.apiError
if (typeof apiError === 'string' && apiError === 'max_output_tokens') {
const reason =
typeof err.message === 'string' ? err.message : 'max_output_tokens'
return { error_code: 413, error_reason: reason }
}
const name = String(err.name ?? 'UnknownError')
const errMessage = typeof err.message === 'string' ? err.message : name
const errorCode = const errorCode =
name === 'ProviderAuthError' name === 'ProviderAuthError'
? 401 ? 401
@ -619,13 +735,28 @@ function extractError(msg: Record<string, unknown>) {
? 413 ? 413
: name === 'MessageAbortedError' : name === 'MessageAbortedError'
? 499 ? 499
: name === 'APIError' && typeof error.statusCode === 'number' : name === 'APIError' && typeof err.statusCode === 'number'
? error.statusCode ? err.statusCode
: 500 : 500
return { error_code: errorCode, error_reason: errMessage }
return { error_code: errorCode, error_reason: message }
} }
// Case 3: no error field, but marked as API error (e.g. image size errors)
if (msg.isApiErrorMessage) {
return { error_code: 500, error_reason: 'api_error_unclassified' }
}
return {}
}
/**
* /raw-store/task-conversation
* - messages assistant message ID assistant
* - conversation payload POST
* - assistant diff
* - state
* @returns false
*/
export async function uploadConversation( export async function uploadConversation(
payload: { payload: {
sessionID: string sessionID: string
@ -783,6 +914,12 @@ export async function uploadConversation(
const SUMMARY_DEDUP_WINDOW_MS = 5 * 60 * 1000 // 同一 session 5 分钟内 summary 只上报一次 const SUMMARY_DEDUP_WINDOW_MS = 5 * 60 * 1000 // 同一 session 5 分钟内 summary 只上报一次
/**
* session /raw-store/task-summary
* session 5 session state.summary
* session
* 使 conversation summary
*/
export async function uploadSummary( export async function uploadSummary(
payload: { payload: {
sessionID: string sessionID: string
@ -833,6 +970,14 @@ export async function uploadSummary(
log.info('summary uploaded', { task_id: payload.sessionID }) log.info('summary uploaded', { task_id: payload.sessionID })
} }
/**
* /raw-store/commit
* - commit_id
* - 50 commit
* - commit POST state.commits
* - 10 commit 500ms
* @returns commit
*/
export async function uploadCommits( export async function uploadCommits(
payload: { payload: {
directory: string directory: string
@ -913,21 +1058,44 @@ export async function uploadCommits(
return commits.length return commits.length
} }
/**
* worker payload
* JSON
*/
function parseWorkerPayload(): RawDumpEventPayload { function parseWorkerPayload(): RawDumpEventPayload {
const raw = process.env[RAW_DUMP_EVENT_ENV_KEY] const raw = process.env[RAW_DUMP_EVENT_ENV_KEY]
if (!raw) throw new Error('missing raw dump payload') if (!raw) throw new Error('missing raw dump payload')
return JSON.parse(raw) as RawDumpEventPayload return JSON.parse(raw) as RawDumpEventPayload
} }
/**
* CoStrict ~.claude CLAUDE_CONFIG_HOME
*/
export function getClaudeConfigHomeDir(): string { export function getClaudeConfigHomeDir(): string {
return process.env.CLAUDE_CONFIG_HOME || path.join(os.homedir(), '.claude') return process.env.CLAUDE_CONFIG_HOME || path.join(os.homedir(), '.claude')
} }
/**
*
* / -
* /Users/linkai/code/csc -Users-linkai-code-csc
*/
function normalizeProjectPath(dir: string): string { function normalizeProjectPath(dir: string): string {
// 将 /Users/linkai/code/csc 转换为 -Users-linkai-code-csc // 将 /Users/linkai/code/csc 转换为 -Users-linkai-code-csc
return dir.replace(/\//g, '-') return dir.replace(/\//g, '-')
} }
/**
* session
*
* 1. ~/.claude/projects/{normalized-path}/
* 2. ~/.claude/transcripts/
* 3. ~/.claude/sessions/
* 4. {directory}/.claude/sessions/
* 5. {directory}/.claude/
* 6. {directory}
* 7. CSC_SESSION_DIR
*/
export function getSessionDirectory( export function getSessionDirectory(
directory: string, directory: string,
sessionID: string, sessionID: string,
@ -947,6 +1115,15 @@ export function getSessionDirectory(
return candidates.find(d => d) || directory return candidates.find(d => d) || directory
} }
/**
* Raw Dump Worker
*
* 1. payload
* 2. session
* 3. / conversationsummarycommits
* 4. state
*
*/
export async function runRawDumpWorker() { export async function runRawDumpWorker() {
try { try {
const payload = parseWorkerPayload() const payload = parseWorkerPayload()
@ -1031,6 +1208,12 @@ export async function runRawDumpWorker() {
} }
} }
/**
*
* - isLocalDumpMode使
* - Authorization header
* 使
*/
export async function authWithFallback(): Promise< export async function authWithFallback(): Promise<
Awaited<ReturnType<typeof auth>> Awaited<ReturnType<typeof auth>>
> { > {