// apps/api/src/queue/runner.ts // Serial queue runner. One worker loop per process; jobs and their items run // strictly in order. Each item flows through: // 1. analyzer (api or mcp) → produces TaskPackage // 2. judge → 6-dimension verdict + pass policy check // 3. zip-export → produces *.task-package.zip on disk // // The actual analyzer / judge / export logic is delegated to apps/api/src/analyzers/* // and apps/api/src/judge/*. This file just orchestrates state transitions. import { join } from "node:path"; import { mkdirSync, readFileSync, writeFileSync } from "node:fs"; import { buildTaskPackage } from "@task-optimizer/core/rl-env"; import { parse as parseApiZip, buildEvidence as buildApiEvidence } from "@task-optimizer/core/importers/api"; import { parse as parseMcpZip, buildEvidence as buildMcpEvidence } from "@task-optimizer/core/importers/mcp"; import { collectMcpConfigUrls, collectMcpToolsFromTask, type PluginMcpTool, } from "@task-optimizer/core/mcp-plugin-match"; import { runLocalRules } from "@task-optimizer/core/rules/mcp"; import { buildAnalysisZipBuffer } from "@task-optimizer/core/zip-export"; import * as apiPrompt from "@task-optimizer/core/prompts/api"; import * as mcpPrompt from "@task-optimizer/core/prompts/mcp"; import type { Config } from "../config.js"; import type { Logger } from "../log.js"; import { AiClient, safePreview } from "../ai/client.js"; import { runJudge, buildGoldenTrajectory } from "../judge/runner.js"; import { evaluateJudgePolicy } from "../judge/policy.js"; import { buildIterationMemoryPrompt, createIterationMemory, optimizeIterationMemory, type IterationMemory, } from "./iteration-memory.js"; import type { DetectedMode, ItemRow, ItemStatus, Store, } from "./store.js"; export interface QueueRunnerOptions { store: Store; config: Config; logger: Logger; } export interface OptionsSnapshot { /** Optional UI-requested grouping mode. Zip format detection remains separate. */ requestedMode?: DetectedMode; model?: string; baseURL?: string; /** * Optional override for the main analyzer's API key (used to generate * recommended_groups + recommended_rubrics). When unset, the backend falls * back to `DEFAULT_API_KEY` from .env. */ apiKey?: string; apiMode?: "chat" | "responses"; /** Parallel item workers per job. Clamped to [1, MAX_CONCURRENCY=8]. */ concurrency?: number; temperature?: number; topP?: number; maxTokens?: number; streaming?: boolean; enableThinking?: boolean; clearThinking?: boolean; ruleProfile?: string; optimizeTask?: boolean; feedback?: string; /** When true, pass policy requires total === max_total AND has_zeros === false. */ strictFullMarks?: boolean; /** When strictFullMarks is true, allow task_complexity to be 1/2 while all other dimensions are full marks. */ ignoreTaskComplexityForFullMarks?: boolean; judge?: { model?: string; baseURL?: string; apiMode?: "chat" | "responses"; promptTemplate?: string; }; rubricGeneration?: { model?: string; baseURL?: string; apiMode?: "chat" | "responses"; }; // NOTE: judge.apiKey / rubricGeneration.apiKey remain backend-only and are // sourced from .env. } interface InFlightJob { jobId: string; cancelRequested: boolean; } type AnalyzerResult = { mode: DetectedMode; parsed: unknown; evidence: Record; flags: Array<{ code: string; note: string; evidence: unknown }>; result: unknown; }; type TaskPackageResult = ReturnType; type JudgeRunResult = Awaited>; /** * Maximum analyzer + judge iterations per item. Each iteration feeds the * previous judge verdict back into the analyzer prompt. */ const MAX_JUDGE_ITERATIONS = 5; /** * Hard cap on parallel item workers per job. Protects upstream model APIs * from rate-limit storms when a user uploads many ZIPs at once. The frontend * setting is clamped to [1, MAX_CONCURRENCY]. */ const MAX_CONCURRENCY = 8; class Semaphore { private active = 0; private readonly queue: Array<() => void> = []; constructor(private readonly limit: number) {} private async acquire(): Promise { if (this.limit <= 0) return; if (this.active < this.limit) { this.active += 1; return; } await new Promise((resolve) => this.queue.push(resolve)); } private release(): void { if (this.limit <= 0) return; const next = this.queue.shift(); if (next) { next(); return; } this.active -= 1; } async run(fn: () => Promise): Promise { if (this.limit <= 0) return fn(); await this.acquire(); try { return await fn(); } finally { this.release(); } } } /** Build the FEEDBACK block injected into the next analyzer prompt. */ function buildIterationFeedback( judgeResult: { total_score?: number; max_score?: number; has_zeros?: boolean; verdict?: string; rationale?: string; dimensions?: Record; }, iter: number ): string { const score = Number(judgeResult.total_score || 0); const maxScore = Number(judgeResult.max_score) || 12; const dims = judgeResult.dimensions || {}; const weakDims = Object.entries(dims) .filter(([, v]) => Number(v?.score ?? 99) <= 1) .map(([key, v]) => { const reason = String(v?.explanation || "").slice(0, 320); return `- ${key} (score=${v?.score ?? "?"}/2): ${reason}`; }); const rationale = String(judgeResult.rationale || "").slice(0, 480); return [ `=== JUDGE FEEDBACK (iteration ${iter}, score ${score}/${maxScore}, has_zeros=${judgeResult.has_zeros ? "true" : "false"}) ===`, "", "TOP PRIORITY — EVIDENCE-FLOOR RULE:", "Every action verb in recommended_instruction MUST map to a real call in EVIDENCE_SUMMARY.", "If the judge cited an unsupported action or missing rubric, DELETE that verb from recommended_instruction", "and shrink the instruction so it describes only what the trajectory actually does.", "", "MANDATORY NEXT ITERATION BEHAVIOR:", "- Explicitly fix every issue listed below before changing unrelated fields.", "- Do NOT repeat any rubric name, checker_key, or task wording the judge criticized.", "- Prefer the smallest targeted edit that removes the cited judge complaint.", "- Task Complexity is allowed to remain below 2; do not invent extra work to inflate it.", "", rationale ? `JUDGE RATIONALE: ${rationale}` : "", weakDims.length ? "WEAK / ZERO DIMENSIONS:\n" + weakDims.join("\n") : "", "", "=== END JUDGE FEEDBACK ===", ] .filter(Boolean) .join("\n"); } export class QueueRunner { private readonly store: Store; private readonly config: Config; private readonly logger: Logger; private readonly aiClient: AiClient; private readonly llmSemaphore: Semaphore; private readonly inFlight = new Map(); private shuttingDown = false; constructor(opts: QueueRunnerOptions) { this.store = opts.store; this.config = opts.config; this.logger = opts.logger; this.aiClient = new AiClient({ upstreamProxy: opts.config.upstreamProxy }); this.llmSemaphore = new Semaphore(Number(opts.config.globalLlmConcurrency || 0)); } /** Schedule a job to run. Idempotent — if already running, no-op. */ scheduleJob(jobId: string): void { if (this.shuttingDown) return; if (this.inFlight.has(jobId)) return; this.inFlight.set(jobId, { jobId, cancelRequested: false }); void this.store.setJobStatus(jobId, "running"); // Fire-and-forget: each job runs in parallel with other jobs. Per-job // concurrency is enforced inside runJob via N item workers. Errors are // logged inside runJob; this top-level catch only guards against the // unhandled-promise-rejection edge case. void this.runJob(jobId).catch((err) => { this.logger.error({ err, jobId }, "Unhandled error in runJob"); }); } cancelJob(jobId: string): void { const handle = this.inFlight.get(jobId); if (handle) handle.cancelRequested = true; // If the job is still queued (not started), mark its remaining items as // stopped immediately. The runner will skip them. } async retryItem(itemId: string): Promise { const item = await this.store.getItem(itemId); if (!item) return; await this.store.resetItemForRetry(itemId); this.scheduleJob(item.job_id); } requestShutdown(): void { this.shuttingDown = true; for (const handle of this.inFlight.values()) handle.cancelRequested = true; } /** Run one job to completion (or cancellation). */ private async runJob(jobId: string): Promise { const log = this.logger.child({ jobId }); const handle = this.inFlight.get(jobId); if (!handle) { log.warn("Job worker invoked but no in-flight handle present"); return; } // Concurrency is read from the job's options snapshot (frontend-supplied). // Falls back to 1 (serial) when unset; capped to MAX_CONCURRENCY. const job = await this.store.getJob(jobId); const snapshot = parseOptionsSnapshot(job?.options_snapshot); const requested = Number(snapshot.concurrency) || 1; const concurrency = Math.max( 1, Math.min(MAX_CONCURRENCY, Math.floor(requested)) ); log.info({ concurrency }, "Job worker starting"); try { // Spawn N parallel workers. Each loops, atomically taking the next // queued item via SQLite transaction (takeNextQueuedItem) until none // remain or cancellation is requested. const workers = Array.from({ length: concurrency }, (_, i) => this.itemWorker(jobId, handle, i + 1) ); await Promise.all(workers); if (handle.cancelRequested || this.shuttingDown) { // Mark all remaining queued items as stopped. for (const it of await this.store.listItemsByJob(jobId)) { if (it.status === "queued") await this.store.setItemStatus(it.id, "stopped"); } await this.store.setJobStatus(jobId, "cancelled"); log.info("Job cancelled"); } else { await this.store.setJobStatus(jobId, "completed"); log.info("Job completed"); } } catch (err) { log.error({ err }, "Job worker crashed"); await this.store.setJobStatus(jobId, "completed"); // surface via item statuses } finally { this.inFlight.delete(jobId); } } /** * Single-worker drain loop. Multiple of these run in parallel within one * job. Each call to takeNextQueuedItem is atomic (SQLite transaction), so * workers never observe the same item. */ private async itemWorker( jobId: string, handle: InFlightJob, workerIndex: number ): Promise { const log = this.logger.child({ jobId, worker: workerIndex }); while (!handle.cancelRequested && !this.shuttingDown) { const next = await this.store.takeNextQueuedItem(jobId); if (!next) return; // queue drained try { await this.runItem(next, jobId); } catch (err) { log.error({ err, itemId: next.id }, "Item runner threw unexpectedly"); } } } /** Run a single item end-to-end. */ private async runItem(item: ItemRow, jobId: string): Promise { const log = this.logger.child({ jobId, itemId: item.id }); log.info({ filename: item.filename, detectedMode: item.detected_mode }, "Item start"); const itemDir = join(this.config.exportDir, jobId, item.id); try { mkdirSync(itemDir, { recursive: true }); } catch (err) { log.error({ err, itemDir }, "Failed to create item export directory"); await this.failItem( item.id, err instanceof Error ? err.message : String(err), "create_export_directory", ); return; } let stage = "initializing"; let uploadBytes: Buffer | null = null; let latestAnalysis: AnalyzerResult | null = null; let latestTaskPackage: TaskPackageResult | null = null; let latestJudge: JudgeRunResult | null = null; let bestAnalysis: AnalyzerResult | null = null; let bestTaskPackage: TaskPackageResult | null = null; let bestJudge: JudgeRunResult | null = null; let bestScore = -1; let lastJudge: JudgeRunResult | null = null; let iter = 0; let completedIterations = 0; let iterationMemory = createIterationMemory(); try { stage = "loading_job_options"; const job = await this.store.getJob(jobId); const snapshot = parseOptionsSnapshot(job?.options_snapshot); stage = "reading_upload"; uploadBytes = readFileSync(item.upload_path); // ── Iterate analyzer + judge up to MAX_JUDGE_ITERATIONS times ────── // Each round feeds the previous judge verdict back into the analyzer // prompt so the model can refine the package. const handle = this.inFlight.get(jobId); let extraFeedback = ""; let passed = false; for (iter = 1; iter <= MAX_JUDGE_ITERATIONS; iter++) { if (handle?.cancelRequested || this.shuttingDown) break; log.info({ iter, max: MAX_JUDGE_ITERATIONS }, "Iteration start"); // ── Stage 1: analyzer ───────────────────────────────────────────── stage = `iteration_${iter}:analyzer`; await this.transitionItem(item.id, "running", undefined, stage); const analysis = await this.runAnalyzer( item, uploadBytes, snapshot, extraFeedback ); stage = `iteration_${iter}:build_task_package`; const taskPackage = buildTaskPackage( analysis.result, analysis.evidence, analysis.parsed, buildGoldenTrajectory ); latestAnalysis = analysis; latestTaskPackage = taskPackage; if (!bestAnalysis || !bestTaskPackage) { bestAnalysis = analysis; bestTaskPackage = taskPackage; } // Persist the analyzer output before calling Judge. If Judge is // unavailable (for example HTTP 403), the generated groups/rubrics still // remain inspectable and downloadable. writeJson(join(itemDir, `evidence-iter-${iter}.json`), { mode: analysis.mode, detectedMode: item.detected_mode, flags: analysis.flags, evidence: analysis.evidence, analyzerResult: analysis.result, }); writeJson(join(itemDir, `task-package-iter-${iter}.json`), taskPackage); // ── Stage 2: judge ──────────────────────────────────────────────── stage = `iteration_${iter}:judge`; await this.transitionItem(item.id, "judging", undefined, stage); const judgeResult = await this.withLlmSlot(() => runJudge({ aiClient: this.aiClient, config: this.config, snapshot, taskPackage, evidence: analysis.evidence, }), ); latestJudge = judgeResult; lastJudge = judgeResult; const score = Number(judgeResult.total_score || 0); if (score > bestScore) { bestScore = score; bestAnalysis = analysis; bestTaskPackage = taskPackage; bestJudge = judgeResult; } writeJson(join(itemDir, `judge-result-iter-${iter}.json`), judgeResult); const iterationFeedback = buildIterationFeedback(judgeResult, iter); iterationMemory = optimizeIterationMemory( iterationMemory, judgeResult, iter, iterationFeedback, ); writeJson(join(itemDir, `iteration-memory-iter-${iter}.json`), iterationMemory); completedIterations = iter; if (isJudgeTargetSatisfied(judgeResult, snapshot)) { log.info( { iter, score, verdict: judgeResult.verdict }, "Judge target satisfied" ); passed = true; break; } log.info( { iter, score, max: MAX_JUDGE_ITERATIONS }, "Judge target not yet satisfied; preparing next iteration" ); extraFeedback = buildIterationMemoryPrompt(iterationMemory); } // Persist canonical artefacts (best round) for downstream consumers. stage = "persisting_final_artifacts"; const finalAnalysis = bestAnalysis ?? latestAnalysis; const finalTaskPackage = bestTaskPackage ?? latestTaskPackage; const finalJudge = bestJudge ?? lastJudge; if (!finalAnalysis || !finalTaskPackage) { throw new Error("迭代未产生任何可用结果(可能在第一次迭代前被取消或失败)。"); } // ── Stage 3: ZIP export ───────────────────────────────────────────── // Export the best task package even when Judge does not pass, so users // can inspect and reuse the generated groups/rubrics. stage = "zip_export"; const exportZipPath = await this.persistExportableArtifacts({ item, itemDir, uploadBytes, analysis: finalAnalysis, taskPackage: finalTaskPackage, judgeResult: finalJudge, iterations: completedIterations || Math.min(iter, MAX_JUDGE_ITERATIONS), passed, iterationMemory, }); if (!finalJudge) { throw new Error("Judge 未产生可用结果,但分组和 rubric 已生成并导出。"); } if (!passed) { const message = buildJudgeFailureMessage(finalJudge, finalTaskPackage, snapshot) + `\n(已迭代 ${completedIterations}/${MAX_JUDGE_ITERATIONS} 轮,最佳 ${bestScore} 分)`; throw new Error(message); } await this.transitionItem(item.id, "completed", undefined, stage); log.info( { exportZipPath, verdict: finalJudge.verdict, score: finalJudge.total_score, iterations: iter, }, "Item completed" ); } catch (err) { if (uploadBytes && latestAnalysis && latestTaskPackage) { try { await this.persistExportableArtifacts({ item, itemDir, uploadBytes, analysis: bestAnalysis ?? latestAnalysis, taskPackage: bestTaskPackage ?? latestTaskPackage, judgeResult: bestJudge ?? latestJudge ?? lastJudge, iterations: completedIterations || Math.min(Math.max(iter, 1), MAX_JUDGE_ITERATIONS), passed: false, iterationMemory, }); } catch (exportErr) { log.error({ err: exportErr }, "Failed to preserve analyzer artifacts after item error"); } } const details = serializeError(err, { stage, jobId, itemId: item.id, ordinal: item.ord, filename: item.filename, }); const message = errorPreview(err, details); writeFileSync(join(itemDir, "error.txt"), message + "\n"); const errorDetailsPath = join(itemDir, "error-details.json"); writeJson(errorDetailsPath, details); if ( isTransientItemError(err) && Number(item.attempt_count || 0) < Number(this.config.itemMaxAttempts || 3) ) { await this.transitionItem( item.id, "queued", `Transient failure, retrying (${item.attempt_count}/${this.config.itemMaxAttempts}): ${message}`, stage, errorDetailsPath, ); log.warn( { preview: message, attempt: item.attempt_count, maxAttempts: this.config.itemMaxAttempts }, "Item failed transiently; requeued for retry", ); return; } await this.failItem(item.id, message, stage, errorDetailsPath); log.error({ err, preview: message }, "Item failed"); } } private async runAnalyzer( item: ItemRow, uploadBytes: Buffer, snapshot: OptionsSnapshot, extraFeedback = "" ): Promise { const mode = resolveRequestedMode(snapshot, item.detected_mode); const isMcp = mode === "mcp"; const parsed = isMcp ? await parseMcpZip(uploadBytes, item.filename) : await parseApiZip(uploadBytes, item.filename); let evidence: Record; if (isMcp) { const toolResolution = await resolvePluginMcpTools(parsed, this.logger.child({ itemId: item.id })); evidence = buildMcpEvidence(parsed as never, { strictPluginMcp: true, pluginMcpTools: toolResolution.tools, pluginMcpToolSource: toolResolution.source, }) as unknown as Record; assertStrictMcpEvidenceReady(evidence, item.filename); } else { evidence = { type: "http", ...buildApiEvidence(parsed as never) } as Record; } const flags = runLocalRules(evidence); const promptModule = isMcp ? mcpPrompt : apiPrompt; const slimEvidence = promptModule.buildSlimEvidence(evidence); const system = promptModule.getSystemPrompt(); const baseFeedback = snapshot.feedback || ""; const mergedFeedback = [baseFeedback, extraFeedback].filter(Boolean).join("\n\n"); const user = promptModule.buildUserPrompt( slimEvidence, flags, mergedFeedback, { optimizeTask: snapshot.optimizeTask } ); const apiMode = snapshot.apiMode || this.config.defaultApiMode; const result = await this.withLlmSlot(() => this.aiClient.requestJson( { apiKey: snapshot.apiKey || this.config.defaultApiKey, baseURL: snapshot.baseURL || this.config.defaultBaseURL, apiMode, model: snapshot.model || this.config.defaultModel, temperature: snapshot.temperature ?? 0, topP: snapshot.topP ?? 1, maxTokens: snapshot.maxTokens || 16384, enableThinking: false, clearThinking: true, }, [ { role: "system", content: system }, { role: "user", content: user }, ], { schema: apiMode === "responses" ? promptModule.RESULT_SCHEMA : undefined, schemaName: "task_optimizer_analyzer_result", stream: false, jsonMode: apiMode === "chat", jsonModeFallback: true, retries: 2, maxTokens: snapshot.maxTokens || 16384, }, ), ); return { mode, parsed, evidence, flags, result }; } private async persistExportableArtifacts(input: { item: ItemRow; itemDir: string; uploadBytes: Buffer; analysis: AnalyzerResult; taskPackage: TaskPackageResult; judgeResult?: JudgeRunResult | null; iterations: number; passed: boolean; iterationMemory?: IterationMemory | null; }): Promise { const evidencePath = join(input.itemDir, "evidence.json"); writeJson(evidencePath, { mode: input.analysis.mode, detectedMode: input.item.detected_mode, flags: input.analysis.flags, evidence: input.analysis.evidence, analyzerResult: input.analysis.result, iterations: input.iterations, passed: input.passed, iterationMemory: input.iterationMemory ?? null, }); const taskPackagePath = join(input.itemDir, "task-package.json"); writeJson(taskPackagePath, input.taskPackage); const paths: Parameters[1] = { evidencePath, taskPackagePath, }; if (input.judgeResult) { const judgeResultPath = join(input.itemDir, "judge-result.json"); writeJson(judgeResultPath, input.judgeResult); paths.judgeResultPath = judgeResultPath; } await this.store.setItemPaths(input.item.id, paths); const builtZip = await buildAnalysisZipBuffer({ originalZipBytes: input.uploadBytes, taskPackage: input.taskPackage, mode: input.analysis.mode, originalFilename: input.item.filename, }); const exportZipPath = join(input.itemDir, builtZip.filename); writeFileSync(exportZipPath, builtZip.buffer); await this.store.setItemPaths(input.item.id, { exportZipPath }); return exportZipPath; } private async failItem( itemId: string, message: string, lastStage?: string, errorDetailsPath?: string, ): Promise { await this.transitionItem(itemId, "failed", message, lastStage, errorDetailsPath); } private async transitionItem( itemId: string, status: ItemStatus, errorPreview?: string, lastStage?: string, errorDetailsPath?: string, ): Promise { await this.store.setItemStatus(itemId, status, errorPreview, { lastStage: lastStage ?? null, errorDetailsPath: errorDetailsPath ?? null, }); } private async withLlmSlot(fn: () => Promise): Promise { return this.llmSemaphore.run(fn); } } function parseOptionsSnapshot(raw: string | undefined): OptionsSnapshot { if (!raw) return {}; try { const parsed = JSON.parse(raw); return parsed && typeof parsed === "object" ? parsed : {}; } catch (_) { return {}; } } function resolveRequestedMode( snapshot: OptionsSnapshot, detectedMode: DetectedMode ): DetectedMode { return snapshot.requestedMode === "api" || snapshot.requestedMode === "mcp" ? snapshot.requestedMode : detectedMode; } function writeJson(path: string, value: unknown): void { writeFileSync(path, JSON.stringify(value, null, 2)); } function isJudgeTargetSatisfied( judgeResult: { total_score?: number; max_score?: number; has_zeros?: boolean; verdict?: string; dimensions?: Record; }, snapshot: OptionsSnapshot ): boolean { return evaluateJudgePolicy(judgeResult, snapshot).passed; } function buildJudgeFailureMessage( judgeResult: { total_score?: number; max_score?: number; dimensions?: Record }, taskPackage: Record, snapshot: OptionsSnapshot = {}, ): string { const maxScore = Number(judgeResult.max_score) || 12; const requirement = snapshot.strictFullMarks ? snapshot.ignoreTaskComplexityForFullMarks ? "要求除 TASK COMPLEXITY 外其他维度满分。" : "要求满分且无 0 分。" : "要求 >=10 且无 0 分。"; const dims = judgeResult.dimensions || {}; const weak = Object.entries(dims) .filter(([, value]) => Number(value && value.score) === 0) .map(([key, value]) => key + ": " + String((value && value.explanation) || "").slice(0, 240)); const limits = Array.isArray(taskPackage.evidence_limits) ? taskPackage.evidence_limits : []; const rerecordReason = String(taskPackage.rerecord_required_reason || ""); return [ "Judge 未达到 Chrome 插件标准:得分 " + Number(judgeResult.total_score || 0) + "/" + maxScore + "," + requirement, weak.length ? "0 分维度:" + weak.join(";") : "", rerecordReason ? "重录建议:" + rerecordReason : "", limits.length ? "证据限制:" + limits.map(String).join(";") : "", ].filter(Boolean).join("\n"); } async function resolvePluginMcpTools(parsed: any, log: Logger): Promise<{ tools: PluginMcpTool[]; source: string }> { const embedded = Array.isArray(parsed?.pluginMcpTools) ? parsed.pluginMcpTools : []; if (embedded.length) return { tools: embedded, source: parsed.pluginMcpToolSource || "task.json" }; const urls = collectMcpConfigUrls({ taskJson: parsed?.taskJson, networkJson: parsed?.networkJson }); const errors: string[] = []; for (const url of urls) { try { const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), 8000); const response = await fetch(url, { signal: controller.signal }); clearTimeout(timer); if (!response.ok) { errors.push(url + " -> HTTP " + response.status); continue; } const payload = await response.json(); const tools = collectMcpToolsFromTask(payload); if (tools.length) return { tools, source: url }; errors.push(url + " -> no tools"); } catch (err) { errors.push(url + " -> " + ((err as Error)?.message || String(err))); } } if (urls.length) { log.warn({ urls, errors: errors.slice(0, 8) }, "Failed to resolve plugin MCP tools from config URLs"); } return { tools: [], source: urls.length ? "mcp_config_unavailable" : (parsed?.pluginMcpToolSource || "unavailable"), }; } function assertStrictMcpEvidenceReady(evidence: Record, filename: string): void { if (evidence?.type !== "mcp") return; const status = evidence.mcpToolsStatus as | { available?: boolean; reason?: string; source?: string; matchedCount?: number; toolCount?: number } | undefined; const calls = Array.isArray(evidence.calls) ? evidence.calls : []; if (status && status.available === false) { throw new Error( [ "严格 MCP 不可用:" + filename + " 没有可用的 Chrome 插件 tools config。", "source=" + (status.source || "unknown"), status.reason ? "reason=" + status.reason : "", "请确认原始 task/metadata 中带有 MCP tools,或重新用插件采集包含 MCP config 的 recording。", ].filter(Boolean).join("\n"), ); } if (!calls.length) { throw new Error( [ "严格 MCP 没有匹配到任何插件可见 MCP call:" + filename, status ? `source=${status.source || "unknown"}, tools=${status.toolCount || 0}, matched=${status.matchedCount || 0}` : "", "后端不会再把 HTTP 请求伪造成 MCP call;请补采 MCP tools config 或重录。", ].filter(Boolean).join("\n"), ); } } export function isTransientItemError(err: unknown): boolean { const anyErr = err as { name?: string; message?: string; status?: number; retryable?: boolean }; const message = String(anyErr?.message || err || ""); if (/严格 MCP|Judge 未达到|zip 缺少|zip 格式未识别|校验失败|schema invalid|not valid JSON/i.test(message)) { return false; } if (anyErr?.retryable === true) return true; const status = Number(anyErr?.status || 0); if (status === 429 || (status >= 500 && status < 600)) return true; return /AbortError|aborted|timeout|ECONNRESET|ECONNREFUSED|ENOTFOUND|EAI_AGAIN|fetch failed|terminated/i.test( String(anyErr?.name || "") + " " + message, ); } function serializeError(err: unknown, context: Record): Record { const anyErr = err as { name?: string; message?: string; stack?: string; rawPreview?: string; status?: number; endpoint?: string; apiMode?: string; cause?: unknown; }; return { ...context, name: anyErr?.name || (err && typeof err === "object" ? err.constructor?.name : typeof err), message: anyErr?.message || String(err), rawPreview: anyErr?.rawPreview, status: anyErr?.status, endpoint: anyErr?.endpoint, apiMode: anyErr?.apiMode, cause: anyErr?.cause instanceof Error ? { name: anyErr.cause.name, message: anyErr.cause.message, stack: anyErr.cause.stack } : anyErr?.cause, stack: anyErr?.stack, aborted: /AbortError|aborted|abort/i.test(String(anyErr?.name || "") + " " + String(anyErr?.message || "")), cancelled: /cancel|stopped/i.test(String(anyErr?.message || "")), terminated: /terminated/i.test(String(anyErr?.message || "")), }; } function errorPreview(err: unknown, details?: Record): string { const anyErr = err as { rawPreview?: string; message?: string }; const stage = details?.stage ? "stage=" + String(details.stage) + "\n" : ""; const message = anyErr?.rawPreview || anyErr?.message || String(err); return safePreview(stage + message).slice(0, 2048); }