fix(rawDump): prevent batch worker concurrency cascade and add fetch timeout

- Add in-process isRunning flag to block reentrant runBatch (file lock
  fails when pid === process.pid)
- Replace setInterval + double immediate trigger with self-scheduling
  setTimeout that only fires after previous batch awaits
- Move clearQueue() to right after readQueue() so any unexpected
  concurrent runBatch sees an empty queue and exits immediately
- Cache repoInfo and workingTreeDiff in processTask, pass into all three
  upload functions to halve git invocations per task
- Add 30s AbortController timeout on postJson fetch so the worker no
  longer hangs indefinitely on unresponsive networks

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
林凯90331 2026-05-11 11:34:44 +08:00
parent 61b1b595bf
commit 1240a1a499
2 changed files with 119 additions and 51 deletions

View File

@ -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<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)
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<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)
}
}
}
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))
}
// 如果直接运行此文件

View File

@ -41,6 +41,10 @@ import type {
const log = createLogger('raw-dump')
const REQUEST_TIMEOUT_MS = 30_000 // 单次 HTTP 请求超时,防止 fetch 永久挂起
type RepoInfo = Awaited<ReturnType<typeof getRepoInfo>>
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<ReturnType<typeof auth>>,
state: Awaited<ReturnType<typeof readState>>,
options?: { workingTreeDiff?: string },
): Promise<boolean> {
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<string, unknown>[]
},
authData: Awaited<ReturnType<typeof auth>>,
options?: { repoInfo?: RepoInfo; workingTreeDiff?: string },
): 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 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<ReturnType<typeof auth>>,
state: Awaited<ReturnType<typeof readState>>,
options?: { repoInfo?: RepoInfo },
): Promise<number> {
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)