diff --git a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts index a76f737f5..d7dba03a6 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts @@ -1,5 +1,4 @@ import { randomUUID } from "node:crypto"; -import { promises as fs } from "node:fs"; import { basename, dirname, join } from "node:path"; import { createDeepScanArtifacts, @@ -11,16 +10,11 @@ import { type ScanDraftInput } from "../artifact-scan-draft.js"; import type { DeepScanArtifacts } from "./artifacts.js"; -import { - DeepScanWorkerRunner, - sha256 -} from "./worker-runner.js"; +import { DeepScanWorkerRunner } from "./worker-runner.js"; import type { AcceptedDiscovery, DedupOutcome, - DiscoveryOutcome, - SuccessfulDedupOutcome, - WorkerExecutionAudit + DiscoveryOutcome } from "./worker-runner.js"; import { boundedDeepScanErrorPair, @@ -48,30 +42,14 @@ type SchedulerSettlement = | { status: "fulfilled"; outcome: SchedulerOutcome } | { status: "rejected"; error: unknown }; -type AcceptedReducer = Omit; - interface SchedulerResult { reason: DeepScanTerminalReason; omittedWorkerIds: string[]; - canceledWorkerIds: string[]; - accepted: AcceptedDiscovery[]; - mergedWorkerIds: string[]; - reducers: AcceptedReducer[]; result?: DeepReductionInput; } type CoordinatorPhase = "setup" | "discovery" | "terminal"; -interface SchedulerAudit { - accepted: AcceptedDiscovery[]; - mergedWorkerIds: string[]; - omittedWorkerIds: string[]; - canceledWorkerIds: string[]; - bufferedWorkerIds: string[]; - reducers: AcceptedReducer[]; - executions: WorkerExecutionAudit[]; -} - export interface CoordinatorOptions { run: DeepScanRunState; store: DeepScanStore; @@ -124,15 +102,6 @@ export class DeepScanCoordinator { private externallyFailed = false; private phase: CoordinatorPhase = "setup"; private discoveryDeadlineReached = false; - private readonly audit: SchedulerAudit = { - accepted: [], - mergedWorkerIds: [], - omittedWorkerIds: [], - canceledWorkerIds: [], - bufferedWorkerIds: [], - reducers: [], - executions: [] - }; private state: DeepScanRunState; constructor(private readonly options: CoordinatorOptions) { @@ -149,10 +118,7 @@ export class DeepScanCoordinator { clock: this.clock, random: options.random ?? Math.random, log: this.log, - retryDelaysMs: options.retryDelaysMs ?? RETRY_DELAYS_MS, - recordExecution: (execution) => { - this.audit.executions.push(execution); - } + retryDelaysMs: options.retryDelaysMs ?? RETRY_DELAYS_MS } satisfies Omit[0], "signal">; this.workers = new DeepScanWorkerRunner({ ...workerOptions, @@ -603,23 +569,12 @@ export class DeepScanCoordinator { .filter((worker) => worker.kind === "discovery" && worker.mergeState === "merged") .map((worker) => worker.id) ); - const mergedDiscoveries: AcceptedDiscovery[] = recovered.filter((worker) => ( - mergedIds.has(worker.id) - )); - const canceledWorkerIds = (this.state.persistedWorkers ?? []) - .filter((worker) => worker.kind === "discovery" && worker.status === "canceled") - .map((worker) => worker.id); const omittedWorkerIds: string[] = []; - this.audit.accepted = [...accepted]; - this.audit.mergedWorkerIds = mergedDiscoveries.map((worker) => worker.id); - this.audit.canceledWorkerIds = [...canceledWorkerIds]; - this.audit.executions = await this.recoverPersistedExecutions(); const recoveredReducers = await this.recoverCompletedReducers(recovered); - const reducerOutcomes = recoveredReducers.reducers; let latestResult = recoveredReducers.result; let buffer: AcceptedDiscovery[] = recovered.filter((worker) => !mergedIds.has(worker.id)); let reducer: Promise | undefined; - let previousReducerResultPath = reducerOutcomes.at(-1)?.resultPath; + let previousReducerResultPath = recoveredReducers.resultPath; let dispatched = this.state.dispatchedCount; let workerSequence = Math.max( dispatched, @@ -636,8 +591,6 @@ export class DeepScanCoordinator { let stopReason: DeepScanTerminalReason | undefined; let lastReplaceableFailure: Extract | undefined; - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); - this.audit.reducers = [...reducerOutcomes]; const errorLimit = config.stopAfterConsecutiveErrors ?? config.stopAfterNoNew; let reducerFailures = persistedReducerFailureStreak(this.state.persistedWorkers ?? []); if (this.state.consecutiveErrors >= errorLimit) { @@ -707,16 +660,12 @@ export class DeepScanCoordinator { const workerId = entries[index]?.[0]; if (workerId) active.delete(workerId); if (result.status === "rejected") { - if (workerId) removeValue(canceledWorkerIds, workerId); firstFailure ??= result.reason; continue; } const outcome = result.value; if (outcome.status === "failed") { - if (outcome.replaceableFailureKind) { - canceledWorkerIds.push(outcome.workerId); - } else { - removeValue(canceledWorkerIds, outcome.workerId); + if (!outcome.replaceableFailureKind) { firstFailure ??= outcome.error; } } else if (outcome.status === "succeeded") { @@ -728,15 +677,8 @@ export class DeepScanCoordinator { } else if (!buffer.some((worker) => worker.id === outcome.worker.id)) { buffer.push(outcome.worker); } - removeValue(canceledWorkerIds, outcome.worker.id); - } else { - canceledWorkerIds.push(outcome.workerId); } } - this.audit.accepted = [...accepted]; - this.audit.omittedWorkerIds = unique(omittedWorkerIds); - this.audit.canceledWorkerIds = unique(canceledWorkerIds); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); return firstFailure; }; const reconcileReducerSettlement = async (): Promise => { @@ -750,18 +692,11 @@ export class DeepScanCoordinator { const outcome = result.value; if ("status" in outcome) { buffer = [...outcome.consumed, ...buffer].sort(compareCompletionSequence); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); return outcome.error; } this.state = outcome.run; previousReducerResultPath = outcome.resultPath; - mergedDiscoveries.push(...outcome.consumed); - const { result: acceptedResult, ...metadata } = outcome; - latestResult = acceptedResult; - reducerOutcomes.push(metadata); - this.audit.reducers = [...reducerOutcomes]; - this.audit.mergedWorkerIds = unique(mergedDiscoveries.map((worker) => worker.id)); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); + latestResult = outcome.result; return undefined; }; @@ -846,8 +781,6 @@ export class DeepScanCoordinator { const consecutiveErrors = outcome.consecutiveErrors ?? (this.state.consecutiveErrors ?? 0) + 1; this.state = { ...this.state, consecutiveErrors }; - canceledWorkerIds.push(outcome.workerId); - this.audit.canceledWorkerIds = unique(canceledWorkerIds); this.log({ event: "discovery_worker_replaced", scanId: this.state.scanId, @@ -874,8 +807,6 @@ export class DeepScanCoordinator { throw outcome.error; } if (outcome.status === "canceled") { - canceledWorkerIds.push(outcome.workerId); - this.audit.canceledWorkerIds = unique(canceledWorkerIds); if ( !this.abortController.signal.aborted && !this.discoveryAbortController.signal.aborted @@ -887,8 +818,6 @@ export class DeepScanCoordinator { accepted.push(outcome.worker); this.state = { ...this.state, consecutiveErrors: 0 }; buffer.push(outcome.worker); - this.audit.accepted = [...accepted]; - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); this.logProgress(accepted.length); continue; } @@ -896,7 +825,6 @@ export class DeepScanCoordinator { reducer = undefined; if ("status" in outcome) { buffer = [...outcome.consumed, ...buffer].sort(compareCompletionSequence); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); reducerFailures += 1; this.log({ event: "dedup_worker_replaced", @@ -919,22 +847,13 @@ export class DeepScanCoordinator { reducerFailures = 0; this.state = outcome.run; previousReducerResultPath = outcome.resultPath; - mergedDiscoveries.push(...outcome.consumed); - const { result: acceptedResult, ...metadata } = outcome; - latestResult = acceptedResult; - reducerOutcomes.push(metadata); - this.audit.reducers = [...reducerOutcomes]; - this.audit.mergedWorkerIds = unique(mergedDiscoveries.map((worker) => worker.id)); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); + latestResult = outcome.result; if ( !this.discoveryDeadlineReached && outcome.run.noNewStreak >= config.stopAfterNoNew && buffer.length === 0 ) { stopReason = "saturated"; - canceledWorkerIds.push(...active.keys()); - this.audit.canceledWorkerIds = unique(canceledWorkerIds); - this.audit.bufferedWorkerIds = []; this.abortController.abort("deep_scan_saturated"); } } @@ -943,12 +862,6 @@ export class DeepScanCoordinator { // the manifest records which results completed and which were canceled. const lateFailure = await reconcileRemainingDiscoveries("omitted"); - this.audit.accepted = [...accepted]; - this.audit.mergedWorkerIds = unique(mergedDiscoveries.map((worker) => worker.id)); - this.audit.omittedWorkerIds = unique(omittedWorkerIds); - this.audit.canceledWorkerIds = unique(canceledWorkerIds); - this.audit.bufferedWorkerIds = buffer.map((worker) => worker.id); - // Once Deep reaches saturation, late worker errors cannot fail the scan. if (lateFailure && stopReason !== "saturated") throw lateFailure; @@ -961,10 +874,6 @@ export class DeepScanCoordinator { return { reason: stopReason, omittedWorkerIds: unique(omittedWorkerIds), - canceledWorkerIds: unique(canceledWorkerIds), - accepted, - mergedWorkerIds: unique(mergedDiscoveries.map((worker) => worker.id)), - reducers: reducerOutcomes, result: latestResult, }; } @@ -981,7 +890,6 @@ export class DeepScanCoordinator { worker.resultManifestPath, this.state.scanId ); - const evidence = await persistedWorkerEvidence(worker); recovered.push({ id: worker.id, label: basename(dirname(worker.promptPath)), @@ -989,8 +897,7 @@ export class DeepScanCoordinator { resultPath: worker.resultManifestPath, completionSequence: worker.completionSequence, attempt: worker.attempt, - ...(worker.threadId ? { threadId: worker.threadId } : {}), - ...evidence + ...(worker.threadId ? { threadId: worker.threadId } : {}) }); } return recovered.sort(compareCompletionSequence); @@ -998,10 +905,10 @@ export class DeepScanCoordinator { private async recoverCompletedReducers( discoveries: AcceptedDiscovery[] - ): Promise<{ reducers: AcceptedReducer[]; result?: DeepReductionInput }> { + ): Promise<{ resultPath?: string; result?: DeepReductionInput }> { const discoveriesById = new Map(discoveries.map((worker) => [worker.id, worker])); const inputs = this.state.persistedDedupInputs ?? []; - const outcomes: AcceptedReducer[] = []; + let resultPath: string | undefined; let latestResult: DeepReductionInput | undefined; const completedReducers = (this.state.persistedWorkers ?? []) .filter((worker) => worker.kind === "dedup" && worker.status === "succeeded") @@ -1009,7 +916,6 @@ export class DeepScanCoordinator { workerLabelSequence(left, "dedup") - workerLabelSequence(right, "dedup") || left.id.localeCompare(right.id) )); - let noNewStreak = 0; for (const worker of completedReducers) { if (!worker.resultManifestPath) { throw new Error(`Completed reducer ${worker.id} has no persisted result manifest.`); @@ -1021,64 +927,17 @@ export class DeepScanCoordinator { if (consumed.length === 0 || consumed.some((value) => !value)) { throw new Error(`Completed reducer ${worker.id} has incomplete persisted inputs.`); } - const accepted = consumed as AcceptedDiscovery[]; - const { newFindings, result } = await validateReducerArtifacts({ + const { result } = await validateReducerArtifacts({ artifacts: this.artifacts, artifactDir: worker.artifactDir, resultPath: worker.resultManifestPath, reducerId: worker.id, - previousReducerResultPath: outcomes.at(-1)?.resultPath + previousReducerResultPath: resultPath }, this.state.scanId); latestResult = result; - noNewStreak = newFindings > 0 ? 0 : noNewStreak + accepted.length; - const evidence = await persistedWorkerEvidence(worker); - outcomes.push({ - type: "dedup", - id: worker.id, - consumed: accepted, - resultPath: worker.resultManifestPath, - newFindings, - attempt: worker.attempt, - ...(worker.threadId ? { threadId: worker.threadId } : {}), - ...evidence, - run: { ...this.state, noNewStreak } - }); + resultPath = worker.resultManifestPath; } - return { reducers: outcomes, result: latestResult }; - } - - private async recoverPersistedExecutions(): Promise { - const executions: WorkerExecutionAudit[] = []; - for (const worker of this.state.persistedWorkers ?? []) { - if ( - worker.kind === "setup" - || worker.status === "queued" - || worker.status === "running" - || (worker.status === "canceled" && worker.attempt === 0) - ) { - continue; - } - const replaceableFailure = persistedReplaceableFailure(worker); - const status = replaceableFailure || worker.status === "failed" - ? "failed" - : worker.status; - executions.push({ - id: worker.id, - label: basename(dirname(worker.promptPath)), - kind: worker.kind, - status, - attempt: worker.attempt, - ...(worker.threadId ? { threadId: worker.threadId } : {}), - promptPath: worker.promptPath, - artifactDir: worker.artifactDir, - ...await persistedWorkerEvidence(worker), - ...(status === "failed" && worker.error - ? { error: replaceableFailure?.message ?? worker.error } - : {}), - ...(replaceableFailure ? { failureKind: replaceableFailure.kind } : {}) - }); - } - return executions; + return { resultPath, result: latestResult }; } private reducerReady( @@ -1186,12 +1045,6 @@ function unique(values: string[]): string[] { return [...new Set(values)]; } -function removeValue(values: string[], value: string): void { - for (let index = values.length - 1; index >= 0; index -= 1) { - if (values[index] === value) values.splice(index, 1); - } -} - function cloneState(state: DeepScanRunState): DeepScanRunState { return { ...state, config: { ...state.config } }; } @@ -1238,30 +1091,6 @@ function persistedReplaceableFailure( return undefined; } -async function persistedWorkerEvidence(worker: PersistedDeepScanWorker): Promise<{ - basePromptSha256: string; - attemptPromptPaths: string[]; -}> { - const attemptPromptPaths = [worker.promptPath]; - for (let attempt = 2; attempt <= worker.attempt; attempt += 1) { - const promptPath = join( - dirname(worker.promptPath), - "prompts", - `attempt-${String(attempt).padStart(2, "0")}.md` - ); - try { - await fs.access(promptPath); - attemptPromptPaths.push(promptPath); - } catch { - // Transient execution retries reuse the original prompt. - } - } - return { - basePromptSha256: sha256(await fs.readFile(worker.promptPath, "utf8")), - attemptPromptPaths - }; -} - function discoveryErrorLimitError( count: number, limit: number, diff --git a/plugins/codex-security/mcp-app/src/deep-scan/worker-runner.ts b/plugins/codex-security/mcp-app/src/deep-scan/worker-runner.ts index 82f119979..51738be8f 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/worker-runner.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/worker-runner.ts @@ -1,4 +1,3 @@ -import { createHash } from "node:crypto"; import { promises as fs } from "node:fs"; import { dirname, join } from "node:path"; import { getCodexSecurityDeepReducerInputs } from "../artifact-deep-reducer.js"; @@ -41,8 +40,6 @@ export interface AcceptedDiscovery { completionSequence: number; attempt: number; threadId?: string; - basePromptSha256: string; - attemptPromptPaths: string[]; } export type DiscoveryOutcome = @@ -66,8 +63,6 @@ export interface SuccessfulDedupOutcome { newFindings: number; attempt: number; threadId?: string; - basePromptSha256: string; - attemptPromptPaths: string[]; run: DeepScanRunState; } @@ -81,22 +76,6 @@ export interface FailedDedupOutcome { export type DedupOutcome = SuccessfulDedupOutcome | FailedDedupOutcome; -/** Audit evidence for every logical SDK execution, including failures and cancellation. */ -export interface WorkerExecutionAudit { - id: string; - label: string; - kind: DeepScanWorkerKind; - status: "succeeded" | "failed" | "canceled"; - attempt: number; - threadId?: string; - promptPath: string; - artifactDir: string; - basePromptSha256: string; - attemptPromptPaths: string[]; - error?: string; - failureKind?: DeepScanReplaceableFailureKind; -} - export interface ReducerRequest { id: string; label: string; @@ -115,13 +94,11 @@ export interface DeepScanWorkerRunnerOptions { log: DeepScanLogger; retryDelaysMs: readonly number[]; signal: AbortSignal; - recordExecution?: (execution: WorkerExecutionAudit) => void; } interface WorkerAttemptEvidence { attempt: number; threadId?: string; - attemptPromptPaths: string[]; } type WorkerAttemptOutcome = @@ -207,15 +184,6 @@ export class DeepScanWorkerRunner { if (!discoveryValidated) { await fs.rm(files.resultPath, { force: true }); } - const basePromptSha256 = sha256(basePrompt); - this.recordExecution({ - id: workerId, - label: workerLabel, - kind: "discovery", - promptPath, - artifactDir, - basePromptSha256 - }, outcome); if (outcome.status === "failed") { return { type: "discovery", @@ -283,9 +251,7 @@ export class DeepScanWorkerRunner { resultPath: files.resultPath, completionSequence: persisted.completionSequence, attempt: outcome.attempt, - threadId: outcome.threadId, - basePromptSha256, - attemptPromptPaths: outcome.attemptPromptPaths + threadId: outcome.threadId } }; } @@ -377,15 +343,6 @@ export class DeepScanWorkerRunner { }, outcome.attempt, outcome.threadId); outcome = { ...outcome, status: "canceled" }; } - const basePromptSha256 = sha256(basePrompt); - this.recordExecution({ - id: reducerId, - label: reducerLabel, - kind: "dedup", - promptPath, - artifactDir, - basePromptSha256 - }, outcome); if (outcome.status === "failed") { if (outcome.error instanceof DeepScanNonRetryableError) throw outcome.error; return { @@ -447,8 +404,6 @@ export class DeepScanWorkerRunner { newFindings: reducerValidation.newFindings, attempt: outcome.attempt, threadId: outcome.threadId, - basePromptSha256, - attemptPromptPaths: outcome.attemptPromptPaths, run: committed }; } @@ -470,10 +425,9 @@ export class DeepScanWorkerRunner { let continuationPrompt: string | undefined; let lastThreadId: string | undefined; let executionPromptPath = input.promptPath; - const attemptPromptPaths = [input.promptPath]; for (let attempt = 1; attempt <= maximumAttempts; attempt += 1) { if (signal.aborted) { - return await this.cancelAttempt(input, attempt, lastThreadId, attemptPromptPaths); + return await this.cancelAttempt(input, attempt, lastThreadId); } let validationStarted = false; let validationCompleted = false; @@ -527,7 +481,7 @@ export class DeepScanWorkerRunner { } }); if (signal.aborted) { - return await this.cancelAttempt(input, attempt, activeThreadId, attemptPromptPaths); + return await this.cancelAttempt(input, attempt, activeThreadId); } validationStarted = true; try { @@ -537,7 +491,7 @@ export class DeepScanWorkerRunner { } validationCompleted = true; if (signal.aborted) { - return await this.cancelAttempt(input, attempt, activeThreadId, attemptPromptPaths); + return await this.cancelAttempt(input, attempt, activeThreadId); } this.options.log({ event: "worker_succeeded", @@ -549,12 +503,11 @@ export class DeepScanWorkerRunner { return { status: "succeeded", attempt, - threadId: result.threadId ?? activeThreadId, - attemptPromptPaths: [...attemptPromptPaths] + threadId: result.threadId ?? activeThreadId }; } catch (error) { if (signal.aborted) { - return await this.cancelAttempt(input, attempt, activeThreadId, attemptPromptPaths); + return await this.cancelAttempt(input, attempt, activeThreadId); } const normalized = asError(error); const policyRefusal = input.kind === "discovery" @@ -588,8 +541,7 @@ export class DeepScanWorkerRunner { ? {} : { consecutiveErrors: persistedFailure.consecutiveErrors }), attempt, - threadId: activeThreadId, - attemptPromptPaths: [...attemptPromptPaths] + threadId: activeThreadId }; } await this.options.store.updateWorker({ @@ -627,7 +579,6 @@ export class DeepScanWorkerRunner { failedAttempt: attempt, error: normalized }); - attemptPromptPaths.push(executionPromptPath); } } const delayMs = Math.ceil( @@ -645,7 +596,7 @@ export class DeepScanWorkerRunner { await this.options.clock.sleep(delayMs, signal); } catch (sleepError) { if (signal.aborted) { - return await this.cancelAttempt(input, attempt, activeThreadId, attemptPromptPaths); + return await this.cancelAttempt(input, attempt, activeThreadId); } throw sleepError; } @@ -714,36 +665,15 @@ export class DeepScanWorkerRunner { artifactDir: string; }, attempt: number, - threadId: string | undefined, - attemptPromptPaths: string[] + threadId: string | undefined ): Promise { await this.persistWorkerCancellation(input, attempt, threadId); return { status: "canceled", attempt, - threadId, - attemptPromptPaths: [...attemptPromptPaths] + threadId }; } - - private recordExecution( - input: Omit, - outcome: WorkerAttemptOutcome - ): void { - this.options.recordExecution?.({ - ...input, - status: outcome.status, - attempt: outcome.attempt, - ...(outcome.threadId ? { threadId: outcome.threadId } : {}), - attemptPromptPaths: [...outcome.attemptPromptPaths], - ...(outcome.status === "failed" ? { - error: outcome.error.message, - ...(outcome.replaceableFailureKind - ? { failureKind: outcome.replaceableFailureKind } - : {}) - } : {}) - }); - } } /** @@ -880,10 +810,6 @@ function validationErrorData(error: Error, message: string): Record ( + [...store.workers.values()].some((worker) => ( + worker.kind === "discovery" && worker.status === "succeeded" + )) + && originalExecutor.discoveryCalls >= 2 + )); + const accepted = [...store.workers.values()].find((worker) => ( + worker.kind === "discovery" && worker.status === "succeeded" + )); + assert.ok(accepted); + + // The waiter detached earlier; an app update now removes its MCP process + // without canceling or finalizing the persisted scan. + original.cancel("mcp server process restarted"); + await eventually(() => originalExecutor.runningDiscovery === 0); + const persistedWorkers = [...store.workers.values()].map((worker) => structuredClone(worker)); + const independentReviews = { + completed: persistedWorkers.filter((worker) => ( + worker.kind === "discovery" && worker.status === "succeeded" + )).length, + active: persistedWorkers.filter((worker) => ( + worker.kind === "discovery" && worker.status === "running" + )).length, + consolidating: persistedWorkers.some((worker) => ( + worker.kind === "dedup" && worker.status === "running" + )) + }; + assert.deepEqual(independentReviews, { completed: 1, active: 0, consolidating: false }); + assert.equal(store.run.status, "running"); + assert.equal(store.run.phase, "discovery"); + assert.equal(store.finishCalls.length, 0); + assert.equal(store.failCalls, 0); + assert.equal(store.run.manifestPath, undefined); + await assert.rejects( + readFile(path.join( + fixture.run.scanDir, + "artifacts", + "deep_discovery", + "coordinator-manifest.json" + )), + { code: "ENOENT" } + ); + + store.run = { + ...store.run, + dispatchedCount: 1, + persistedWorkers + }; + const continuationClaims = []; + store.claimCoordinator = async (input) => { + continuationClaims.push(structuredClone(input)); + assert.equal(input.handoffClaimToken, handoffClaimToken); + store.run = { + ...store.run, + coordinatorGeneration: 3, + updatedAt: new Date().toISOString() + }; + return { acquired: true, run: structuredClone(store.run) }; + }; + store.heartbeatCoordinator = async () => structuredClone(store.run); + const replacementExecutor = new FakeExecutor({ discoveryCandidateId: "candidate-next" }); + const acceptedResult = await readFile(accepted.resultManifestPath, "utf8"); + const acceptedWorkerId = await workerIdFromPrompt(accepted.promptPath); + if (removeHistoricalPrompts) { + await Promise.all(persistedWorkers.map((worker) => rm(worker.promptPath, { force: true }))); + } + const resumed = await startOrJoinDeepScanCoordinator({ + begin: { run: structuredClone(store.run), shouldStart: false }, + registry: new DeepScanCoordinatorRegistry(), + options: { + store, + executor: replacementExecutor, + pluginRoot: fixture.pluginRoot, + clock: immediateClock, + threadId: "track-c-owning-thread", + handoffClaimToken + } + }); + const terminal = await resumed.coordinator.wait(undefined, 5_000); + + assert.equal(continuationClaims.length, 1); + assert.equal(continuationClaims[0].handoffClaimToken, handoffClaimToken); + assert.equal(terminal?.status, "succeeded", terminal?.error); + assert.equal(store.failCalls, 0); + assert.equal(replacementExecutor.logicalDiscoveryWorkers.size, 1); + assert.equal( + replacementExecutor.logicalDiscoveryWorkers.has(acceptedWorkerId), + false + ); + assert.equal(store.dedupClaims.length, 1); + assert.equal(store.dedupClaims[0].workerIds.includes(accepted.id), true); + const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); + assert.equal(manifest.scan.scanId, fixture.run.scanId); + assert.equal(store.dedupClaims[0].workerIds.length, 2); + assert.equal(await readFile(accepted.resultManifestPath, "utf8"), acceptedResult); + assert.deepEqual(store.dedupClaims[0].workerIds, [ + accepted.id, + ...[...store.workers.values()] + .filter((worker) => worker.kind === "discovery" && worker.status === "succeeded" && worker.id !== accepted.id) + .map((worker) => worker.id) + ]); + assert.deepEqual(manifest.findings.map((finding) => finding.provenance.candidateId), [ + "candidate-original", "candidate-next" + ]); + const newPrompt = [...replacementExecutor.discoveryPromptPaths.values()][0].values().next().value; + assert.equal((await promptContext(newPrompt)).userContext, originalInput.userContext); + assert.deepEqual(store.run.config, originalInput.config); + assert.equal(store.run.createdAt, originalInput.createdAt); + assert.equal(store.run.targetPath, originalInput.targetPath); + assert.equal(store.run.scope, originalInput.scope); + } + + async function testResumedDiscoveryDeadlineUsesPersistedCreationTime( + alreadyExpired = false, + maxTimeHours + ) { + const fixture = await fixtureRun({ + workers: 1, + subagents: 0, + stopAfterNoNew: 99, + maxDiscoveryRuns: 8, + ...(maxTimeHours === undefined ? {} : { maxTimeHours }) + }); + const discoveryTimeoutMs = (maxTimeHours ?? 96) * 60 * 60 * 1_000; + let currentTime = immediateClock.now(); + const clock = { + now: () => currentTime, + sleep: immediateClock.sleep + }; + const createdAt = new Date(currentTime - discoveryTimeoutMs + 30_000).toISOString(); + const store = new FakeStore({ + ...fixture.run, + createdAt, + phase: "setup", + coordinatorGeneration: 2 + }); + const originalExecutor = new FakeExecutor({ + blockDiscoveryAfterCalls: 1, + discoveryCandidateId: "candidate-1" + }); + const original = new DeepScanCoordinator({ + run: store.run, + store, + executor: originalExecutor, + pluginRoot: fixture.pluginRoot, + clock + }); + original.start(); + await eventually(() => ( + originalExecutor.discoveryCalls === 2 + && originalExecutor.runningDiscovery === 1 + && [...store.workers.values()].some((worker) => ( + worker.kind === "discovery" && worker.status === "succeeded" + )) + )); + + original.cancel("mcp server process restarted"); + await eventually(() => originalExecutor.runningDiscovery === 0); + assert.equal(store.run.status, "running"); + assert.equal(store.run.createdAt, createdAt); + currentTime = Date.parse(createdAt) + discoveryTimeoutMs + (alreadyExpired ? 1_000 : -1_000); + store.run = { + ...store.run, + persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)) + }; + await Promise.all(store.run.persistedWorkers.map((worker) => rm(worker.promptPath, { force: true }))); + + const resumedExecutor = new FakeExecutor({ + blockDiscovery: true, + canonicalCandidateId: "candidate-1", + dedupNewFindings: [1] + }); + const resumed = new DeepScanCoordinator({ + run: store.run, + store, + executor: resumedExecutor, + pluginRoot: fixture.pluginRoot, + clock + }); + resumed.start(); + if (!alreadyExpired) await resumedExecutor.discoveryStarted; + + const terminal = await resumed.wait(undefined, 5_000); + assert.equal(terminal?.status, "succeeded"); + assert.equal(terminal?.terminalReason, "capped"); + assert.equal(store.failCalls, 0); + assert.equal(resumedExecutor.discoveryCalls, alreadyExpired ? 0 : 1); + assert.equal(resumedExecutor.dedupCalls, 1); + assert.equal(resumedExecutor.runningDiscovery, 0); + + const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); + assert.equal(store.run.config.maxTimeHours, maxTimeHours); + assert.equal(store.dedupClaims[0].workerIds.length, 1); + assert.equal([...store.workers.values()].some((worker) => worker.status === "canceled"), true); + assert.deepEqual(manifest.findings.map((finding) => finding.provenance.candidateId), ["candidate-1"]); + } + + async function testResumedManifestPreservesCompletedReducer( + includeUnstartedReducer = false, + removeHistoricalPrompts = true + ) { + const fixture = await fixtureRun({ + workers: 2, + subagents: 0, + stopAfterNoNew: 2, + maxDiscoveryRuns: 3 + }); + const store = new FakeStore({ ...fixture.run, phase: "setup" }); + store.blockDedupCommitResponse = true; + const original = new DeepScanCoordinator({ + run: store.run, + store, + executor: new FakeExecutor({ + discoveryCandidateId: "candidate-accepted", + blockDiscoveryAfterCalls: 2 + }), + pluginRoot: fixture.pluginRoot, + clock: immediateClock + }); + original.start(); + await store.dedupCommitPersisted.promise; + original.cancel("mcp server process restarted"); + store.releaseDedupCommitResponse(); + await original.settled(); + await eventually(() => [...store.workers.values()].every((worker) => ( + worker.status !== "queued" && worker.status !== "running" + ))); + + let unstartedReducer; + if (includeUnstartedReducer) { + const artifactDir = path.join( + fixture.run.scanDir, + "artifacts", + "deep_discovery", + "dedup", + "dedup-unstarted", + "output" + ); + await mkdir(artifactDir, { recursive: true }); + const promptPath = path.join(path.dirname(artifactDir), "prompt.md"); + await writeFile(promptPath, "Reducer claimed before coordinator restart.\n"); + unstartedReducer = { + id: randomUUID(), + kind: "dedup", + status: "canceled", + promptPath, + artifactDir, + attempt: 0, + mergeState: "none" + }; + store.workers.set(unstartedReducer.id, unstartedReducer); + } + + store.run = { + ...store.run, + status: "running", + phase: "discovery", + persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)), + persistedDedupInputs: store.dedupClaims.flatMap((claim) => ( + claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ + dedupWorkerId: claim.id, + discoveryWorkerId, + inputOrder + })) + )) + }; + const acceptedWorkers = store.run.persistedWorkers.filter((worker) => worker.status === "succeeded"); + const acceptedResults = await Promise.all(acceptedWorkers.map((worker) => readFile(worker.resultManifestPath, "utf8"))); + const persistedInputs = structuredClone(store.run.persistedDedupInputs); + const persistedWorkers = structuredClone(store.run.persistedWorkers); + const committedReducer = store.workers.get(store.dedupClaims[0].id); + const committedResult = JSON.parse(await readFile(committedReducer.resultManifestPath, "utf8")); + if (removeHistoricalPrompts) await rm(committedReducer.promptPath); + const replacementExecutor = new FakeExecutor(); + const replacement = new DeepScanCoordinator({ + run: store.run, + store, + executor: replacementExecutor, + pluginRoot: fixture.pluginRoot, + clock: immediateClock + }); + replacement.start(); + + const terminal = await replacement.wait(undefined, 5_000); + assert.equal(terminal?.status, "succeeded", terminal?.error); + const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); + assert.equal(manifest.scan.scanId, fixture.run.scanId); + assert.equal(store.dedupCommits.length, 1); + assert.equal(store.dedupClaims.length, 1); + assert.equal(replacementExecutor.calls, 0, "accepted work needs no new worker launch"); + assert.deepEqual(manifest.findings, committedResult.findings); + assert.deepEqual(store.run.persistedDedupInputs, persistedInputs); + assert.deepEqual(store.run.persistedWorkers, persistedWorkers); + assert.deepEqual( + await Promise.all(acceptedWorkers.map((worker) => readFile(worker.resultManifestPath, "utf8"))), + acceptedResults + ); + assert.equal(store.run.persistedWorkers.some((worker) => worker.id === unstartedReducer?.id), + includeUnstartedReducer); + } + + async function testResumeUsesHistoricalCandidateSnapshotForEachReducer() { + const fixture = await fixtureRun({ + workers: 3, + subagents: 0, + stopAfterNoNew: 10, + maxDiscoveryRuns: 5 + }); + const store = new FakeStore(fixture.run); + const original = new DeepScanCoordinator({ + run: fixture.run, + store, + executor: new FakeExecutor({ + dedupNewFindings: [1, 0], + dedupEvidenceByCall: ["first reducer evidence", "final reducer evidence"] + }), + pluginRoot: fixture.pluginRoot, + clock: immediateClock + }); + original.start(); + assert.equal((await original.wait(undefined, 5_000))?.status, "succeeded"); + assert.equal(store.dedupClaims.length >= 2, true); + + store.run = { + ...store.run, + status: "running", + phase: "discovery", + terminalReason: undefined, + manifestPath: undefined, + persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)), + persistedDedupInputs: store.dedupClaims.flatMap((claim) => ( + claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ + dedupWorkerId: claim.id, + discoveryWorkerId, + inputOrder + })) + )) + }; + const persistedInputs = structuredClone(store.run.persistedDedupInputs); + await Promise.all(store.run.persistedWorkers.map((worker) => rm(worker.promptPath, { force: true }))); + const replacementExecutor = new FakeExecutor(); + const replacement = new DeepScanCoordinator({ + run: store.run, + store, + executor: replacementExecutor, + pluginRoot: fixture.pluginRoot, + clock: immediateClock + }); + replacement.start(); + + const resumed = await replacement.wait(undefined, 5_000); + assert.equal(resumed?.status, "succeeded", resumed?.error); + const firstReducer = store.run.persistedWorkers.find((worker) => ( + worker.kind === "dedup" && worker.promptPath.includes("dedup-0001") + )); + assert.ok(firstReducer); + const firstResult = JSON.parse(await readFile(firstReducer.resultManifestPath, "utf8")); + assert.equal(firstResult.findings[0]?.rootCause.summary, "first reducer evidence"); + const lastReducer = store.run.persistedWorkers.filter((worker) => ( + worker.kind === "dedup" && worker.status === "succeeded" + )).at(-1); + const latestResult = JSON.parse(await readFile(lastReducer.resultManifestPath, "utf8")); + assert.equal(latestResult.findings[0]?.rootCause.summary, "final reducer evidence"); + const manifest = JSON.parse(await readFile(resumed.manifestPath, "utf8")); + assert.deepEqual(manifest.findings, latestResult.findings); + assert.deepEqual(store.run.persistedDedupInputs, persistedInputs); + assert.equal(replacementExecutor.calls, 0); + } + + await testPausedDiscoverySurvivesCoordinatorRestart(false); + await testPausedDiscoverySurvivesCoordinatorRestart(); + await testResumedDiscoveryDeadlineUsesPersistedCreationTime(); + await testResumedDiscoveryDeadlineUsesPersistedCreationTime(true); + await testResumedDiscoveryDeadlineUsesPersistedCreationTime(false, 2.5); + await testResumedDiscoveryDeadlineUsesPersistedCreationTime(true, 96); + await testResumedManifestPreservesCompletedReducer(false, false); + await testResumedManifestPreservesCompletedReducer(); + await testResumedManifestPreservesCompletedReducer(true); + await testResumeUsesHistoricalCandidateSnapshotForEachReducer(); +} diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_coordinator.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_coordinator.mjs index 0625bbbb1..474cf9fd5 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_coordinator.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_coordinator.mjs @@ -5,6 +5,7 @@ import { tmpdir } from "node:os"; import path from "node:path"; import { build } from "esbuild"; import { testDeepScanPublication } from "./deep_scan_publication_cases.mjs"; +import { testDeepScanResumeCases } from "./deep_scan_resume_cases.mjs"; const bundle = await build({ bundle: true, @@ -2824,353 +2825,6 @@ async function testJoinAndOrphanRules() { assert.equal(failures, 0); } -async function testPausedDiscoverySurvivesCoordinatorRestart() { - const fixture = await fixtureRun({ - workers: 1, - subagents: 0, - stopAfterNoNew: 2, - maxDiscoveryRuns: 2 - }); - const handoffClaimToken = randomUUID(); - const store = new FakeStore({ - ...fixture.run, - phase: "setup", - coordinatorGeneration: 2, - updatedAt: "2026-08-03T13:33:08Z" - }); - const originalExecutor = new FakeExecutor({ blockDiscoveryAfterCalls: 1 }); - const original = new DeepScanCoordinator({ - run: store.run, - store, - executor: originalExecutor, - pluginRoot: fixture.pluginRoot, - clock: immediateClock, - handoffClaimToken - }); - original.start(); - await eventually(() => ( - [...store.workers.values()].some((worker) => ( - worker.kind === "discovery" && worker.status === "succeeded" - )) - && originalExecutor.discoveryCalls >= 2 - )); - const accepted = [...store.workers.values()].find((worker) => ( - worker.kind === "discovery" && worker.status === "succeeded" - )); - assert.ok(accepted); - - // The waiter detached earlier; an app update now removes its MCP process - // without canceling or finalizing the persisted scan. - original.cancel("mcp server process restarted"); - await eventually(() => originalExecutor.runningDiscovery === 0); - const persistedWorkers = [...store.workers.values()].map((worker) => structuredClone(worker)); - const independentReviews = { - completed: persistedWorkers.filter((worker) => ( - worker.kind === "discovery" && worker.status === "succeeded" - )).length, - active: persistedWorkers.filter((worker) => ( - worker.kind === "discovery" && worker.status === "running" - )).length, - consolidating: persistedWorkers.some((worker) => ( - worker.kind === "dedup" && worker.status === "running" - )) - }; - assert.deepEqual(independentReviews, { completed: 1, active: 0, consolidating: false }); - assert.equal(store.run.status, "running"); - assert.equal(store.run.phase, "discovery"); - assert.equal(store.finishCalls.length, 0); - assert.equal(store.failCalls, 0); - assert.equal(store.run.manifestPath, undefined); - await assert.rejects( - readFile(path.join( - fixture.run.scanDir, - "artifacts", - "deep_discovery", - "coordinator-manifest.json" - )), - { code: "ENOENT" } - ); - - store.run = { - ...store.run, - dispatchedCount: 1, - persistedWorkers - }; - const continuationClaims = []; - store.claimCoordinator = async (input) => { - continuationClaims.push(structuredClone(input)); - assert.equal(input.handoffClaimToken, handoffClaimToken); - store.run = { - ...store.run, - coordinatorGeneration: 3, - updatedAt: new Date().toISOString() - }; - return { acquired: true, run: structuredClone(store.run) }; - }; - store.heartbeatCoordinator = async () => structuredClone(store.run); - const replacementExecutor = new FakeExecutor({ dedupNewFindings: [0] }); - const acceptedResult = await readFile(accepted.resultManifestPath, "utf8"); - const resumed = await startOrJoinDeepScanCoordinator({ - begin: { run: structuredClone(store.run), shouldStart: false }, - registry: new DeepScanCoordinatorRegistry(), - options: { - store, - executor: replacementExecutor, - pluginRoot: fixture.pluginRoot, - clock: immediateClock, - threadId: "track-c-owning-thread", - handoffClaimToken - } - }); - const terminal = await resumed.coordinator.wait(undefined, 5_000); - - assert.equal(continuationClaims.length, 1); - assert.equal(continuationClaims[0].handoffClaimToken, handoffClaimToken); - assert.equal(terminal?.status, "succeeded"); - assert.equal(store.failCalls, 0); - assert.equal(replacementExecutor.logicalDiscoveryWorkers.size, 1); - assert.equal( - replacementExecutor.logicalDiscoveryWorkers.has(await workerIdFromPrompt(accepted.promptPath)), - false - ); - assert.equal(store.dedupClaims.length, 1); - assert.equal(store.dedupClaims[0].workerIds.includes(accepted.id), true); - const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); - assert.equal(manifest.scan.scanId, fixture.run.scanId); - assert.equal(store.dedupClaims[0].workerIds.length, 2); - assert.equal(await readFile(accepted.resultManifestPath, "utf8"), acceptedResult); -} - -async function testResumedDiscoveryDeadlineUsesPersistedCreationTime( - alreadyExpired = false, - maxTimeHours -) { - const fixture = await fixtureRun({ - workers: 1, - subagents: 0, - stopAfterNoNew: 99, - maxDiscoveryRuns: 8, - ...(maxTimeHours === undefined ? {} : { maxTimeHours }) - }); - const discoveryTimeoutMs = (maxTimeHours ?? 96) * 60 * 60 * 1_000; - let currentTime = immediateClock.now(); - const clock = { - now: () => currentTime, - sleep: immediateClock.sleep - }; - const createdAt = new Date(currentTime - discoveryTimeoutMs + 30_000).toISOString(); - const store = new FakeStore({ - ...fixture.run, - createdAt, - phase: "setup", - coordinatorGeneration: 2 - }); - const originalExecutor = new FakeExecutor({ - blockDiscoveryAfterCalls: 1, - discoveryCandidateId: "candidate-1" - }); - const original = new DeepScanCoordinator({ - run: store.run, - store, - executor: originalExecutor, - pluginRoot: fixture.pluginRoot, - clock - }); - original.start(); - await eventually(() => ( - originalExecutor.discoveryCalls === 2 - && originalExecutor.runningDiscovery === 1 - && [...store.workers.values()].some((worker) => ( - worker.kind === "discovery" && worker.status === "succeeded" - )) - )); - - original.cancel("mcp server process restarted"); - await eventually(() => originalExecutor.runningDiscovery === 0); - assert.equal(store.run.status, "running"); - assert.equal(store.run.createdAt, createdAt); - currentTime = Date.parse(createdAt) + discoveryTimeoutMs + (alreadyExpired ? 1_000 : -1_000); - store.run = { - ...store.run, - persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)) - }; - - const resumedExecutor = new FakeExecutor({ - blockDiscovery: true, - canonicalCandidateId: "candidate-1", - dedupNewFindings: [1] - }); - const resumed = new DeepScanCoordinator({ - run: store.run, - store, - executor: resumedExecutor, - pluginRoot: fixture.pluginRoot, - clock - }); - resumed.start(); - if (!alreadyExpired) await resumedExecutor.discoveryStarted; - - const terminal = await resumed.wait(undefined, 5_000); - assert.equal(terminal?.status, "succeeded"); - assert.equal(terminal?.terminalReason, "capped"); - assert.equal(store.failCalls, 0); - assert.equal(resumedExecutor.discoveryCalls, alreadyExpired ? 0 : 1); - assert.equal(resumedExecutor.dedupCalls, 1); - assert.equal(resumedExecutor.runningDiscovery, 0); - - const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); - assert.equal(store.run.config.maxTimeHours, maxTimeHours); - assert.equal(store.dedupClaims[0].workerIds.length, 1); - assert.equal([...store.workers.values()].some((worker) => worker.status === "canceled"), true); - assert.deepEqual(manifest.findings.map((finding) => finding.provenance.candidateId), ["candidate-1"]); -} - -async function testResumedManifestPreservesCompletedReducer(includeUnstartedReducer = false) { - const fixture = await fixtureRun({ - workers: 2, - subagents: 0, - stopAfterNoNew: 2, - maxDiscoveryRuns: 3 - }); - const store = new FakeStore({ ...fixture.run, phase: "setup" }); - store.blockDedupCommitResponse = true; - const original = new DeepScanCoordinator({ - run: store.run, - store, - executor: new FakeExecutor({ - dedupNewFindings: [0], - blockDiscoveryAfterCalls: 2 - }), - pluginRoot: fixture.pluginRoot, - clock: immediateClock - }); - original.start(); - await store.dedupCommitPersisted.promise; - original.cancel("mcp server process restarted"); - store.releaseDedupCommitResponse(); - await original.settled(); - await eventually(() => [...store.workers.values()].every((worker) => ( - worker.status !== "queued" && worker.status !== "running" - ))); - - let unstartedReducer; - if (includeUnstartedReducer) { - const artifactDir = path.join( - fixture.run.scanDir, - "artifacts", - "deep_discovery", - "dedup", - "dedup-unstarted", - "output" - ); - await mkdir(artifactDir, { recursive: true }); - const promptPath = path.join(path.dirname(artifactDir), "prompt.md"); - await writeFile(promptPath, "Reducer claimed before coordinator restart.\n"); - unstartedReducer = { - id: randomUUID(), - kind: "dedup", - status: "canceled", - promptPath, - artifactDir, - attempt: 0, - mergeState: "none" - }; - store.workers.set(unstartedReducer.id, unstartedReducer); - } - - store.run = { - ...store.run, - status: "running", - phase: "discovery", - persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)), - persistedDedupInputs: store.dedupClaims.flatMap((claim) => ( - claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ - dedupWorkerId: claim.id, - discoveryWorkerId, - inputOrder - })) - )) - }; - const replacement = new DeepScanCoordinator({ - run: store.run, - store, - executor: new FakeExecutor(), - pluginRoot: fixture.pluginRoot, - clock: immediateClock - }); - replacement.start(); - - const terminal = await replacement.wait(undefined, 5_000); - assert.equal(terminal?.status, "succeeded"); - const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); - assert.equal(manifest.scan.scanId, fixture.run.scanId); - assert.equal(store.dedupCommits.length, 1); - assert.equal(store.dedupClaims.length, 1); - assert.equal(store.run.persistedWorkers.some((worker) => worker.id === unstartedReducer?.id), - includeUnstartedReducer); -} - -async function testResumeUsesHistoricalCandidateSnapshotForEachReducer(legacyLayout = false) { - const fixture = await fixtureRun({ - workers: 3, - subagents: 0, - stopAfterNoNew: 10, - maxDiscoveryRuns: 5 - }); - const store = new FakeStore(fixture.run); - const original = new DeepScanCoordinator({ - run: fixture.run, - store, - executor: new FakeExecutor({ - dedupNewFindings: [1, 0], - dedupEvidenceByCall: ["first reducer evidence", "final reducer evidence"] - }), - pluginRoot: fixture.pluginRoot, - clock: immediateClock - }); - original.start(); - assert.equal((await original.wait(undefined, 5_000))?.status, "succeeded"); - assert.equal(store.dedupClaims.length >= 2, true); - - store.run = { - ...store.run, - status: "running", - phase: "discovery", - terminalReason: undefined, - manifestPath: undefined, - persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker)), - persistedDedupInputs: store.dedupClaims.flatMap((claim) => ( - claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ - dedupWorkerId: claim.id, - discoveryWorkerId, - inputOrder - })) - )) - }; - const replacement = new DeepScanCoordinator({ - run: store.run, - store, - executor: new FakeExecutor(), - pluginRoot: fixture.pluginRoot, - clock: immediateClock - }); - replacement.start(); - - const resumed = await replacement.wait(undefined, 5_000); - assert.equal(resumed?.status, "succeeded", resumed?.error); - const firstReducer = store.run.persistedWorkers.find((worker) => ( - worker.kind === "dedup" && worker.promptPath.includes("dedup-0001") - )); - assert.ok(firstReducer); - const firstResult = JSON.parse(await readFile(firstReducer.resultManifestPath, "utf8")); - assert.equal(firstResult.findings[0]?.rootCause.summary, "first reducer evidence"); - const lastReducer = store.run.persistedWorkers.filter((worker) => ( - worker.kind === "dedup" && worker.status === "succeeded" - )).at(-1); - const latestResult = JSON.parse(await readFile(lastReducer.resultManifestPath, "utf8")); - assert.equal(latestResult.findings[0]?.rootCause.summary, "final reducer evidence"); -} - async function testPersistedErrorLimitStopsBeforeRescheduling() { const fixture = await fixtureRun({ workers: 1, @@ -4032,15 +3686,10 @@ try { await testCoordinatorHeartbeatsContinueDuringBlockedOwnershipRead(); await testRemoteObserverRetriesTransientPersistenceFailures(); await testJoinAndOrphanRules(); - await testPausedDiscoverySurvivesCoordinatorRestart(); - await testResumedDiscoveryDeadlineUsesPersistedCreationTime(); - await testResumedDiscoveryDeadlineUsesPersistedCreationTime(true); - await testResumedDiscoveryDeadlineUsesPersistedCreationTime(false, 2.5); - await testResumedDiscoveryDeadlineUsesPersistedCreationTime(true, 96); - await testResumedManifestPreservesCompletedReducer(); - await testResumedManifestPreservesCompletedReducer(true); - await testResumeUsesHistoricalCandidateSnapshotForEachReducer(); - await testResumeUsesHistoricalCandidateSnapshotForEachReducer(true); + await testDeepScanResumeCases({ + fixtureRun, FakeStore, FakeExecutor, DeepScanCoordinator, DeepScanCoordinatorRegistry, + startOrJoinDeepScanCoordinator, immediateClock, eventually, promptContext, workerIdFromPrompt + }); await testPersistedErrorLimitStopsBeforeRescheduling(); await testPersistedReducerErrorLimitStopsBeforeRescheduling(); } finally {