diff --git a/src/services/rawDump/batchWorker.ts b/src/services/rawDump/batchWorker.ts index f84c9a90e..070e8b8e8 100644 --- a/src/services/rawDump/batchWorker.ts +++ b/src/services/rawDump/batchWorker.ts @@ -1,18 +1,21 @@ /** * Raw Dump Batch Worker * 顺序消费队列,避免并发 429 - * 独立进程,通过 setInterval 定期执行 + * 独立进程,通过自循环 setTimeout 严格串行执行 */ 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 { getRepoInfo, getWorkingTreeDiff } from './git.js' import { createLogger } from './logger.js' const log = createLogger('raw-dump-batch') -const BATCH_INTERVAL_MS = 30_000 // 30 秒检查一次队列 +const BATCH_INTERVAL_MS = 30_000 // 每轮间隔 +// 进程内重入保护:文件锁不防同进程重入,必须用内存 flag 兜底 +let isRunning = false async function processTask(task: QueueTask) { log('info', 'processing task', { sessionID: task.sessionID, messageID: task.messageID }) @@ -27,19 +30,28 @@ async function processTask(task: QueueTask) { const authData = await auth() const state = await readState() + // 预加载 git 信息,三次上传共享,避免每个 task 重复 spawn 8+ 个 git 进程 + const repoInfo = await getRepoInfo(task.directory) + const workingTreeDiff = await getWorkingTreeDiff(task.directory) + try { // conversation const conversationUploaded = await uploadConversation( { sessionID: task.sessionID, messageID: task.messageID, directory: task.directory, messages }, authData, state, + { workingTreeDiff }, ) // summary(每个 turn 都报,但内容会累积) - await uploadSummary({ sessionID: task.sessionID, directory: task.directory, messages }, authData) + await uploadSummary( + { sessionID: task.sessionID, directory: task.directory, messages }, + authData, + { repoInfo, workingTreeDiff }, + ) // commits(限制频率,避免重复上报) - await uploadCommits({ directory: task.directory }, authData, state) + await uploadCommits({ directory: task.directory }, authData, state, { repoInfo }) log('info', 'task completed', { sessionID: task.sessionID, conversationUploaded }) } finally { @@ -49,62 +61,86 @@ async function processTask(task: QueueTask) { } async function runBatch() { - if (!acquireLock()) { - log('debug', 'another worker is running, skip') + // 第一道防线:同进程重入保护 + if (isRunning) { + log('debug', 'runBatch already running in-process, skip') return } + isRunning = true try { - const tasks = readQueue() - if (tasks.length === 0) { - log('debug', 'queue empty') + // 第二道防线:跨进程文件锁 + if (!acquireLock()) { + log('debug', 'another worker process holds the lock, skip') return } - log('info', `processing ${tasks.length} tasks`) - - // 去重:同一个 session 的多个 task,只保留最新的一个 - const deduped = new Map() - 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) + try { + const tasks = readQueue() + if (tasks.length === 0) { + log('debug', 'queue empty') + return } - } - const uniqueTasks = Array.from(deduped.values()).sort((a, b) => a.enqueuedAt - b.enqueuedAt) - log('info', `deduped to ${uniqueTasks.length} unique tasks`) + // 第三道防线:读完立刻清空队列 + // - 处理期间新进来的任务会在下一轮处理 + // - 即使有意外的并发 runBatch 拿到锁,也只会看到空队列直接返回 + clearQueue() - 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 }) + log('info', `processing ${tasks.length} tasks`) + + // 去重:同一个 session 的多个 task,只保留最新的一个 + const deduped = new Map() + 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) + } } - } - clearQueue() - log('info', 'batch completed') + 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, + }) + } + } + + log('info', 'batch completed') + } finally { + releaseLock() + } } finally { - releaseLock() + isRunning = false } } export function startBatchWorker() { log('info', 'batch worker started', { interval: BATCH_INTERVAL_MS }) - // 立即执行一次 - void runBatch() + // 自循环 setTimeout:上一轮跑完才安排下一轮,从源头消除并发 + // 即便 runBatch 抛错也确保下一轮被排上,避免 worker 卡死 + const scheduleNext = (delay: number) => { + setTimeout(async () => { + try { + await runBatch() + } catch (err) { + log('error', 'runBatch threw', { error: err instanceof Error ? err.message : String(err) }) + } + const jitter = Math.floor(Math.random() * 5_000) + scheduleNext(BATCH_INTERVAL_MS + jitter) + }, delay) + } - // 定期执行,添加随机抖动避免规律性 429 - const jitter = Math.floor(Math.random() * 10_000) - setTimeout(() => { - void runBatch() - setInterval(() => { - void runBatch() - }, BATCH_INTERVAL_MS) - }, jitter) + // 启动时随机抖动 0~10s,避免多个 csc 实例同时起 worker 撞 API + scheduleNext(Math.floor(Math.random() * 10_000)) } // 如果直接运行此文件 diff --git a/src/services/rawDump/worker.ts b/src/services/rawDump/worker.ts index be7c25f35..d98f824ad 100644 --- a/src/services/rawDump/worker.ts +++ b/src/services/rawDump/worker.ts @@ -41,6 +41,10 @@ import type { const log = createLogger('raw-dump') +const REQUEST_TIMEOUT_MS = 30_000 // 单次 HTTP 请求超时,防止 fetch 永久挂起 + +type RepoInfo = Awaited> + function formatIso(ms: number | undefined): string { if (!ms) return '' return new Date(ms).toISOString().replace(/\.\d{3}Z$/, 'Z') @@ -86,11 +90,14 @@ async function postJson( await new Promise((r) => setTimeout(r, delay)) } + const controller = new AbortController() + const timer = setTimeout(() => controller.abort(), REQUEST_TIMEOUT_MS) try { const res = await fetch(url, { method: 'POST', headers, body: JSON.stringify(body), + signal: controller.signal, }) if (res.ok) { @@ -108,8 +115,15 @@ async function postJson( 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 }) + const isAbort = lastError.name === 'AbortError' + // 网络错误 / 超时也重试 + log('warn', `${endpoint} ${isAbort ? 'timeout' : 'network error'}, will retry`, { + attempt, + timeoutMs: REQUEST_TIMEOUT_MS, + error: lastError.message, + }) + } finally { + clearTimeout(timer) } } @@ -343,6 +357,7 @@ export async function uploadConversation( }, authData: Awaited>, state: Awaited>, + options?: { workingTreeDiff?: string }, ): Promise { log('debug', 'uploadConversation start', { messageID: payload.messageID, messageCount: payload.messages.length }) @@ -374,12 +389,12 @@ export async function uploadConversation( const userMsgTime = (user?.timestamp as number) || Date.now() const assistantMsgTime = (assistant.timestamp as number) || Date.now() - // diff: 优先从 tool_use 提取,fallback 到 git diff HEAD + // 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 rawDiff = toolDiff.diff || options?.workingTreeDiff || (await getWorkingTreeDiff(payload.directory)) + log('debug', 'final diff', { diffLength: rawDiff.length, hasToolDiff: !!toolDiff.diff, fromCache: !toolDiff.diff && !!options?.workingTreeDiff }) const diffLines = rawDiff ? countDiffLines(rawDiff) : 0 const files = rawDiff ? extractFilesFromDiff(rawDiff) : [] @@ -425,11 +440,17 @@ export async function uploadSummary( messages: Record[] }, authData: Awaited>, + options?: { repoInfo?: RepoInfo; workingTreeDiff?: string }, ): Promise { 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 repoInfo = options?.repoInfo ?? (await getRepoInfo(payload.directory)) + const rawDiff = options?.workingTreeDiff ?? (await getWorkingTreeDiff(payload.directory)) + log('debug', 'summary repo info', { + repo_addr: repoInfo.repo_addr, + repo_branch: repoInfo.repo_branch, + diffLength: rawDiff.length, + fromCache: { repo: !!options?.repoInfo, diff: !!options?.workingTreeDiff }, + }) const assistants = payload.messages.filter((m) => m.type === 'assistant') const { upstream_tokens, downstream_tokens } = assistants.reduce( @@ -477,9 +498,10 @@ export async function uploadCommits( }, authData: Awaited>, state: Awaited>, + options?: { repoInfo?: RepoInfo }, ): Promise { log('debug', 'uploadCommits start', { directory: payload.directory }) - const repoInfo = await getRepoInfo(payload.directory) + const repoInfo = options?.repoInfo ?? (await getRepoInfo(payload.directory)) if (!repoInfo.repo_addr || !repoInfo.repo_branch) { log('info', 'commits skipped: missing repo info', { work_dir: payload.directory, repo_addr: repoInfo.repo_addr, repo_branch: repoInfo.repo_branch }) return 0 @@ -584,20 +606,30 @@ export async function runRawDumpWorker() { const state = await readState() log('debug', 'state loaded', { conversationCount: Object.keys(state.conversation).length, commitCount: Object.keys(state.commits).length }) + // 预加载 git 信息,三次上传共享,避免重复 spawn git + const repoInfo = await getRepoInfo(payload.directory) + const workingTreeDiff = await getWorkingTreeDiff(payload.directory) + log('debug', 'preloaded git info', { repo_branch: repoInfo.repo_branch, diffLength: workingTreeDiff.length }) + log('debug', 'starting uploadConversation...') const conversationUploaded = await uploadConversation( { ...payload, messages }, authData, state, + { workingTreeDiff }, ) log('debug', 'uploadConversation done', { conversationUploaded }) log('debug', 'starting uploadSummary...') - await uploadSummary({ sessionID: payload.sessionID, directory: payload.directory, messages }, authData) + await uploadSummary( + { sessionID: payload.sessionID, directory: payload.directory, messages }, + authData, + { repoInfo, workingTreeDiff }, + ) log('debug', 'uploadSummary done') log('debug', 'starting uploadCommits...') - const commitCount = await uploadCommits({ directory: payload.directory }, authData, state) + const commitCount = await uploadCommits({ directory: payload.directory }, authData, state, { repoInfo }) log('debug', 'uploadCommits done', { commitCount }) await writeState(state)