Rl-Auto / apps /api /src /queue /runner.ts
Lazywords's picture
Deploy RL Auto Docker Space
c4ae742
Raw
History Blame Contribute Delete
32.1 kB
// 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<string, unknown>;
flags: Array<{ code: string; note: string; evidence: unknown }>;
result: unknown;
};
type TaskPackageResult = ReturnType<typeof buildTaskPackage>;
type JudgeRunResult = Awaited<ReturnType<typeof runJudge>>;
/**
* 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<void> {
if (this.limit <= 0) return;
if (this.active < this.limit) {
this.active += 1;
return;
}
await new Promise<void>((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<T>(fn: () => Promise<T>): Promise<T> {
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<string, { score?: number; explanation?: string }>;
},
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<string, InFlightJob>();
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<void> {
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<void> {
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<void> {
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<void> {
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<AnalyzerResult> {
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<string, unknown>;
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<string, unknown>;
assertStrictMcpEvidenceReady(evidence, item.filename);
} else {
evidence = { type: "http", ...buildApiEvidence(parsed as never) } as Record<string, unknown>;
}
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<string> {
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<Store["setItemPaths"]>[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<void> {
await this.transitionItem(itemId, "failed", message, lastStage, errorDetailsPath);
}
private async transitionItem(
itemId: string,
status: ItemStatus,
errorPreview?: string,
lastStage?: string,
errorDetailsPath?: string,
): Promise<void> {
await this.store.setItemStatus(itemId, status, errorPreview, {
lastStage: lastStage ?? null,
errorDetailsPath: errorDetailsPath ?? null,
});
}
private async withLlmSlot<T>(fn: () => Promise<T>): Promise<T> {
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<string, { score?: number; explanation?: string }>;
},
snapshot: OptionsSnapshot
): boolean {
return evaluateJudgePolicy(judgeResult, snapshot).passed;
}
function buildJudgeFailureMessage(
judgeResult: { total_score?: number; max_score?: number; dimensions?: Record<string, { score?: number; explanation?: string }> },
taskPackage: Record<string, unknown>,
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<string, unknown>, 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<string, unknown>): Record<string, unknown> {
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, unknown>): 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);
}