From 6216186eb1ca4af9db5d28c00d0d35ccf17c83c3 Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 3 Oct 2026 19:14:41 +0000 Subject: [PATCH 1/4] refactor(plugin): simplify artifact storage boundaries Artifact storage and server entrypoints repeat context and file-handling work. Consolidate that work at its actual owners while preserving storage boundaries. --- .../mcp-app/artifact-writer-main.ts | 60 ++++---- plugins/codex-security/mcp-app/server.ts | 3 +- .../mcp-app/src/artifact-context.ts | 53 +------ .../mcp-app/src/artifact-inventory.ts | 18 --- .../codex-security/mcp-app/src/artifact-io.ts | 32 ++--- .../mcp-app/src/artifact-scan-draft.ts | 62 ++++---- .../src/server/compact-artifact-tools.ts | 3 +- .../mcp-app/src/server/handoff-tools.ts | 6 +- plugins/codex-security/mcp-app/src/types.ts | 48 ------- .../reducer-paging/controlled-code-mode.mjs | 8 +- .../deep-reducer-paging-fixture.mjs | 9 +- .../reducer-paging/deep-reducer-paging.mjs | 9 +- .../tests/support/source-references.mjs | 7 + .../tests/support/temporary-directories.d.mts | 4 + .../tests/support/temporary-directories.mjs | 24 ++++ .../tests/test_artifact_attack_path.mjs | 37 +---- .../tests/test_artifact_deep_reducer.mjs | 25 +--- .../test_artifact_deep_reducer_pages.mjs | 15 +- .../mcp-app/tests/test_artifact_discovery.mjs | 22 +-- .../tests/test_artifact_foundation.mjs | 76 ++++++---- .../mcp-app/tests/test_artifact_inventory.mjs | 52 +------ .../tests/test_artifact_scan_draft.mjs | 132 ++++++++---------- .../mcp-app/tests/test_artifact_storage.mjs | 6 +- .../tests/test_artifact_validation_phase.mjs | 20 +-- .../tests/test_compact_artifact_server.mjs | 36 ++--- 25 files changed, 281 insertions(+), 486 deletions(-) create mode 100644 plugins/codex-security/mcp-app/tests/support/source-references.mjs create mode 100644 plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts create mode 100644 plugins/codex-security/mcp-app/tests/support/temporary-directories.mjs diff --git a/plugins/codex-security/mcp-app/artifact-writer-main.ts b/plugins/codex-security/mcp-app/artifact-writer-main.ts index 505fc29b14..be11c3df20 100644 --- a/plugins/codex-security/mcp-app/artifact-writer-main.ts +++ b/plugins/codex-security/mcp-app/artifact-writer-main.ts @@ -1,11 +1,14 @@ import { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import { - createWorkerArtifactContext, + canonicalDirectory, + defined, + type ArtifactContext, type DeepReducerContext, } from "./src/artifact-context.js"; import { CODEX_SANDBOX_STATE_META_CAPABILITY } from "./src/deep-scan/parent-sandbox.js"; import { registerCompactWorkerArtifactTools } from "./src/server/compact-artifact-tools.js"; import { MCP_APP_VERSION } from "./src/version.js"; +import { isRecord } from "./src/record.js"; /** Build the narrow worker-only MCP from coordinator-inherited state. */ export async function createCodexSecurityArtifactWriterServer( @@ -30,24 +33,35 @@ export async function createCodexSecurityArtifactWriterServer( ); } - const context = await createWorkerArtifactContext({ - root, - repoRoot, + /** + * Bind a lightweight worker to host-supplied state, never model-supplied paths. + */ + const scanId = environment.CODEX_SECURITY_SCAN_ID + ? environment.CODEX_SECURITY_SCAN_ID + : undefined; + const scope = environment.CODEX_SECURITY_SCOPE + ? environment.CODEX_SECURITY_SCOPE + : undefined; + const pluginRoot = environment.CODEX_SECURITY_PLUGIN_ROOT + ? environment.CODEX_SECURITY_PLUGIN_ROOT + : undefined; + const pythonCommand = environment.CODEX_SECURITY_PYTHON_COMMAND + ? environment.CODEX_SECURITY_PYTHON_COMMAND + : undefined; + // Preserve the asynchronous context boundary before server construction. + const context: ArtifactContext = await (async () => ({ + root: await canonicalDirectory(root, "Codex Security worker artifact root"), + repoRoot: await canonicalDirectory( + repoRoot, + "Codex Security worker target root", + ), layout, - ...(environment.CODEX_SECURITY_SCAN_ID - ? { scanId: environment.CODEX_SECURITY_SCAN_ID } - : {}), - ...(environment.CODEX_SECURITY_SCOPE - ? { scope: environment.CODEX_SECURITY_SCOPE } - : {}), - ...(environment.CODEX_SECURITY_PLUGIN_ROOT - ? { pluginRoot: environment.CODEX_SECURITY_PLUGIN_ROOT } - : {}), - ...(environment.CODEX_SECURITY_PYTHON_COMMAND - ? { pythonCommand: environment.CODEX_SECURITY_PYTHON_COMMAND } - : {}), - ...(deepReducer ? { deepReducer } : {}), - }); + ...defined("scanId", scanId), + ...defined("scope", scope), + ...defined("pluginRoot", pluginRoot), + ...defined("pythonCommand", pythonCommand), + ...defined("deepReducer", deepReducer), + }))(); const server = new McpServer( { name: "codex-security-artifacts", version: MCP_APP_VERSION }, { @@ -84,15 +98,13 @@ function parseReducerContext(value: string): DeepReducerContext { ); } if ( - !parsed || - typeof parsed !== "object" || - Array.isArray(parsed) || - typeof (parsed as Record).scanRoot !== "string" || - !Array.isArray((parsed as Record).claimedWorkers) + !isRecord(parsed) || + typeof parsed.scanRoot !== "string" || + !Array.isArray(parsed.claimedWorkers) ) { throw new Error( "The coordinator-bound Deep reducer context is incomplete.", ); } - return parsed as DeepReducerContext; + return parsed as unknown as DeepReducerContext; } diff --git a/plugins/codex-security/mcp-app/server.ts b/plugins/codex-security/mcp-app/server.ts index 230e471024..48ee45d44c 100644 --- a/plugins/codex-security/mcp-app/server.ts +++ b/plugins/codex-security/mcp-app/server.ts @@ -11,7 +11,6 @@ import { missingPythonHelperMessage, resolvePythonCommand, } from "./src/python_command.js"; -import type { ScanResults } from "./src/types.js"; import { MCP_APP_VERSION } from "./src/version.js"; import { handoffClaimTokenSchema, @@ -129,7 +128,7 @@ async function scanRoot(): Promise { interface WorkspaceState extends JsonObject { id: string; - results?: ScanResults & JsonObject; + results?: JsonObject; setup: { submitted: boolean; }; diff --git a/plugins/codex-security/mcp-app/src/artifact-context.ts b/plugins/codex-security/mcp-app/src/artifact-context.ts index 688e2b4a7b..22b87eacee 100644 --- a/plugins/codex-security/mcp-app/src/artifact-context.ts +++ b/plugins/codex-security/mcp-app/src/artifact-context.ts @@ -22,13 +22,6 @@ export interface ScanArtifactContextOptions { pythonCommand?: string; } -export interface WorkerArtifactContextInput extends Omit< - ArtifactContext, - "layout" -> { - layout?: "worker" | "reducer"; -} - /** * Resolve parent artifacts from their authoritative, persisted workbench scan. */ @@ -108,45 +101,6 @@ export async function createScanArtifactContext( }; } -/** - * Bind a lightweight worker to host-supplied state, never model-supplied paths. - */ -export async function createWorkerArtifactContext( - input: WorkerArtifactContextInput, -): Promise { - const layout = input.layout ?? "worker"; - if (layout !== "worker" && layout !== "reducer") { - throw new Error("Codex Security worker artifact layout is invalid."); - } - const context: ArtifactContext = { - root: await canonicalDirectory( - input.root, - "Codex Security worker artifact root", - ), - repoRoot: await canonicalDirectory( - input.repoRoot, - "Codex Security worker target root", - ), - layout, - ...defined("scanId", input.scanId), - ...defined("scope", input.scope), - ...defined("pluginRoot", input.pluginRoot), - ...defined("pythonCommand", input.pythonCommand), - ...defined("targetContract", input.targetContract), - ...defined("targetRevision", input.targetRevision), - ...defined("handoffClaimToken", input.handoffClaimToken), - ...defined("status", input.status), - ...defined("mode", input.mode), - ...defined("deepReducer", input.deepReducer), - }; - if (context.deepReducer && layout !== "reducer") { - throw new Error( - "Codex Security reducer state requires a reducer-bound context.", - ); - } - return context; -} - function scanRecord( result: Record, scanId: string, @@ -162,7 +116,7 @@ function scanRecord( return scan; } -async function canonicalDirectory( +export async function canonicalDirectory( value: string, label: string, ): Promise { @@ -170,8 +124,7 @@ async function canonicalDirectory( throw new Error(label + " must be an absolute directory."); } const requested = resolve(value); - const metadata = await fs.lstat(requested).catch(() => undefined); - if (!metadata || metadata.isSymbolicLink() || !metadata.isDirectory()) { + if (!(await fs.lstat(requested).catch(() => undefined))?.isDirectory()) { throw new Error(label + " is not a safe regular directory."); } try { @@ -191,7 +144,7 @@ function optionalString(value: unknown): string | undefined { return typeof value === "string" && value.trim() ? value : undefined; } -function defined( +export function defined( key: Key, value: Value | undefined, ): Partial> { diff --git a/plugins/codex-security/mcp-app/src/artifact-inventory.ts b/plugins/codex-security/mcp-app/src/artifact-inventory.ts index 007720791f..f174c63758 100644 --- a/plugins/codex-security/mcp-app/src/artifact-inventory.ts +++ b/plugins/codex-security/mcp-app/src/artifact-inventory.ts @@ -44,30 +44,12 @@ export const prepareReviewItemsInputSchema = loadArtifactZodSchema( "prepareInput", ) as z.ZodType<{ scanId: string; handoffClaimToken?: string }>; -export const prepareReviewItemsOutputSchema = loadArtifactZodSchema( - documents, - reviewItemsSchema.$id, - "prepareOutput", -) as z.ZodType; - export const reviewItemsReaderInputSchema = loadArtifactZodSchema( documents, reviewItemsSchema.$id, "reviewItemsInput", ) as z.ZodType<{ scanId: string; handoffClaimToken?: string } & ArtifactPage>; -export const reviewItemsWorkerReaderInputSchema = loadArtifactZodSchema( - documents, - reviewItemsSchema.$id, - "reviewItemsWorkerInput", -) as z.ZodType; - -export const reviewItemsReaderOutputSchema = loadArtifactZodSchema( - documents, - reviewItemsSchema.$id, - "reviewItemsOutput", -) as z.ZodType; - const reviewItemSchema = loadArtifactZodSchema( documents, reviewItemsSchema.$id, diff --git a/plugins/codex-security/mcp-app/src/artifact-io.ts b/plugins/codex-security/mcp-app/src/artifact-io.ts index 2d6e6a6d46..d8e19a61db 100644 --- a/plugins/codex-security/mcp-app/src/artifact-io.ts +++ b/plugins/codex-security/mcp-app/src/artifact-io.ts @@ -1,7 +1,8 @@ +import { canonicalDirectory } from "./artifact-context.js"; import { isRecord } from "./record.js"; import { randomUUID } from "node:crypto"; import { constants as fsConstants, promises as fs } from "node:fs"; -import { dirname, isAbsolute, join, resolve, sep } from "node:path"; +import { dirname, isAbsolute, join, sep } from "node:path"; export interface DeepReducerWorkerContext { id: string; @@ -101,10 +102,7 @@ async function artifactSourcePath( throw new Error(label + ": the requested artifact is unavailable."); } const isLast = index === components.length - 1; - if ( - metadata.isSymbolicLink() || - (isLast ? !metadata.isFile() : !metadata.isDirectory()) - ) { + if (isLast ? !metadata.isFile() : !metadata.isDirectory()) { throw new Error( label + ": the requested artifact is not a safe regular file.", ); @@ -226,7 +224,7 @@ export async function artifactDestination( } metadata = await inspectOptionalPath(directory, label); } - if (!metadata || metadata.isSymbolicLink() || !metadata.isDirectory()) { + if (!metadata || !metadata.isDirectory()) { throw new Error( label + ": destination directory is not a regular directory.", ); @@ -241,8 +239,7 @@ export async function artifactDestination( if (!destination.startsWith(root + sep)) { throw new Error(label + ": destination escaped its bound context."); } - const metadata = await inspectOptionalPath(destination, label); - if (metadata && (metadata.isSymbolicLink() || !metadata.isFile())) { + if ((await inspectOptionalPath(destination, label))?.isFile() === false) { throw new Error(label + ": destination is not a regular file."); } return destination; @@ -338,27 +335,16 @@ function validateArtifactComponents( } } -export async function requireArtifactRoot( +export function requireArtifactRoot( artifactRoot: string, label: string, ): Promise { if (!artifactRoot || !isAbsolute(artifactRoot)) { - throw new Error( - label + ": artifact context must have an absolute bound root.", + return Promise.reject( + new Error(label + ": artifact context must have an absolute bound root."), ); } - const requested = resolve(artifactRoot); - const metadata = await fs.lstat(requested).catch(() => undefined); - if (!metadata || metadata.isSymbolicLink() || !metadata.isDirectory()) { - throw new Error( - label + ": artifact context is not a safe regular directory.", - ); - } - try { - return await fs.realpath(requested); - } catch { - throw new Error(label + ": artifact context cannot be resolved."); - } + return canonicalDirectory(artifactRoot, label + ": artifact context"); } async function inspectOptionalPath( diff --git a/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts b/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts index 9e5793e981..d0e661a410 100644 --- a/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts +++ b/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts @@ -1,3 +1,4 @@ +import type { JsonObject } from "./types.js"; import { isRecord as isObject } from "./record.js"; import { createHash, randomUUID } from "node:crypto"; import { promises as fs } from "node:fs"; @@ -21,8 +22,6 @@ import { } from "./artifact-schema-loader.js"; import { saveThreatModelDocument } from "./threat-model-document.js"; -type JsonObject = Record; - export interface ScanDraftInput { scanId: string; complete?: boolean; @@ -1126,10 +1125,7 @@ async function readSavedCheckpoints( let checkpointRoot = join(context.root, "checkpoints"); const checkpointRootMetadata = await lstatIfExists(checkpointRoot); if (checkpointRootMetadata === undefined) return []; - if ( - checkpointRootMetadata.isSymbolicLink() || - !checkpointRootMetadata.isDirectory() - ) { + if (!checkpointRootMetadata.isDirectory()) { throw new Error( `scan checkpoint: ${kind} checkpoint set is not a safe directory.`, ); @@ -1165,7 +1161,7 @@ async function readSavedCheckpoints( for (const entry of entries) { const checkpointPath = join(checkpointRoot, entry.name); const checkpointMetadata = await fs.lstat(checkpointPath); - if (checkpointMetadata.isSymbolicLink() || !checkpointMetadata.isFile()) { + if (!checkpointMetadata.isFile()) { throw new Error( `scan checkpoint: ${kind} checkpoint is not a safe file.`, ); @@ -1326,7 +1322,7 @@ async function readArchivedWorkerCheckpoints( const attemptsRoot = join(workerRoot, "attempts"); const attemptsMetadata = await lstatIfExists(attemptsRoot); if (attemptsMetadata === undefined) return []; - if (attemptsMetadata.isSymbolicLink() || !attemptsMetadata.isDirectory()) { + if (!attemptsMetadata.isDirectory()) { throw new Error( "scan checkpoint: archived attempts are not a safe directory.", ); @@ -1345,10 +1341,11 @@ async function readArchivedWorkerCheckpoints( const attempts = ( await fs.readdir(canonicalAttemptsRoot, { withFileTypes: true }) ) - .filter((entry) => entry.isDirectory() && !entry.isSymbolicLink()) + .filter((entry) => entry.isDirectory()) .sort( (left, right) => - archivedAttemptNumber(right.name) - archivedAttemptNumber(left.name) || + Number(/^attempt-(\d+)$/.exec(right.name)?.[1] ?? -1) - + Number(/^attempt-(\d+)$/.exec(left.name)?.[1] ?? -1) || right.name.localeCompare(left.name), ); for (const attempt of attempts) { @@ -1382,7 +1379,7 @@ async function readArchivedWorkerCheckpoints( join(attemptRoot, "result.json"), ); if (resultMetadata !== undefined) { - if (resultMetadata.isSymbolicLink() || !resultMetadata.isFile()) { + if (!resultMetadata.isFile()) { throw new Error("scan checkpoint: archived result is not a safe file."); } const saved = await readArtifactTextWithMetadata( @@ -1436,11 +1433,6 @@ async function readArchivedWorkerCheckpoints( return archived; } -function archivedAttemptNumber(name: string): number { - const match = /^attempt-(\d+)$/.exec(name); - return match ? Number(match[1]) : -1; -} - async function lstatIfExists( path: string, ): Promise> | undefined> { @@ -2241,29 +2233,27 @@ function buildScope( ): JsonObject { const includePaths = trustedScope.requiredIncludePaths; const excludePaths = trustedScope.requiredExcludePaths; - const resolvedIncludePaths = - includePaths === undefined - ? [ - typeof trustedScope.requestedPath === "string" - ? trustedScope.requestedPath - : (context.scope ?? "."), - ] - : requireTextArray( - includePaths, - "scan draft: authoritative included scope", - ); - const resolvedExcludePaths = - excludePaths === undefined - ? [] - : requireTextArray( - excludePaths, - "scan draft: authoritative excluded scope", - ); return { ...semanticScope, - includePaths: resolvedIncludePaths, - excludePaths: resolvedExcludePaths, + includePaths: + includePaths === undefined + ? [ + typeof trustedScope.requestedPath === "string" + ? trustedScope.requestedPath + : (context.scope ?? "."), + ] + : requireTextArray( + includePaths, + "scan draft: authoritative included scope", + ), + excludePaths: + excludePaths === undefined + ? [] + : requireTextArray( + excludePaths, + "scan draft: authoritative excluded scope", + ), }; } diff --git a/plugins/codex-security/mcp-app/src/server/compact-artifact-tools.ts b/plugins/codex-security/mcp-app/src/server/compact-artifact-tools.ts index 3dacd96dfa..121359c1c7 100644 --- a/plugins/codex-security/mcp-app/src/server/compact-artifact-tools.ts +++ b/plugins/codex-security/mcp-app/src/server/compact-artifact-tools.ts @@ -1,5 +1,6 @@ import type { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import type { ZodType } from "zod/v4"; +import type { JsonObject as JsonRecord } from "../types.js"; import { createScanArtifactContext, type ArtifactContext, @@ -52,8 +53,6 @@ import { type ArtifactLocation, } from "../artifact-storage.js"; -type JsonRecord = Record; - export interface CompactArtifactToolOptions { runWorkbench: RunArtifactWorkbench; pluginRoot: string; diff --git a/plugins/codex-security/mcp-app/src/server/handoff-tools.ts b/plugins/codex-security/mcp-app/src/server/handoff-tools.ts index 863850c814..0a68592d0f 100644 --- a/plugins/codex-security/mcp-app/src/server/handoff-tools.ts +++ b/plugins/codex-security/mcp-app/src/server/handoff-tools.ts @@ -1,13 +1,11 @@ import type { McpServer } from "@modelcontextprotocol/sdk/server/mcp.js"; import type { CallToolResult } from "@modelcontextprotocol/sdk/types.js"; import * as z from "zod/v4"; -import type { ScanResults } from "../types.js"; - -type JsonObject = Record; +import type { JsonObject } from "../types.js"; export interface HandoffWorkspaceState extends JsonObject { id: string; - results?: ScanResults & JsonObject; + results?: JsonObject; setup: { submitted: boolean; }; diff --git a/plugins/codex-security/mcp-app/src/types.ts b/plugins/codex-security/mcp-app/src/types.ts index 6808e8f49f..84f16550ad 100644 --- a/plugins/codex-security/mcp-app/src/types.ts +++ b/plugins/codex-security/mcp-app/src/types.ts @@ -1,49 +1 @@ -export type Mode = "diff" | "standard" | "deep"; -export type DiffTargetKind = "working_tree" | "commit" | "range"; export type JsonObject = Record; - -export interface DiffTarget { - baseRevision?: string; - contentDigest?: string; - headRevision?: string; - kind: DiffTargetKind; -} - -export interface ScanArtifacts { - coverage?: string; - findings?: string; - manifest?: string; - markdownReport?: string; - sarifReport?: string; - threatModel?: string; -} - -export interface ScanResults { - artifacts: ScanArtifacts; - canceledAt?: string; - diffTarget?: DiffTarget; - findingCount: number; - findings: JsonObject[]; - findingsTruncated?: boolean; - failureMessage?: string; - handoffClaimedAt?: string; - handoffClaimToken?: string; - handoffStatus: "delivered" | "pending"; - mode: Mode; - progress?: JsonObject; - remediationAvailable?: boolean; - remediationUnavailableReason?: string; - resultsRecoveryNeeded: boolean; - scanDir: string; - scanId: string; - scope?: string; - severityCounts?: Record; - targetPath: string; - targetRevision?: string; - targetSummary?: string | null; - threatModelAvailable: boolean; - threatModelPath?: string; - threatModelProvenance?: JsonObject; - updatedAt?: string; - userContext?: string; -} diff --git a/plugins/codex-security/mcp-app/tests/support/reducer-paging/controlled-code-mode.mjs b/plugins/codex-security/mcp-app/tests/support/reducer-paging/controlled-code-mode.mjs index 6ab9852802..7af70e079b 100644 --- a/plugins/codex-security/mcp-app/tests/support/reducer-paging/controlled-code-mode.mjs +++ b/plugins/codex-security/mcp-app/tests/support/reducer-paging/controlled-code-mode.mjs @@ -1,8 +1,8 @@ +import { temporaryDirectory } from "../temporary-directories.mjs"; import { spawn } from "node:child_process"; -import { mkdir, mkdtemp, rm } from "node:fs/promises"; +import { mkdir, rm } from "node:fs/promises"; import { createServer } from "node:http"; import { createRequire } from "node:module"; -import { tmpdir } from "node:os"; import path from "node:path"; function responseEvents(id, item) { @@ -41,9 +41,7 @@ export async function runControlledCodeMode({ "bin", "codex.js", ); - const fixtureRoot = await mkdtemp( - path.join(tmpdir(), "codex-controlled-ipc-"), - ); + const fixtureRoot = await temporaryDirectory("codex-controlled-ipc-"); const codexHome = path.join(fixtureRoot, "home"); const toolOutputs = new Map(); const serverErrors = []; diff --git a/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging-fixture.mjs b/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging-fixture.mjs index 293f183495..60a76de6c4 100644 --- a/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging-fixture.mjs +++ b/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging-fixture.mjs @@ -2,6 +2,7 @@ import assert from "node:assert/strict"; import { createHash } from "node:crypto"; import { mkdir, readFile, realpath, writeFile } from "node:fs/promises"; import path from "node:path"; +import { sourceReferences } from "../source-references.mjs"; // This is the existing code-mode transport ceiling, used only to size the eval. const IPC_FRAME_LIMIT_BYTES = 64 * 1024 * 1024; @@ -103,13 +104,7 @@ export async function createReducerPagingFixture(root) { workerId, result: { ...result, - findings: currentFindings.map((value, index) => ({ - ...value, - provenance: { - ...value.provenance, - sourceFindingIds: [`${workerId}:${index}`], - }, - })), + findings: currentFindings.map(sourceReferences({ id: workerId })), }, }, ], diff --git a/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging.mjs b/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging.mjs index 7d22a94852..ad9a501f95 100644 --- a/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging.mjs +++ b/plugins/codex-security/mcp-app/tests/support/reducer-paging/deep-reducer-paging.mjs @@ -54,9 +54,16 @@ export async function runReducerPagingEval({ const { renderDedupPrompt } = await import( pathToFileURL(promptModulePath).href ); + const claimedWorkerIds = fixture.context.deepReducer.claimedWorkers.map( + (worker) => worker.id, + ); const prompt = renderDedupPrompt({ reducerLabel: "paging-eval", - discoveries: fixture.context.deepReducer.claimedWorkers, + claimedWorkerIds, + }); + assert.deepEqual(JSON.parse(prompt.match(/```json\n([\s\S]*?)\n```/)[1]), { + reducerLabel: "paging-eval", + claimedWorkerIds, }); const mcpServers = { cs_artifacts: { diff --git a/plugins/codex-security/mcp-app/tests/support/source-references.mjs b/plugins/codex-security/mcp-app/tests/support/source-references.mjs new file mode 100644 index 0000000000..037fab5839 --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/support/source-references.mjs @@ -0,0 +1,7 @@ +export const sourceReferences = (worker) => (finding, index) => ({ + ...finding, + provenance: { + ...finding.provenance, + sourceFindingIds: [`${worker.id}:${index}`], + }, +}); diff --git a/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts b/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts new file mode 100644 index 0000000000..20e0e9a407 --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts @@ -0,0 +1,4 @@ +export function temporaryDirectory( + prefix: string, + canonicalize?: boolean, +): Promise; diff --git a/plugins/codex-security/mcp-app/tests/support/temporary-directories.mjs b/plugins/codex-security/mcp-app/tests/support/temporary-directories.mjs new file mode 100644 index 0000000000..89d977306d --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/support/temporary-directories.mjs @@ -0,0 +1,24 @@ +import { mkdtemp, realpath, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; + +export function createTemporaryDirectories(canonicalize = false) { + const roots = []; + return { + async create(prefix) { + const root = await temporaryDirectory(prefix, canonicalize); + roots.push(root); + return root; + }, + cleanup() { + return Promise.all( + roots.map((root) => rm(root, { recursive: true, force: true })), + ); + }, + }; +} + +export async function temporaryDirectory(prefix, canonicalize = false) { + const directory = await mkdtemp(join(tmpdir(), prefix)); + return canonicalize ? realpath(directory) : directory; +} diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_attack_path.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_attack_path.mjs index 6001838929..5546c69117 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_attack_path.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_attack_path.mjs @@ -1,13 +1,6 @@ +import { createTemporaryDirectories } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; -import { - mkdir, - mkdtemp, - readFile, - realpath, - rm, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, writeFile } from "node:fs/promises"; import path from "node:path"; import { importSource } from "./import-module.mjs"; @@ -19,7 +12,7 @@ const { ); const scanId = "11111111-1111-4111-8111-111111111111"; -const temporaryRoots = []; +const temporaryDirectories = createTemporaryDirectories(true); try { await testSchemaMatchesDocumentedAttackPathDecisions(); @@ -32,14 +25,7 @@ try { await testMalformedLedgerIsNotReplaced(); await testEmptyLedgerAcceptsAnEmptyBatch(); } finally { - await Promise.all( - temporaryRoots.map((root) => - rm(root, { - recursive: true, - force: true, - }), - ), - ); + await temporaryDirectories.cleanup(); } async function testSchemaMatchesDocumentedAttackPathDecisions() { @@ -293,15 +279,9 @@ async function testEmptyLedgerAcceptsAnEmptyBatch() { } async function createFixture(label, originalRows) { - const root = await realpath( - await mkdtemp( - path.join( - tmpdir(), - `codex-security-attack-path-${label.replace(/\s+/gu, "-")}-`, - ), - ), + const root = await temporaryDirectories.create( + `codex-security-attack-path-${label.replace(/\s+/gu, "-")}-`, ); - temporaryRoots.push(root); const ledgerPath = path.join( root, "artifacts", @@ -374,10 +354,7 @@ function attackPath(decision = "reportable") { async function readRows(fixture) { const content = await readFile(fixture.ledgerPath, "utf8"); - return content - .split(/\r?\n/u) - .filter(Boolean) - .map((line) => JSON.parse(line)); + return content.split(/\r?\n/u).filter(Boolean).map(JSON.parse); } async function assertUnchanged(fixture) { diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer.mjs index 3f7a177252..257dcef2f0 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer.mjs @@ -1,15 +1,8 @@ +import { sourceReferences } from "./support/source-references.mjs"; +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import { finding, scanId, workerDraft } from "./scan-draft-fixture.mjs"; import assert from "node:assert/strict"; -import { - mkdir, - mkdtemp, - readFile, - readdir, - realpath, - rm, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, readdir, rm, writeFile } from "node:fs/promises"; import path from "node:path"; import { importSource } from "./import-module.mjs"; @@ -56,9 +49,7 @@ assert.equal( false, ); -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "codex-security-deep-reducer-")), -); +const root = await temporaryDirectory("codex-security-deep-reducer-", true); try { const scanRoot = path.join(root, "scan"); const workersRoot = path.join( @@ -536,13 +527,7 @@ function withSourceRefs(worker) { const { coverage: _coverage, ...result } = worker.result; return { ...result, - findings: worker.result.findings.map((finding, index) => ({ - ...finding, - provenance: { - ...finding.provenance, - sourceFindingIds: [`${worker.id}:${index}`], - }, - })), + findings: worker.result.findings.map(sourceReferences(worker)), }; } diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer_pages.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer_pages.mjs index b5ce6f9227..e745102045 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer_pages.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_deep_reducer_pages.mjs @@ -1,13 +1,6 @@ +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; -import { - mkdir, - mkdtemp, - readFile, - realpath, - rm, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, rm, writeFile } from "node:fs/promises"; import path from "node:path"; import { importModule } from "./import-module.mjs"; @@ -45,9 +38,7 @@ assert.equal( ); const scanId = "7fc17317-9594-49e0-b06a-d72fd7e14bba"; -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "reducer-pages-")), -); +const root = await temporaryDirectory("reducer-pages-", true); try { const scanRoot = path.join(root, "scan"); const workerRoot = path.join( diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_discovery.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_discovery.mjs index 5938d163ad..8a582d4c6d 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_discovery.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_discovery.mjs @@ -1,16 +1,13 @@ +import { createTemporaryDirectories } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { copyFile, mkdir, - mkdtemp, readFile, readdir, - realpath, - rm, symlink, writeFile, } from "node:fs/promises"; -import { tmpdir } from "node:os"; import path from "node:path"; import { fileURLToPath } from "node:url"; import { build } from "esbuild"; @@ -86,9 +83,8 @@ assert.deepEqual(toolSchemas.$defs.workbenchListCandidatesInput.required, [ "scanId", ]); -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "security-artifact-discovery-")), -); +const rootDirectories = createTemporaryDirectories(true); +const root = await rootDirectories.create("security-artifact-discovery-"); const runtimePluginRoot = path.join(root, "plugin"); try { await build({ @@ -134,7 +130,7 @@ try { await verifyMalformedLedgerIsNotModified(root, repoRoot); await verifySymlinkRejection(root, repoRoot); } finally { - await rm(root, { recursive: true, force: true }); + await rootDirectories.cleanup(); } async function verifyInputSchema() { @@ -446,10 +442,7 @@ async function verifyReaderPreservesSharedPhaseRecords(context) { "candidate_ledger.jsonl", ); const original = await readFile(destination, "utf8"); - const rows = original - .trimEnd() - .split("\n") - .map((line) => JSON.parse(line)); + const rows = original.trimEnd().split("\n").map(JSON.parse); rows[0].validation = { disposition: "reportable", evidence: "The existing validation phase confirmed the affected code path.", @@ -458,10 +451,7 @@ async function verifyReaderPreservesSharedPhaseRecords(context) { disposition: "reportable", evidence: "The existing attack-path phase confirmed request reachability.", }; - await writeFile( - destination, - `${rows.map((row) => JSON.stringify(row)).join("\n")}\n`, - ); + await writeFile(destination, `${rows.map(JSON.stringify).join("\n")}\n`); const page = await listCodexSecurityCandidates({}, context); assert.deepEqual(page.rows, rows); diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_foundation.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_foundation.mjs index 5a23f5727f..0cd79c6609 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_foundation.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_foundation.mjs @@ -18,12 +18,25 @@ const compiled = await build({ entryPoints: [ new URL("../src/artifact-io.ts", import.meta.url).pathname, new URL("../src/artifact-context.ts", import.meta.url).pathname, + new URL("../artifact-writer-main.ts", import.meta.url).pathname, new URL("../src/artifact-schema-loader.ts", import.meta.url).pathname, ], format: "esm", outdir: "codex-security-artifact-foundation", platform: "node", write: false, + plugins: [ + { + name: "observe-worker-context", + setup(builder) { + builder.onLoad({ filter: /compact-artifact-tools\.ts$/ }, () => ({ + contents: + "export function registerCompactWorkerArtifactTools(server, context) { server.context = context; }", + loader: "js", + })); + }, + }, + ], }); const modules = new Map( compiled.outputFiles.map((file) => [ @@ -34,6 +47,7 @@ const modules = new Map( ); const io = await import(modules.get("artifact-io.js")); const contextApi = await import(modules.get("artifact-context.js")); +const writerApi = await import(modules.get("artifact-writer-main.js")); const schemas = await import(modules.get("artifact-schema-loader.js")); const fixture = await realpath( await mkdtemp(path.join(tmpdir(), "codex-security-artifact-foundation-")), @@ -272,12 +286,14 @@ async function testWorkerStandardLayout() { const root = path.join(fixture, "worker", "output"); const repoRoot = path.join(fixture, "repository"); await mkdir(root, { recursive: true }); - const context = await contextApi.createWorkerArtifactContext({ - root, - repoRoot, - scope: ".", - pluginRoot: "/fixture/plugin", - }); + const environment = { + CODEX_SECURITY_ARTIFACT_ROOT: root, + CODEX_SECURITY_REPO_ROOT: repoRoot, + CODEX_SECURITY_SCOPE: ".", + CODEX_SECURITY_PLUGIN_ROOT: "/fixture/plugin", + }; + const { context } = + await writerApi.createCodexSecurityArtifactWriterServer(environment); const inventory = await io.artifactDestination( context, ["artifacts", "02_discovery", "in_scope_files.txt"], @@ -310,24 +326,26 @@ async function testWorkerStandardLayout() { assert.equal(context.scope, "."); await assert.rejects( - contextApi.createWorkerArtifactContext({ - root, - repoRoot, - deepReducer: { scanRoot: root, claimedWorkers: [] }, + writerApi.createCodexSecurityArtifactWriterServer({ + ...environment, + CODEX_SECURITY_REDUCER_CONTEXT_JSON: JSON.stringify({ + scanRoot: root, + claimedWorkers: [], + }), }), - /reducer-bound context/, - ); - const reducer = await contextApi.createWorkerArtifactContext({ - root, - repoRoot, - layout: "reducer", - deepReducer: { - scanRoot: root, - claimedWorkers: [ - { id: "worker-1", resultPath: path.join(root, "worker-result.json") }, - ], - }, - }); + /coordinator-bound reducer context/, + ); + const { context: reducer } = + await writerApi.createCodexSecurityArtifactWriterServer({ + ...environment, + CODEX_SECURITY_ARTIFACT_LAYOUT: "reducer", + CODEX_SECURITY_REDUCER_CONTEXT_JSON: JSON.stringify({ + scanRoot: root, + claimedWorkers: [ + { id: "worker-1", resultPath: path.join(root, "worker-result.json") }, + ], + }), + }); assert.equal(reducer.layout, "reducer"); assert.equal(reducer.deepReducer.claimedWorkers[0].id, "worker-1"); } @@ -480,6 +498,8 @@ async function testUnsafeArtifacts() { const outside = path.join(fixture, "outside"); await mkdir(outside, { recursive: true }); + const outsideFile = path.join(outside, "candidate_ledger.jsonl"); + await writeFile(outsideFile, "outside remains unchanged\n"); const symlinkPath = path.join(context.root, "linked"); await symlink(outside, symlinkPath); await assert.rejects( @@ -491,12 +511,20 @@ async function testUnsafeArtifacts() { /not a regular directory/, ); + assert.equal( + await readFile(outsideFile, "utf8"), + "outside remains unchanged\n", + ); const linkedFile = path.join(context.root, "linked.json"); - await symlink(path.join(outside, "outside.json"), linkedFile); + await symlink(outsideFile, linkedFile); await assert.rejects( io.artifactDestination(context, ["linked.json"], "scan_manifest"), /not a regular file/, ); + assert.equal( + await readFile(outsideFile, "utf8"), + "outside remains unchanged\n", + ); const linkedRoot = path.join(fixture, "linked-root"); await symlink(context.root, linkedRoot); diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_inventory.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_inventory.mjs index db5265d144..6a759d39c3 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_inventory.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_inventory.mjs @@ -1,22 +1,13 @@ +import { createTemporaryDirectories } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { execFile as nodeExecFile } from "node:child_process"; -import { - mkdir, - mkdtemp, - readFile, - realpath, - rm, - symlink, - unlink, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, symlink, unlink, writeFile } from "node:fs/promises"; import path from "node:path"; import { promisify } from "node:util"; import { importSource } from "./import-module.mjs"; const execFile = promisify(nodeExecFile); -const temporaryRoots = []; +const temporaryDirectories = createTemporaryDirectories(true); const inventory = await importSource( new URL("../src/artifact-inventory.ts", import.meta.url).pathname, @@ -39,16 +30,13 @@ try { await testBoundScopeFailurePreservesPreviousInventory(); await testInvalidDiffTargetPreservesPreviousInventory(); } finally { - await Promise.all( - temporaryRoots.map((root) => rm(root, { force: true, recursive: true })), - ); + await temporaryDirectories.cleanup(); } async function testSchemasAreBoundAndExact() { const scanId = "f84c8312-a602-4660-8e01-518a176cd75a"; const prepare = inventory.prepareReviewItemsInputSchema; const parent = inventory.reviewItemsReaderInputSchema; - const worker = inventory.reviewItemsWorkerReaderInputSchema; assert.equal(prepare.safeParse({ scanId }).success, true); assert.equal( @@ -65,33 +53,6 @@ async function testSchemasAreBoundAndExact() { assert.equal(parent.safeParse({ scanId, limit: 0 }).success, false); assert.equal(parent.safeParse({ scanId, limit: 1001 }).success, false); assert.equal(parent.safeParse({ scanId, cursor: "-1" }).success, false); - assert.equal(worker.safeParse({ limit: 2, cursor: "0" }).success, true); - assert.equal(worker.safeParse({ scanId }).success, false); - assert.equal(worker.safeParse({ scope: "." }).success, false); - assert.equal( - inventory.prepareReviewItemsOutputSchema.safeParse({ reviewItemsTotal: 0 }) - .success, - true, - ); - assert.equal( - inventory.prepareReviewItemsOutputSchema.safeParse({ - reviewItemsTotal: 0, - path: "leaked", - }).success, - false, - ); - assert.equal( - inventory.reviewItemsReaderOutputSchema.safeParse({ - items: [{ path: "src/a.ts" }], - }).success, - true, - ); - assert.equal( - inventory.reviewItemsReaderOutputSchema.safeParse({ - items: [{ path: "src/a.ts", area: "src" }], - }).success, - false, - ); } async function testPrepareUsesTheExistingStandardGenerator() { @@ -462,10 +423,9 @@ async function testInvalidDiffTargetPreservesPreviousInventory() { } async function createFixture(label) { - const root = await realpath( - await mkdtemp(path.join(tmpdir(), "security-artifact-inventory-")), + const root = await temporaryDirectories.create( + "security-artifact-inventory-", ); - temporaryRoots.push(root); const fixtureRoot = path.join(root, label); const repoRoot = path.join(fixtureRoot, "repository"); const scanRoot = path.join(fixtureRoot, "scan"); diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_scan_draft.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_scan_draft.mjs index 952ad27853..c95bcb9c24 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_scan_draft.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_scan_draft.mjs @@ -1,19 +1,18 @@ +import { mock } from "node:test"; +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { createHash } from "node:crypto"; import { promises as fsPromises } from "node:fs"; import { mkdir, - mkdtemp, readdir, readFile, - realpath, rename, rm, symlink, utimes, writeFile, } from "node:fs/promises"; -import { tmpdir } from "node:os"; import path from "node:path"; import { claimToken, @@ -32,9 +31,9 @@ const { scanDraftInputSchema, } = draftApi; -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "codex-security-scan-draft-")), -); +const surfaceDisposition = ({ id, disposition }) => ({ id, disposition }); + +const root = await temporaryDirectory("codex-security-scan-draft-", true); try { const { context } = draftFixture(root, "standard"); @@ -337,14 +336,9 @@ try { ...workerInput, complete: false, }); - const { - scope: _scope, - threatModel: _threatModel, - ...workerWithoutContext - } = workerInput; await recordCodexSecurityWorkerScanDraft( carriedContext, - workerWithoutContext, + withoutScanContext(workerInput), ); const carriedWorker = JSON.parse( await readFile(path.join(carriedContextRoot, "result.json"), "utf8"), @@ -394,13 +388,9 @@ try { const progressedCoverage = JSON.parse( await readFile(path.join(coverageProgressRoot, "result.json"), "utf8"), ).coverage; - assert.deepEqual( - progressedCoverage.surfaces.map(({ id, disposition }) => ({ - id, - disposition, - })), - [{ id: "surface-archive", disposition: "reported" }], - ); + assert.deepEqual(progressedCoverage.surfaces.map(surfaceDisposition), [ + { id: "surface-archive", disposition: "reported" }, + ]); assert.deepEqual(progressedCoverage.openQuestions ?? [], []); assert.equal(progressedCoverage.completeness, "complete"); @@ -719,14 +709,9 @@ try { true, ); - const { - scope: _parentScope, - threatModel: _parentThreatModel, - ...parentWithoutContext - } = input; await recordCodexSecurityScanDraft( { ...context, root: parentCheckpointRoot }, - parentWithoutContext, + withoutScanContext(input), ); const carriedParentManifest = await readJson( parentCheckpointRoot, @@ -866,34 +851,31 @@ try { "obsolete.json", ); await writeFile(obsoleteCheckpointPath, "{malformed obsolete checkpoint\n"); - let deepWorkbenchWrites = 0; + const deepWorkbenchWrites = mock.fn(async (arguments_) => { + assert.deepEqual(arguments_.slice(0, 3), [ + "write-scan-draft", + "--scan-id", + scanId, + ]); + assert.equal(arguments_.includes("--expected-draft-digest"), false); + assert.deepEqual(arguments_.slice(-2), ["--claim-token", claimToken]); + const draftPath = arguments_[arguments_.indexOf("--draft-path") + 1]; + const checkpointPath = + arguments_[arguments_.indexOf("--checkpoint-path") + 1]; + const staged = JSON.parse(await readFile(draftPath, "utf8")); + const stagedCheckpoint = JSON.parse(await readFile(checkpointPath, "utf8")); + assert.deepEqual(staged.findings, acceptedDeepFindings); + assert.deepEqual(staged.coverage, acceptedDeepCoverage); + assert.deepEqual(stagedCheckpoint.findings, acceptedDeepDraft.findings); + assert.equal(stagedCheckpoint.handoffClaimToken, undefined); + }); await recordCodexSecurityScanDraftViaWorkbench( deepParentContext, acceptedDeepDraft, - async (arguments_) => { - deepWorkbenchWrites += 1; - assert.deepEqual(arguments_.slice(0, 3), [ - "write-scan-draft", - "--scan-id", - scanId, - ]); - assert.equal(arguments_.includes("--expected-draft-digest"), false); - assert.deepEqual(arguments_.slice(-2), ["--claim-token", claimToken]); - const draftPath = arguments_[arguments_.indexOf("--draft-path") + 1]; - const checkpointPath = - arguments_[arguments_.indexOf("--checkpoint-path") + 1]; - const staged = JSON.parse(await readFile(draftPath, "utf8")); - const stagedCheckpoint = JSON.parse( - await readFile(checkpointPath, "utf8"), - ); - assert.deepEqual(staged.findings, acceptedDeepFindings); - assert.deepEqual(staged.coverage, acceptedDeepCoverage); - assert.deepEqual(stagedCheckpoint.findings, acceptedDeepDraft.findings); - assert.equal(stagedCheckpoint.handoffClaimToken, undefined); - }, + deepWorkbenchWrites, ); assert.equal( - deepWorkbenchWrites, + deepWorkbenchWrites.mock.callCount(), 1, "terminal Deep drafts still publish through the workbench lock despite obsolete malformed checkpoints", ); @@ -1782,23 +1764,22 @@ try { assert.equal(retried.status, "draft_written"); const conflictAbort = new AbortController(); - let abortedConflictAttempts = 0; + const abortedConflictAttempts = mock.fn(async () => { + conflictAbort.abort(new Error("draft publication canceled")); + throw Object.assign(new Error("scan_draft_conflict"), { + code: "scan_draft_conflict", + }); + }); await assert.rejects( recordCodexSecurityScanDraft( context, input, - async () => { - abortedConflictAttempts += 1; - conflictAbort.abort(new Error("draft publication canceled")); - throw Object.assign(new Error("scan_draft_conflict"), { - code: "scan_draft_conflict", - }); - }, + abortedConflictAttempts, conflictAbort.signal, ), /draft publication canceled/, ); - assert.equal(abortedConflictAttempts, 1); + assert.equal(abortedConflictAttempts.mock.callCount(), 1); const monotonicRoot = path.join(root, "monotonic-final-draft"); await mkdir(monotonicRoot); @@ -1849,13 +1830,9 @@ try { ); } assert.notEqual(staged.manifest.scan.complete, false); - assert.deepEqual( - staged.coverage.surfaces.map(({ id, disposition }) => ({ - id, - disposition, - })), - [{ id: "surface-archive", disposition: "reported" }], - ); + assert.deepEqual(staged.coverage.surfaces.map(surfaceDisposition), [ + { id: "surface-archive", disposition: "reported" }, + ]); }, ); assert.equal(monotonicWrites, 2); @@ -1884,13 +1861,8 @@ try { path.join(archivedRetryRoot, "attempts", "attempt-01"), ); await mkdir(archivedRetryOutput); - const { - scope: _archivedScope, - threatModel: _archivedThreatModel, - ...replacementAttempt - } = workerInput; await recordCodexSecurityWorkerScanDraft(archivedRetryContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }); const archivedRetryResult = JSON.parse( @@ -1947,7 +1919,7 @@ try { ); await mkdir(archivedResolutionOutput); await recordCodexSecurityWorkerScanDraft(archivedResolutionContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }); const archivedResolutionResult = JSON.parse( @@ -2023,7 +1995,7 @@ try { ); await mkdir(repeatedCheckpointOutput); await recordCodexSecurityWorkerScanDraft(repeatedCheckpointContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }); const repeatedCheckpointResult = JSON.parse( @@ -2053,7 +2025,7 @@ try { ); await mkdir(multiAttemptOutput); await recordCodexSecurityWorkerScanDraft(multiAttemptContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), complete: false, findings: [], coverage: { @@ -2069,7 +2041,7 @@ try { ); await mkdir(multiAttemptOutput); await recordCodexSecurityWorkerScanDraft(multiAttemptContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }); const multiAttemptResult = JSON.parse( @@ -2111,7 +2083,7 @@ try { ); await mkdir(malformedArchivedOutput); await recordCodexSecurityWorkerScanDraft(malformedArchivedContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }); assert.deepEqual( @@ -2146,7 +2118,7 @@ try { await mkdir(crossScanRetryOutput); await assert.rejects( recordCodexSecurityWorkerScanDraft(crossScanRetryContext, { - ...replacementAttempt, + ...withoutScanContext(workerInput), findings: [], }), /scanId does not match the authoritative workbench scan/u, @@ -3463,3 +3435,11 @@ async function recordFreshScanDraft(context, input) { ]); return recordCodexSecurityScanDraft(context, input); } + +function withoutScanContext({ + scope: _scope, + threatModel: _threatModel, + ...input +}) { + return input; +} diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_storage.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_storage.mjs index 030004e832..048974c9ed 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_storage.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_storage.mjs @@ -1,8 +1,8 @@ +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { execFileSync } from "node:child_process"; import { mkdir, - mkdtemp, readFile, realpath, rm, @@ -17,9 +17,7 @@ import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js" import { applicationRoot, buildServer } from "./build-server.mjs"; -const fixture = await realpath( - await mkdtemp(path.join(tmpdir(), "codex-security-storage-test-")), -); +const fixture = await temporaryDirectory("codex-security-storage-test-", true); const stateRoot = path.join(fixture, "state"); const repository = path.join(fixture, "repository"); const bundle = path.join(fixture, "server.cjs"); diff --git a/plugins/codex-security/mcp-app/tests/test_artifact_validation_phase.mjs b/plugins/codex-security/mcp-app/tests/test_artifact_validation_phase.mjs index d99107fbe1..4146c7fd12 100644 --- a/plugins/codex-security/mcp-app/tests/test_artifact_validation_phase.mjs +++ b/plugins/codex-security/mcp-app/tests/test_artifact_validation_phase.mjs @@ -1,15 +1,7 @@ +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { randomUUID } from "node:crypto"; -import { - mkdir, - mkdtemp, - readFile, - realpath, - rm, - symlink, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, rm, symlink, writeFile } from "node:fs/promises"; import path from "node:path"; import { importSource } from "./import-module.mjs"; @@ -93,9 +85,7 @@ assert.equal( false, ); -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "codex-security-validation-phase-")), -); +const root = await temporaryDirectory("codex-security-validation-phase-", true); try { const context = await scanContext(root, "scan", scanId); const ledger = path.join( @@ -297,9 +287,7 @@ async function assertNoMutation(context, ledger, input, expectedError) { async function writeJsonl(file, rows) { await writeFile( file, - rows.length > 0 - ? `${rows.map((row) => JSON.stringify(row)).join("\n")}\n` - : "", + rows.length > 0 ? `${rows.map(JSON.stringify).join("\n")}\n` : "", ); } diff --git a/plugins/codex-security/mcp-app/tests/test_compact_artifact_server.mjs b/plugins/codex-security/mcp-app/tests/test_compact_artifact_server.mjs index 09f245ba6c..25c7a22013 100644 --- a/plugins/codex-security/mcp-app/tests/test_compact_artifact_server.mjs +++ b/plugins/codex-security/mcp-app/tests/test_compact_artifact_server.mjs @@ -1,3 +1,4 @@ +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import { readOnlyParentSandboxState } from "./sandbox-state.mjs"; import assert from "node:assert/strict"; import { execFileSync } from "node:child_process"; @@ -10,7 +11,6 @@ import { rm, writeFile, } from "node:fs/promises"; -import { tmpdir } from "node:os"; import path from "node:path"; import { Client } from "@modelcontextprotocol/sdk/client/index.js"; import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; @@ -21,9 +21,7 @@ const pluginRoot = path.resolve(applicationRoot, ".."); const bundledPluginRoot = process.env.CODEX_SECURITY_TEST_PLUGIN_ROOT ? path.resolve(process.env.CODEX_SECURITY_TEST_PLUGIN_ROOT) : path.resolve(applicationRoot, "../../../sdk/typescript/_bundled_plugin"); -const temporaryRoot = await mkdtemp( - path.join(tmpdir(), "codex-security-artifact-mcp-"), -); +const temporaryRoot = await temporaryDirectory("codex-security-artifact-mcp-"); try { const runtimeBundle = path.join(temporaryRoot, "server.cjs"); @@ -93,12 +91,7 @@ async function testCompactDiffScanCompletion(bundle, runtimeLabel) { CODEX_SECURITY_STATE_DIR: stateRoot, }); const ownerThread = `compact-diff-owner-${runtimeLabel}`; - const call = (name, arguments_) => - client.callTool({ - name, - arguments: arguments_, - _meta: { "openai/threadId": ownerThread }, - }); + const call = toolCaller(client, ownerThread); try { const selection = { @@ -303,12 +296,7 @@ async function testSemanticScanDraftCompletion(bundle, runtimeLabel) { CODEX_SECURITY_STATE_DIR: stateRoot, }); const ownerThread = `semantic-draft-owner-${runtimeLabel}`; - const call = (name, arguments_) => - client.callTool({ - name, - arguments: arguments_, - _meta: { "openai/threadId": ownerThread }, - }); + const call = toolCaller(client, ownerThread); try { const opened = requireSuccessfulTool( @@ -904,12 +892,7 @@ async function testClaimedParentArtifactOperations(bundle, runtimeLabel) { }); const ownerThread = `compact-artifact-owner-${runtimeLabel}`; const otherThread = `compact-artifact-other-${runtimeLabel}`; - const call = (name, arguments_, threadId = ownerThread) => - client.callTool({ - name, - arguments: arguments_, - ...(threadId == null ? {} : { _meta: { "openai/threadId": threadId } }), - }); + const call = toolCaller(client, ownerThread); try { const opened = requireSuccessfulTool( @@ -1738,3 +1721,12 @@ async function startClient(bundle, environment) { await client.connect(transport); return client; } + +function toolCaller(client, ownerThread) { + return (name, arguments_, threadId = ownerThread) => + client.callTool({ + name, + arguments: arguments_, + ...(threadId == null ? {} : { _meta: { "openai/threadId": threadId } }), + }); +} From 5f195eed821cd34a79ce6bfc24d8e8227b30cbdc Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 3 Oct 2026 20:31:05 +0000 Subject: [PATCH 2/4] test(plugin): declare temporary directory factory --- .../mcp-app/tests/support/temporary-directories.d.mts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts b/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts index 20e0e9a407..6b4a225dbc 100644 --- a/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts +++ b/plugins/codex-security/mcp-app/tests/support/temporary-directories.d.mts @@ -1,3 +1,8 @@ +export function createTemporaryDirectories(canonicalize?: boolean): { + create(prefix: string): Promise; + cleanup(): Promise; +}; + export function temporaryDirectory( prefix: string, canonicalize?: boolean, From c7fcfd8d28cad167f420725941e2f831714d63a5 Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 3 Oct 2026 19:14:48 +0000 Subject: [PATCH 3/4] refactor(plugin): simplify Deep Scan coordinator ownership The Deep Scan coordinator repeats lifecycle, failure-forwarding and store plumbing. Consolidate those operations at their existing owners. --- plugins/codex-security/mcp-app/server.ts | 72 ++--- .../mcp-app/src/deep-scan/coordinator.ts | 45 +-- .../mcp-app/src/deep-scan/registry.ts | 110 +++---- .../mcp-app/src/deep-scan/store.ts | 15 +- .../mcp-app/src/deep-scan/types.ts | 64 +--- .../mcp-app/src/deep-scan/worker-runner.ts | 4 +- .../tests/deep_scan_deadline_cases.mjs | 20 +- .../tests/deep_scan_publication_cases.mjs | 12 +- .../tests/deep_scan_worker_failure_cases.mjs | 50 ++- .../mcp-app/tests/support/streams.mjs | 94 ++++++ .../test_deep_scan_artifact_validation.mjs | 16 +- .../tests/test_deep_scan_coordinator.mjs | 296 ++++++++---------- .../tests/test_deep_scan_stdio_lifecycle.mjs | 99 +----- .../mcp-app/tests/test_deep_scan_store.mjs | 80 +++-- .../test_deep_scan_store_integration.mjs | 15 +- .../tests/test_workbench_state_fallback.mjs | 75 +---- 16 files changed, 451 insertions(+), 616 deletions(-) create mode 100644 plugins/codex-security/mcp-app/tests/support/streams.mjs diff --git a/plugins/codex-security/mcp-app/server.ts b/plugins/codex-security/mcp-app/server.ts index 48ee45d44c..01a4124b33 100644 --- a/plugins/codex-security/mcp-app/server.ts +++ b/plugins/codex-security/mcp-app/server.ts @@ -1,3 +1,4 @@ +import type { JsonObject } from "./src/types.js"; import { isRecord as isJsonObject } from "./src/record.js"; import { execFile } from "node:child_process"; import { randomUUID } from "node:crypto"; @@ -16,13 +17,14 @@ import { handoffClaimTokenSchema, recoveryHandoffClaimTokenSchema, registerScanHandoffTools, + type HandoffWorkspaceState as WorkspaceState, } from "./src/server/handoff-tools.js"; import { registerCompactArtifactTools } from "./src/server/compact-artifact-tools.js"; import { createScanArtifactContext } from "./src/artifact-context.js"; import { recordCodexSecurityScanDraftViaWorkbench } from "./src/artifact-scan-draft.js"; import { DeepScanCoordinatorRegistry, - DeepScanStartLock, + AsyncLock, startOrJoinDeepScanCoordinator, } from "./src/deep-scan/registry.js"; import { CodexSdkWorkerExecutor } from "./src/deep-scan/executor.js"; @@ -48,12 +50,10 @@ const WORKBENCH_COMMANDS_WITHOUT_DATABASE = new Set([ "read-artifact", ]); -type JsonObject = Record; - let fallbackWorkbenchStateDir: Promise | undefined; let fallbackWorkbenchStateLogged = false; let persistentWorkbenchStateSucceeded = false; -let workbenchStateSelectionTail: Promise = Promise.resolve(); +const workbenchStateSelectionLock = new AsyncLock(); const userContextSchema = z.string().trim().min(1); const editableUserContextSchema = z.string().trim(); @@ -126,14 +126,6 @@ async function scanRoot(): Promise { return result.scanRoot; } -interface WorkspaceState extends JsonObject { - id: string; - results?: JsonObject; - setup: { - submitted: boolean; - }; -} - const diffTargetSchema = z.discriminatedUnion("kind", [ z .object({ @@ -613,7 +605,8 @@ export function createCodexSecurityServer(): McpServer { }, ); const deepScanCoordinators = new DeepScanCoordinatorRegistry(); - const deepScanStartLock = new DeepScanStartLock(); + // Serialize start-or-join so a scan creates only one coordinator. + const deepScanStartLock = new AsyncLock(); const deepScanStore = new WorkbenchDeepScanStore(runWorkbench); const authenticatedArtifactClaims = new Map< string, @@ -1799,7 +1792,7 @@ export function createCodexSecurityServer(): McpServer { message, ...optionalArg("--claim-token", handoffClaimToken), ]); - deepScanCoordinators.failExternallyPersisted(scanId, message); + deepScanCoordinators.get(scanId)?.failExternallyPersisted(message); return scanActionResult( failed, "Recorded the Codex Security scan failure.", @@ -2323,12 +2316,11 @@ function promptOnlyScanResult(promptOnly: JsonObject) { "Codex Security prompt-only scan returned malformed context; no prompt-driven scan was started.", ); } - const disposition = startDisposition === "joined" ? "Rejoined" : "Started"; return { content: [ { type: "text" as const, - text: `${disposition} prompt-driven scan ${scanId}. Use the returned scanId and scanDir for every phase. Author scan-manifest.json as an unsealed draft: omit scan.sealedAt and scan.artifacts because completion supplies the exact workbench timestamps, seal, artifact digests, and derived finding identities. Then call complete_codex_security_scan once to index the completed findings.`, + text: `${startDisposition === "joined" ? "Rejoined" : "Started"} prompt-driven scan ${scanId}. Use the returned scanId and scanDir for every phase. Author scan-manifest.json as an unsealed draft: omit scan.sealedAt and scan.artifacts because completion supplies the exact workbench timestamps, seal, artifact digests, and derived finding identities. Then call complete_codex_security_scan once to index the completed findings.`, }, ], structuredContent: promptOnly, @@ -2424,15 +2416,14 @@ async function logUserInputFailure( function boundedErrorData(error: unknown): { message: string; name: string } { const name = error instanceof Error && error.name.trim() ? error.name : "UnknownError"; - const message = - error instanceof Error + return { + name: name.slice(0, 128), + message: (error instanceof Error ? error.message : typeof error === "string" ? error - : "Unknown user-input elicitation failure."; - return { - name: name.slice(0, 128), - message: message.slice(0, 1000), + : "Unknown user-input elicitation failure." + ).slice(0, 1000), }; } @@ -2532,7 +2523,7 @@ async function executeWorkbenchWithStateSelection( if (persistentWorkbenchStateSucceeded) { return await executeWorkbench(pythonCommand, args, undefined, input); } - return await withWorkbenchStateSelectionLock(async () => { + return await workbenchStateSelectionLock.run(async () => { if (fallbackWorkbenchStateDir) { return await executeWorkbench( pythonCommand, @@ -2567,22 +2558,6 @@ async function executeWorkbenchWithStateSelection( }); } -async function withWorkbenchStateSelectionLock( - operation: () => Promise, -): Promise { - const predecessor = workbenchStateSelectionTail; - let release!: () => void; - workbenchStateSelectionTail = new Promise((resolvePromise) => { - release = resolvePromise; - }); - await predecessor; - try { - return await operation(); - } finally { - release(); - } -} - async function executeWorkbench( pythonCommand: string, args: string[], @@ -2807,12 +2782,10 @@ function deepScanInvocationFailureMessage(error: unknown): string { } function deepScanFailureMessage(run: DeepScanRunState): string { - const manifest = run.manifestPath - ? ` Failure manifest: ${run.manifestPath}.` - : ""; - const diagnostic = `${run.error ?? `Deep Scan ${run.scanId} ${run.status}.`}${manifest}`; return [ - diagnostic, + `${run.error ?? `Deep Scan ${run.scanId} ${run.status}.`}${ + run.manifestPath ? ` Failure manifest: ${run.manifestPath}.` : "" + }`, "This is a terminal failure of this logical Deep Scan; no successful discovery manifest was returned.", "Stop further scanning and surface this exact stable MCP failure. Read the existing scan context to report saved findings and pending candidates separately, with incomplete coverage.", "Do not call start_codex_security_deep_scan again in this response.", @@ -2822,12 +2795,11 @@ function deepScanFailureMessage(run: DeepScanRunState): string { } function isUnwritableSqliteOpenError(error: unknown): boolean { - const diagnostic = isExecError(error) - ? error.stderr - : error instanceof Error - ? error.message - : ""; return /sqlite3\.OperationalError:\s*unable to open database file/i.test( - diagnostic, + isExecError(error) + ? error.stderr + : error instanceof Error + ? error.message + : "", ); } 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 3e4925ed0c..29e492ff85 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts @@ -1,3 +1,4 @@ +/// import { randomUUID } from "node:crypto"; import { promises as fs } from "node:fs"; import { join } from "node:path"; @@ -162,7 +163,7 @@ export class DeepScanCoordinator { scanId: this.state.scanId, reason: errorKind(error), }); - this.failLocally(error); + if (this.stopLocally()) this.rejectTerminal(error); }); } @@ -325,7 +326,7 @@ export class DeepScanCoordinator { scanId: this.state.scanId, reason: schedulerResult.reason, }); - this.finishLocally(this.state); + if (this.stopLocally()) this.resolveTerminal(cloneState(this.state)); } catch (error) { if (this.canceled || this.externallyFailed) { return; @@ -447,29 +448,17 @@ export class DeepScanCoordinator { }); } if (this.canceled) this.state = { ...this.state, status: "canceled" }; - this.finishLocally(this.state); + if (this.stopLocally()) this.resolveTerminal(cloneState(this.state)); } } - private finishLocally(state: DeepScanRunState): void { - if (this.markTerminal()) this.resolveTerminal(cloneState(state)); - } - - private failLocally(error: unknown): void { - if (this.markTerminal()) this.rejectTerminal(error); - } - - private markTerminal(): boolean { + private stopLocally(): boolean { if (this.terminal) return false; this.terminal = true; - if (this.heartbeatTimeout !== undefined) { - clearTimeout(this.heartbeatTimeout); - this.heartbeatTimeout = undefined; - } - if (this.discoveryTimeout !== undefined) { - clearTimeout(this.discoveryTimeout); - this.discoveryTimeout = undefined; - } + clearTimeout(this.heartbeatTimeout); + this.heartbeatTimeout = undefined; + clearTimeout(this.discoveryTimeout); + this.discoveryTimeout = undefined; return true; } @@ -601,7 +590,7 @@ export class DeepScanCoordinator { } else { this.state = current; } - this.finishLocally(this.state); + if (this.stopLocally()) this.resolveTerminal(cloneState(this.state)); return true; } @@ -667,11 +656,9 @@ export class DeepScanCoordinator { ); } if (reducerFailures >= errorLimit) { - const failure = [...(this.state.persistedWorkers ?? [])] - .reverse() - .find( - (worker) => worker.kind === "dedup" && worker.status === "failed", - ); + const failure = (this.state.persistedWorkers ?? []).findLast( + (worker) => worker.kind === "dedup" && worker.status === "failed", + ); throw reducerErrorLimitError( reducerFailures, errorLimit, @@ -1046,10 +1033,8 @@ export class DeepScanCoordinator { private trackSchedulerWork(promise: Promise): Promise { this.schedulerWork.add(promise); - void promise.then( - () => this.schedulerWork.delete(promise), - () => this.schedulerWork.delete(promise), - ); + const cleanup = () => this.schedulerWork.delete(promise); + void promise.then(cleanup, cleanup); return promise; } diff --git a/plugins/codex-security/mcp-app/src/deep-scan/registry.ts b/plugins/codex-security/mcp-app/src/deep-scan/registry.ts index b0f47ab173..6b957cef14 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/registry.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/registry.ts @@ -23,7 +23,6 @@ export class DeepScanCoordinatorRegistry { start(options: CoordinatorOptions): DeepScanCoordinator { const existing = this.coordinators.get(options.run.scanId); if (existing) return existing; - const { observeReplacement: _unused, ...remoteOptions } = options; let coordinator!: DeepScanCoordinator; coordinator = new DeepScanCoordinator({ ...options, @@ -31,12 +30,11 @@ export class DeepScanCoordinatorRegistry { if (this.coordinators.get(run.scanId) === coordinator) { this.coordinators.delete(run.scanId); } - const observer = new DeepScanRemoteCoordinator({ + return await new DeepScanRemoteCoordinator({ run, registry: this, - options: remoteOptions, - }); - return await observer.wait(undefined); + options, + }).wait(undefined); }, }); this.coordinators.set(options.run.scanId, coordinator); @@ -61,21 +59,14 @@ export class DeepScanCoordinatorRegistry { return true; } - failExternallyPersisted(scanId: string, reason: string): boolean { - const coordinator = this.coordinators.get(scanId); - if (!coordinator) return false; - coordinator.failExternallyPersisted(reason); - return true; - } - shutdown(reason: string): void { for (const coordinator of this.coordinators.values()) coordinator.cancel(reason); } } -/** Serializes start-or-join calls so one scan can create only one coordinator. */ -export class DeepScanStartLock { +/** Serializes asynchronous operations in invocation order. */ +export class AsyncLock { private tail: Promise = Promise.resolve(); async run(operation: () => Promise): Promise { @@ -125,61 +116,58 @@ export class DeepScanRemoteCoordinator { timeoutMs === undefined ? undefined : Date.now() + timeoutMs; while (true) { if (signal?.aborted) throw remoteAbortError(signal.reason); - let run: DeepScanRunState; + let run: DeepScanRunState | undefined; try { run = await options.store.get(this.input.run.scanId, threadId); } catch (error) { if (!isTransientPersistenceError(error)) throw error; - const remaining = - deadline === undefined ? COORDINATOR_POLL_MS : deadline - Date.now(); - if (remaining <= 0) return undefined; - await delay(Math.min(COORDINATOR_POLL_MS, remaining), undefined, { - signal, - }); - continue; } - if (run.status !== "running") return run; - - const heartbeat = run.updatedAt ? Date.parse(run.updatedAt) : Number.NaN; - if ( - (!Number.isFinite(heartbeat) || - Date.now() - heartbeat >= COORDINATOR_LEASE_MS) && - Date.now() >= this.nextClaimAt - ) { - const local = registry.get(run.scanId); - if (local) { - return deadline === undefined - ? await local.wait(signal) - : await local.wait(signal, Math.max(0, deadline - Date.now())); - } - let claim: DeepScanCoordinatorClaim | undefined; - try { - claim = await options.store.claimCoordinator({ - scanId: run.scanId, - threadId, - handoffClaimToken: options.handoffClaimToken, - }); - } catch (error) { - if (!isTransientPersistenceError(error)) { - try { - const latest = await options.store.get(run.scanId, threadId); - if (latest.status !== "running") return latest; - throw error; - } catch (readError) { - if (!isTransientPersistenceError(readError)) throw readError; + if (run !== undefined) { + if (run.status !== "running") return run; + + const heartbeat = run.updatedAt + ? Date.parse(run.updatedAt) + : Number.NaN; + if ( + (!Number.isFinite(heartbeat) || + Date.now() - heartbeat >= COORDINATOR_LEASE_MS) && + Date.now() >= this.nextClaimAt + ) { + const local = registry.get(run.scanId); + if (local) { + return deadline === undefined + ? await local.wait(signal) + : await local.wait(signal, Math.max(0, deadline - Date.now())); + } + let claim: DeepScanCoordinatorClaim | undefined; + try { + claim = await options.store.claimCoordinator({ + scanId: run.scanId, + threadId, + handoffClaimToken: options.handoffClaimToken, + }); + } catch (error) { + if (!isTransientPersistenceError(error)) { + try { + const latest = await options.store.get(run.scanId, threadId); + if (latest.status !== "running") return latest; + throw error; + } catch (readError) { + if (!isTransientPersistenceError(readError)) throw readError; + } } } + if (claim?.acquired) { + const coordinator = registry.start({ ...options, run: claim.run }); + return deadline === undefined + ? await coordinator.wait(signal) + : await coordinator.wait( + signal, + Math.max(0, deadline - Date.now()), + ); + } + if (claim) this.nextClaimAt = Date.now() + COORDINATOR_LEASE_MS; } - if (claim?.acquired) { - const coordinator = registry.start({ ...options, run: claim.run }); - return deadline === undefined - ? await coordinator.wait(signal) - : await coordinator.wait( - signal, - Math.max(0, deadline - Date.now()), - ); - } - if (claim) this.nextClaimAt = Date.now() + COORDINATOR_LEASE_MS; } const remaining = diff --git a/plugins/codex-security/mcp-app/src/deep-scan/store.ts b/plugins/codex-security/mcp-app/src/deep-scan/store.ts index 901b7daa50..7e4487cbb9 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/store.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/store.ts @@ -15,7 +15,6 @@ import type { DeepScanMergeState, DeepScanRunState, DeepScanRunStatus, - DeepScanStore, DeepScanTerminalReason, DeepScanWorkerKind, DeepScanWorkerMutation, @@ -110,7 +109,7 @@ class DeepScanPersistenceError extends Error { } } -export class WorkbenchDeepScanStore implements DeepScanStore { +export class WorkbenchDeepScanStore { private writeTail: Promise = Promise.resolve(); private readonly coordinatorLeases = new Map< string, @@ -260,16 +259,6 @@ export class WorkbenchDeepScanStore implements DeepScanStore { return { ...lease.run, updatedAt }; } - async cancel(scanId: string, threadId: string): Promise { - return this.enqueueWrite([ - "cancel-scan", - "--scan-id", - scanId, - "--thread-id", - threadId, - ]); - } - async updateWorker( update: DeepScanWorkerMutation, ): Promise { @@ -475,7 +464,7 @@ export class WorkbenchDeepScanStore implements DeepScanStore { /** * Run mutations in call order. Callers receive their own operation's result * or error, while the stored tail always resolves so one failed write cannot - * prevent later cancellation or cleanup from reaching the workbench. + * prevent later persistence or cleanup from reaching the workbench. * * This orders one Node store instance; SQLite still provides transactions for * other workbench processes. Reads remain concurrent and observe a committed diff --git a/plugins/codex-security/mcp-app/src/deep-scan/types.ts b/plugins/codex-security/mcp-app/src/deep-scan/types.ts index 0ba1acc1d6..33f63cbd25 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/types.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/types.ts @@ -1,4 +1,5 @@ import type { DeepReducerContext } from "../artifact-io.js"; +import type { WorkbenchDeepScanStore } from "./store.js"; export type DeepScanTerminalReason = "saturated" | "capped"; @@ -105,65 +106,10 @@ export interface DedupCommit { } /** Durable operations implemented by the Python workbench. */ -export interface DeepScanStore { - begin(input: { - scanId?: string; - targetPath?: string; - scope?: string; - userContext?: string; - handoffClaimToken?: string; - model?: string; - reasoningEffort?: string; - threadId: string; - scanRoot: string; - }): Promise; - get(scanId: string, threadId: string): Promise; - claimCoordinator( - input: DeepScanCoordinatorLeaseInput, - ): Promise; - heartbeatCoordinator( - input: DeepScanCoordinatorLeaseInput, - ): Promise; - cancel(scanId: string, threadId: string): Promise>; - updateWorker( - update: DeepScanWorkerMutation, - ): Promise; - claimDedup(input: { - id: string; - scanId: string; - workerIds: string[]; - promptPath: string; - artifactDir: string; - }): Promise; - commitDedup(commit: DedupCommit): Promise; - finish(input: { - scanId: string; - reason: DeepScanTerminalReason; - manifestPath: string; - stagedManifestPath?: string; - omittedWorkerIds: string[]; - }): Promise; - fail( - scanId: string, - message: string, - status?: "failed" | "interrupted", - manifestPath?: string, - stagedManifestPath?: string, - ): Promise; - recordStoppedPublicationFailure( - scanId: string, - message: string, - coordinatorGeneration?: number, - ): Promise; - updateProgress(input: { - scanId: string; - handoffClaimToken?: string; - phase?: "preflight" | "discovery"; - deepReviewPass?: number; - reviewItemsTotal?: number; - reviewItemsCompleted?: number; - }): Promise; -} +export type DeepScanStore = Omit< + WorkbenchDeepScanStore, + "begin" | "coordinatorLeaseArgs" +>; /** Host-bound worker artifact state; never populate this from model input. */ export interface CodexWorkerArtifactContext { 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 790449a859..5a3f469352 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 @@ -428,8 +428,8 @@ export class DeepScanWorkerRunner { const result = await this.options.executor.run({ kind: input.kind, promptPath: executionPromptPath, - // Discovery workers write only to their isolated directory. Setup and - // dedup workers own shared scan artifacts; the target remains read-only. + // Discovery workers write only to their isolated directory. Reducers + // own shared scan artifacts; the target remains read-only. workingDirectory: input.kind === "discovery" ? input.artifactDir diff --git a/plugins/codex-security/mcp-app/tests/deep_scan_deadline_cases.mjs b/plugins/codex-security/mcp-app/tests/deep_scan_deadline_cases.mjs index c3a5f5021d..e25190cd30 100644 --- a/plugins/codex-security/mcp-app/tests/deep_scan_deadline_cases.mjs +++ b/plugins/codex-security/mcp-app/tests/deep_scan_deadline_cases.mjs @@ -30,7 +30,7 @@ export async function testDeepScanDeadlines({ }); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; await eventually(() => executor.runningDiscovery > 0); await eventually(() => executor.runningDiscovery === 0); @@ -47,7 +47,7 @@ export async function testDeepScanDeadlines({ const terminal = await coordinator.wait(undefined, 5_000); assert.equal(terminal?.status, "succeeded"); assert.equal(terminal?.terminalReason, "capped"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(executor.discoveryCalls, discoveryCallsAtDeadline); assert.equal( terminal.dispatchedCount < fixture.run.config.maxDiscoveryRuns, @@ -119,7 +119,7 @@ export async function testDeepScanDeadlines({ terminal.dispatchedCount < fixture.run.config.maxDiscoveryRuns, true, ); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(executor.dedupCalls, 1); assert.equal(executor.runningDiscovery, 0); @@ -147,9 +147,9 @@ export async function testDeepScanDeadlines({ maxDiscoveryRuns: 8, }); const store = new FakeStore(fixture.run); - const acceptancePersisted = deferred(); - const releaseAcceptance = deferred(); - const discoveryDeadlineReached = deferred(); + const acceptancePersisted = Promise.withResolvers(); + const releaseAcceptance = Promise.withResolvers(); + const discoveryDeadlineReached = Promise.withResolvers(); const updateWorker = store.updateWorker.bind(store); store.updateWorker = async (update) => { const persisted = await updateWorker(update); @@ -181,7 +181,7 @@ export async function testDeepScanDeadlines({ const terminal = await terminalWait; assert.equal(terminal?.status, "succeeded"); assert.equal(terminal?.terminalReason, "capped"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(executor.discoveryCalls, 1); assert.equal(executor.dedupCalls, 1); @@ -215,12 +215,12 @@ export async function testDeepScanDeadlines({ onComplete: async (draft) => completedDrafts.push(structuredClone(draft)), }); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; const terminal = await coordinator.wait(undefined, 5_000); assert.equal(terminal?.status, "succeeded"); assert.equal(terminal?.terminalReason, "capped"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(store.finishCalls.length, 1); assert.equal(executor.runningDiscovery, 0); assert.equal(executor.dedupCalls, 0); @@ -261,7 +261,7 @@ export async function testDeepScanDeadlines({ const terminal = await coordinator.wait(undefined, 5_000); assert.equal(terminal?.status, "succeeded"); assert.equal(terminal?.terminalReason, "capped"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(store.finishCalls.length, 1); assert.equal(executor.discoveryCalls, 0); assert.equal(executor.dedupCalls, 0); diff --git a/plugins/codex-security/mcp-app/tests/deep_scan_publication_cases.mjs b/plugins/codex-security/mcp-app/tests/deep_scan_publication_cases.mjs index 062443f1d7..e331060ef8 100644 --- a/plugins/codex-security/mcp-app/tests/deep_scan_publication_cases.mjs +++ b/plugins/codex-security/mcp-app/tests/deep_scan_publication_cases.mjs @@ -18,9 +18,9 @@ export async function testDeepScanPublication({ maxDiscoveryRuns: 3, }); const store = new FakeStore(fixture.run); - const releaseLateWorker = deferred(); - const lateAcceptance = deferred(); - const releaseAcceptance = deferred(); + const releaseLateWorker = Promise.withResolvers(); + const lateAcceptance = Promise.withResolvers(); + const releaseAcceptance = Promise.withResolvers(); const updateWorker = store.updateWorker.bind(store); let acceptedLateWorker; store.updateWorker = async (update) => { @@ -48,7 +48,7 @@ export async function testDeepScanPublication({ onComplete: async (draft) => completed.push(structuredClone(draft)), }); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; releaseLateWorker.resolve(); await lateAcceptance.promise; executor.releaseDedup(); @@ -166,7 +166,7 @@ export async function testDeepScanPublication({ onComplete: async (draft) => completed.push(structuredClone(draft)), }); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; await eventually( () => executor.discoveryCalls === 4 && executor.runningDiscovery === 2, ); @@ -180,7 +180,7 @@ export async function testDeepScanPublication({ ); assert.equal(terminal?.status, "succeeded", terminal?.error); assert.equal(terminal.terminalReason, "saturated"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(store.finishCalls.length, 1); assert.equal(store.finishCalls[0].reason, "saturated"); assert.equal( diff --git a/plugins/codex-security/mcp-app/tests/deep_scan_worker_failure_cases.mjs b/plugins/codex-security/mcp-app/tests/deep_scan_worker_failure_cases.mjs index a478166036..7d6dedf5ed 100644 --- a/plugins/codex-security/mcp-app/tests/deep_scan_worker_failure_cases.mjs +++ b/plugins/codex-security/mcp-app/tests/deep_scan_worker_failure_cases.mjs @@ -1,4 +1,5 @@ import assert from "node:assert/strict"; +import path from "node:path"; import { readFile } from "node:fs/promises"; export function createDeepScanWorkerFailureCases({ @@ -6,6 +7,7 @@ export function createDeepScanWorkerFailureCases({ FakeStore, FakeExecutor, createCoordinator, + DeepScanCoordinator, DeepScanNonRetryableError, classifyCodexWorkerError, deferred, @@ -13,6 +15,43 @@ export function createDeepScanWorkerFailureCases({ workerIdFromPrompt, promptContext, }) { + async function testResumeRequiresHistoricalWorkerPrompt(status) { + const fixture = await fixtureRun({ + workers: 1, + subagents: 0, + stopAfterNoNew: 2, + maxDiscoveryRuns: 2, + }); + const promptPath = path.join(fixture.run.scanDir, "missing-prompt.md"); + const store = new FakeStore({ + ...fixture.run, + persistedWorkers: [ + { + id: "historical-worker", + kind: "discovery", + status, + attempt: 1, + promptPath, + artifactDir: fixture.run.scanDir, + }, + ], + }); + const executor = new FakeExecutor(); + const coordinator = new DeepScanCoordinator({ + run: store.run, + store, + executor, + pluginRoot: fixture.pluginRoot, + clock: immediateClock, + }); + coordinator.start(); + const terminal = await coordinator.wait(undefined, 5_000); + assert.equal(terminal?.status, "failed"); + assert.match(terminal.error, /ENOENT/); + assert.ok(terminal.error.includes(promptPath)); + assert.equal(executor.discoveryCalls, 0); + } + async function testRecoverableWorkerErrorsCannotFailScan() { const failures = ["config unknown", "authentication required"].map( (output) => @@ -99,7 +138,7 @@ export function createDeepScanWorkerFailureCases({ .length, 1, ); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); } } @@ -112,8 +151,8 @@ export function createDeepScanWorkerFailureCases({ maxDiscoveryRuns: 4, }); const store = new FakeStore(fixture.run); - const nextDiscovery = deferred(); - const siblingDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); + const siblingDiscovery = Promise.withResolvers(); const normalExecutor = new FakeExecutor({ discoveryCandidateId: "candidate-1", canonicalCandidateId: "candidate-1", @@ -191,7 +230,7 @@ export function createDeepScanWorkerFailureCases({ manifest.findings.map((finding) => finding.provenance.candidateId), ["candidate-1"], ); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); } async function testNonRetryableReducerAbortsScanWithoutRetry( @@ -236,11 +275,12 @@ export function createDeepScanWorkerFailureCases({ assert.deepEqual(attempts, [undefined]); assert.deepEqual(sleeps, []); assert.equal(normalExecutor.runningDiscovery, 0); - assert.equal(store.failCalls, 1); + assert.equal(store.failureInputs.length, 1); assert.equal(store.dedupCommits.length, 0); } return { + testResumeRequiresHistoricalWorkerPrompt, testRecoverableWorkerErrorsCannotFailScan, testPolicyRefusedReducerPreservesInputsAndCommittedAggregate, testNonRetryableReducerAbortsScanWithoutRetry, diff --git a/plugins/codex-security/mcp-app/tests/support/streams.mjs b/plugins/codex-security/mcp-app/tests/support/streams.mjs new file mode 100644 index 0000000000..03d7ba9af1 --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/support/streams.mjs @@ -0,0 +1,94 @@ +import { spawn } from "node:child_process"; +import { setTimeout as delay } from "node:timers/promises"; + +function consumeStreamLines(stream, consume) { + let buffer = ""; + stream.setEncoding("utf8"); + stream.on("data", (chunk) => { + buffer += chunk; + const lines = buffer.split("\n"); + buffer = lines.pop(); + for (const line of lines) { + const trimmed = line.trim(); + if (trimmed) consume(trimmed); + } + }); +} + +export function startServer( + serverPath, + env, + { cwd, component, withTimeout, responseLabel, stderrLines, checkSignalCode }, +) { + const child = spawn(process.execPath, [serverPath, "--stdio"], { + cwd, + env, + stdio: ["pipe", "pipe", "pipe"], + }); + const responses = new Map(); + const waiters = new Map(); + const stderrEvents = []; + consumeStreamLines(child.stdout, (line) => { + const response = JSON.parse(line); + responses.set(response.id, response); + waiters.get(response.id)?.(response); + waiters.delete(response.id); + }); + consumeStreamLines(child.stderr, (line) => { + stderrLines?.push(line); + try { + const event = JSON.parse(line); + if (event.component === component) stderrEvents.push(event); + } catch { + // Callers choose whether to retain non-structured diagnostics. + } + }); + return { + pid: child.pid, + notify(method, params = {}) { + writeMessage(child, { jsonrpc: "2.0", method, params }); + }, + sendRequest(id, method, params = {}) { + writeMessage(child, { jsonrpc: "2.0", id, method, params }); + }, + request(id, method, params = {}) { + this.sendRequest(id, method, params); + return this.waitForResponse(id); + }, + waitForResponse(id, timeoutMs = 15_000) { + const existing = responses.get(id); + if (existing) return Promise.resolve(existing); + return withTimeout( + new Promise((resolve) => waiters.set(id, resolve)), + timeoutMs, + `${responseLabel} ${id}`, + ); + }, + stderrEvents() { + return [...stderrEvents]; + }, + stderrText() { + return stderrLines.join("\n"); + }, + response(id) { + return responses.get(id); + }, + async stop() { + if ( + child.exitCode !== null || + (checkSignalCode && child.signalCode !== null) + ) + return; + child.stdin.end(); + const exited = new Promise((resolve) => child.once("exit", resolve)); + await Promise.race([exited, delay(2_000)]); + if (child.exitCode === null && child.signalCode === null) + child.kill("SIGKILL"); + await exited; + }, + }; +} + +export function writeMessage(child, message) { + child.stdin.write(`${JSON.stringify(message)}\n`); +} diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_artifact_validation.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_artifact_validation.mjs index 9e3e66e938..71532a760c 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_artifact_validation.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_artifact_validation.mjs @@ -3,17 +3,9 @@ import { scanId, workerDraft as draft, } from "./scan-draft-fixture.mjs"; +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; -import { - mkdir, - mkdtemp, - readFile, - realpath, - rm, - symlink, - writeFile, -} from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, rm, symlink, writeFile } from "node:fs/promises"; import path from "node:path"; import { importSource } from "./import-module.mjs"; @@ -24,9 +16,7 @@ const { validateDiscoveryArtifacts, validateReducerArtifacts } = ); const otherScanId = "12c17317-9594-49e0-b06a-d72fd7e14bba"; -const root = await realpath( - await mkdtemp(path.join(tmpdir(), "deep-scan-artifact-validation-")), -); +const root = await temporaryDirectory("deep-scan-artifact-validation-", true); try { await testDiscoveryValidation(root); await testReducerValidation(root); 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 dfeb5ddd2a..7bbda94c3f 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 @@ -1,9 +1,12 @@ +import { sourceReferences } from "./support/source-references.mjs"; import { deferred } from "./deferred.mjs"; +import { createTemporaryDirectories } from "./support/temporary-directories.mjs"; +import { once } from "node:events"; +import { mock } from "node:test"; import assert from "node:assert/strict"; import { createHash, randomUUID } from "node:crypto"; import { mkdir, - mkdtemp, readFile, readdir, realpath, @@ -11,7 +14,6 @@ import { rm, writeFile, } from "node:fs/promises"; -import { tmpdir } from "node:os"; import path from "node:path"; import { importModule } from "./import-module.mjs"; import { testDeepScanDeadlines } from "./deep_scan_deadline_cases.mjs"; @@ -23,7 +25,7 @@ const { DeepScanCoordinatorRegistry, DeepScanNonRetryableError, DeepScanRemoteCoordinator, - DeepScanStartLock, + AsyncLock, classifyCodexWorkerError, startOrJoinDeepScanCoordinator, } = await importModule({ @@ -36,7 +38,19 @@ const { }, loader: { ".md": "text" }, }); -const temporaryRoots = []; +const claimDedupInputs = (claim) => + claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ + dedupWorkerId: claim.id, + discoveryWorkerId, + inputOrder, + })); + +const recordSleeps = (sleeps) => async (delayMs, signal) => { + assert.equal(signal.aborted, false); + sleeps.push(delayMs); +}; + +const temporaryDirectories = createTemporaryDirectories(true); async function testCappedQueueAndSerialDedup() { const fixture = await fixtureRun({ workers: 3, @@ -202,7 +216,7 @@ async function testDiscoveryWorkersKeepOneContextAfterPersistedUpdate() { maxDiscoveryRuns: 2, }); fixture.run.userContext = "Initial context."; - const firstWorkerGate = deferred(); + const firstWorkerGate = Promise.withResolvers(); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ dedupNewFindings: [1], @@ -210,7 +224,7 @@ async function testDiscoveryWorkersKeepOneContextAfterPersistedUpdate() { }); const coordinator = createCoordinator(fixture, store, executor, {}); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; store.run.userContext = "Updated context."; firstWorkerGate.resolve(); await coordinator.wait(undefined, 5_000); @@ -235,7 +249,7 @@ async function testPersistedContextDoesNotChangeAnotherProcessDiscoverySnapshot( maxDiscoveryRuns: 2, }); fixture.run.userContext = "Initial cross-process context."; - const firstWorkerGate = deferred(); + const firstWorkerGate = Promise.withResolvers(); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ dedupNewFindings: [1], @@ -250,7 +264,7 @@ async function testPersistedContextDoesNotChangeAnotherProcessDiscoverySnapshot( clock: immediateClock, }); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; store.run.userContext = "Updated cross-process context."; firstWorkerGate.resolve(); const terminal = await coordinator.wait(undefined, 5_000); @@ -397,10 +411,7 @@ async function testRetryKeepsLogicalWorker() { random: () => 0.5, clock: { now: () => 1_700_000_000_000, - sleep: async (delayMs, signal) => { - assert.equal(signal.aborted, false); - sleeps.push(delayMs); - }, + sleep: recordSleeps(sleeps), }, }); coordinator.start(); @@ -501,7 +512,7 @@ async function testCompletionOrdering() { maxDiscoveryRuns: 3, }); const store = new FakeStore(fixture.run); - const firstWorkerGate = deferred(); + const firstWorkerGate = Promise.withResolvers(); const executor = new FakeExecutor({ dedupNewFindings: [1, 0], discoveryGates: { "discovery-0001": firstWorkerGate.promise }, @@ -543,7 +554,7 @@ async function testSaturationDrainsBufferedAndCancelsInflight() { }); const coordinator = createCoordinator(fixture, store, executor, {}); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; await eventually( () => executor.discoveryCalls === 6 && executor.runningDiscovery === 2, ); @@ -586,7 +597,7 @@ async function testSaturationPreservesFindingAlreadyBuffered() { maxDiscoveryRuns: 4, }); const store = new FakeStore(fixture.run); - const laterDiscoveries = deferred(); + const laterDiscoveries = Promise.withResolvers(); const executor = new FakeExecutor({ blockDedup: true, discoveryGates: { @@ -603,7 +614,7 @@ async function testSaturationPreservesFindingAlreadyBuffered() { onComplete: async (draft) => completed.push(draft), }); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; laterDiscoveries.resolve(); await eventually( () => @@ -681,7 +692,7 @@ async function testSaturationIgnoresWorkerFailureSettledAfterStop() { coordinator.start(); await Promise.all([ - executor.dedupStarted, + executor.dedupStarted.promise, store.discoveryFailureBlocked.promise, ]); executor.releaseDedup(); @@ -701,7 +712,7 @@ async function testSaturationIgnoresWorkerFailureSettledAfterStop() { ); assert.ok(failedWorker); assert.equal(failedWorker.status, "failed"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(store.finishCalls.length, 1); assert.equal(executor.discoveryCalls, 3); assert.equal(executor.dedupCalls, 1); @@ -723,7 +734,7 @@ async function testSettledReducerIsNotStarvedByDiscoveryBacklog() { const coordinator = createCoordinator(fixture, store, executor, {}); coordinator.start(); - await executor.dedupStarted; + await executor.dedupStarted.promise; await eventually( () => executor.discoveryCalls >= 8 && executor.runningDiscovery === 4, ); @@ -1210,7 +1221,7 @@ async function testFailureManifestWriteDoesNotMaskOriginalError() { assert.equal(terminal?.status, "failed"); assert.match(terminal?.error ?? "", /fixture configuration failure/); assert.equal(terminal?.manifestPath, undefined); - assert.equal(store.failCalls, 1); + assert.equal(store.failureInputs.length, 1); } async function testFinishPersistenceFailureRewritesManifestAsFailure() { @@ -1256,7 +1267,7 @@ async function testLostFinishResponseReplaysWithoutOverwritingSuccessManifest() assert.equal(terminal?.status, "succeeded"); assert.equal(store.finishCalls.length, 2); assert.deepEqual(store.finishCalls[1], store.finishCalls[0]); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); const manifest = JSON.parse(await readFile(terminal.manifestPath, "utf8")); assert.equal(manifest.scan.scanId, fixture.run.scanId); } @@ -1285,7 +1296,7 @@ async function testLostWorkerCommitResponsesReplayIdempotently() { assert.equal(store.dedupCommitResponseLosses, 1); assert.equal(store.dedupCommitCalls.length, 2); assert.equal(store.dedupCommits.length, 1); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); } async function testCommittedReducerIsReconciledBeforeDiscoveryFailureManifest() { @@ -1295,7 +1306,7 @@ async function testCommittedReducerIsReconciledBeforeDiscoveryFailureManifest() stopAfterNoNew: 10, maxDiscoveryRuns: 3, }); - const thirdWorkerGate = deferred(); + const thirdWorkerGate = Promise.withResolvers(); const store = new FakeStore(fixture.run); store.blockDedupCommitResponse = true; const executor = new FakeExecutor({ @@ -1385,7 +1396,7 @@ async function testCancellationClearsRetryWait() { failFirstDiscoveryAttempt: true, blockDiscoveryAfterCalls: 1, }); - const sleepStarted = deferred(); + const sleepStarted = Promise.withResolvers(); const coordinator = createCoordinator(fixture, store, executor, { clock: { now: immediateClock.now, @@ -1498,10 +1509,7 @@ async function testInvalidArtifactsRetry() { retryDelaysMs: [1, 3, 9], clock: { now: immediateClock.now, - sleep: async (delayMs, signal) => { - assert.equal(signal.aborted, false); - sleeps.push(delayMs); - }, + sleep: recordSleeps(sleeps), }, }); coordinator.start(); @@ -1847,7 +1855,7 @@ async function testExhaustedReducerPreservesCommittedArtifacts() { stopAfterConsecutiveErrors: 2, maxDiscoveryRuns: 3, }); - const nextDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ discoveryCandidateId: "candidate-1", @@ -1891,7 +1899,7 @@ async function testCommittedAggregateIsNotSalvagedWhenUntrusted(failure) { stopAfterConsecutiveErrors: 1, maxDiscoveryRuns: 3, }); - const nextDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ discoveryCandidateId: "candidate-1", @@ -1941,7 +1949,7 @@ async function testCancellationAfterCommittedAggregateRemainsCanceled() { stopAfterNoNew: 99, maxDiscoveryRuns: 3, }); - const nextDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); const store = new FakeStore(fixture.run); const completedDrafts = []; const coordinator = createCoordinator( @@ -2020,7 +2028,7 @@ async function testRejectedStaleReducerCommitPreservesReplacementCandidates() { stopAfterNoNew: 99, maxDiscoveryRuns: 3, }); - const nextDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); const store = new FakeStore(fixture.run); store.failDedupCommitFromCall = 2; store.replacementCandidatesBeforeDedupRejection = JSON.stringify( @@ -2095,7 +2103,7 @@ async function testAmbiguousReducerCommitPreservesPublishedCandidates() { stopAfterNoNew: 99, maxDiscoveryRuns: 3, }); - const nextDiscovery = deferred(); + const nextDiscovery = Promise.withResolvers(); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ discoveryCandidateId: "candidate-1", @@ -2231,7 +2239,7 @@ async function testWaiterDetachAndCancellation() { const executor = new FakeExecutor({ blockDiscovery: true }); const coordinator = createCoordinator(fixture, store, executor, {}); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; assert.equal(coordinator.snapshot().dispatchedCount, 2); const waiterAbort = new AbortController(); @@ -2306,8 +2314,8 @@ async function testCancellationDuringDiscoveryAcceptanceRejectsLateSuccess() { store.blockDiscoverySuccess = true; const executor = new FakeExecutor({ dedupNewFindings: [0] }); const events = []; - const preserved = deferred(); - const releasePreservation = deferred(); + const preserved = Promise.withResolvers(); + const releasePreservation = Promise.withResolvers(); const coordinator = createCoordinator(fixture, store, executor, { log: (event) => events.push(event), threadId: "checkpoint-owner", @@ -2414,23 +2422,19 @@ async function testRegistryEvictionAndExternalFailure() { pluginRoot: failedFixture.pluginRoot, clock: immediateClock, }); - await failedExecutor.discoveryStarted; + await failedExecutor.discoveryStarted.promise; failedStore.run.status = "failed"; failedStore.run.error = "failure already persisted by fail-scan"; - assert.equal( - registry.failExternallyPersisted( - failedFixture.run.scanId, - failedStore.run.error, - ), - true, - ); + registry + .get(failedFixture.run.scanId) + ?.failExternallyPersisted(failedStore.run.error); const terminal = await failed.wait(undefined, 5_000); assert.equal(terminal?.status, "failed"); assert.equal(terminal?.error, "failure already persisted by fail-scan"); await eventually(() => failedExecutor.runningDiscovery === 0); await eventually(() => registry.get(failedFixture.run.scanId) === undefined); assert.equal( - failedStore.failCalls, + failedStore.failureInputs.length, 0, "an externally persisted failure must not be persisted again", ); @@ -2456,7 +2460,7 @@ async function testStoppedPublicationFailurePreservesOriginalDiagnostic() { }, }); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; store.run.status = "failed"; store.run.error = `authoritative worker failure diagnostic ${"x".repeat(2_350)}`; coordinator.failExternallyPersisted( @@ -2509,7 +2513,7 @@ async function testStoppedPublicationFailureBoundsPrefixedDiagnostic() { }, }); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; store.run.status = "canceled"; coordinator.cancel("fixture persisted cancellation"); @@ -2534,15 +2538,13 @@ async function testTerminalReadFailureIsNotRecordedAsPublicationFailure() { }); const store = new FakeStore(fixture.run); const executor = new FakeExecutor({ blockDiscovery: true }); - let publicationAttempts = 0; + const publicationAttempts = mock.fn(async () => {}); const coordinator = createCoordinator(fixture, store, executor, { threadId: "checkpoint-owner", - onStopped: async () => { - publicationAttempts += 1; - }, + onStopped: publicationAttempts, }); coordinator.start(); - await executor.discoveryStarted; + await executor.discoveryStarted.promise; store.run.status = "failed"; store.run.error = "original worker failure diagnostic"; store.failNextTerminalGet = true; @@ -2551,7 +2553,7 @@ async function testTerminalReadFailureIsNotRecordedAsPublicationFailure() { const terminal = await coordinator.wait(undefined, 5_000); assert.equal(terminal?.error, "original worker failure diagnostic"); - assert.equal(publicationAttempts, 0); + assert.equal(publicationAttempts.mock.callCount(), 0); assert.deepEqual(store.publicationFailureMessages, []); } @@ -2568,11 +2570,9 @@ async function testCoordinatorHeartbeatsStopAfterOwnershipChanges() { updatedAt: new Date().toISOString(), }; const store = new FakeStore(run); - let heartbeats = 0; - store.heartbeatCoordinator = async () => { - heartbeats += 1; + store.heartbeatCoordinator = mock.fn(async () => { return structuredClone(store.run); - }; + }); let reads = 0; store.get = async () => { reads += 1; @@ -2604,7 +2604,7 @@ async function testCoordinatorHeartbeatsStopAfterOwnershipChanges() { await new Promise((resolve) => setTimeout(resolve, 25)); assert.equal( - heartbeats, + store.heartbeatCoordinator.mock.callCount(), 3, "heartbeat writes continue until a newer generation is confirmed", ); @@ -2616,7 +2616,7 @@ async function testCoordinatorHeartbeatsStopAfterOwnershipChanges() { assert.equal(current?.status, "succeeded"); assert.equal(current?.coordinatorGeneration, 3); assert.equal( - store.failCalls, + store.failureInputs.length, 0, "a stale coordinator must never fail the new owner", ); @@ -2631,7 +2631,7 @@ async function testCoordinatorHeartbeatsContinueDuringBlockedOwnershipRead() { }); const run = { ...fixture.run, coordinatorGeneration: 2 }; const store = new FakeStore(run); - const ownershipRead = deferred(); + const ownershipRead = Promise.withResolvers(); let heartbeats = 0; let reads = 0; store.heartbeatCoordinator = async () => { @@ -2815,7 +2815,7 @@ async function testStaleMutationObservesReplacement() { assert.equal(terminal?.status, "succeeded"); assert.equal(terminal?.coordinatorGeneration, 3); assert.equal( - store.failCalls, + store.failureInputs.length, 0, "a fenced mutation must observe rather than fail the replacement", ); @@ -2834,30 +2834,25 @@ async function testJoinAndOrphanRules() { pluginRoot: fixture.pluginRoot, threadId: "thread-fixture", }; - let starts = 0; - let failures = 0; + const starts = mock.fn(() => existingCoordinator); + const failures = mock.fn(async () => {}); const existing = await startOrJoinDeepScanCoordinator({ run: { ...fixture.run, persistedWorkerCount: 3 }, registry: { get: () => existingCoordinator, - start: () => { - starts += 1; - return existingCoordinator; - }, + start: starts, }, options: { ...defaults, store: { - fail: async () => { - failures += 1; - }, + fail: failures, }, }, }); assert.equal(existing.coordinator, existingCoordinator); assert.equal(existing.joined, true); - assert.equal(starts, 0); - assert.equal(failures, 0); + assert.equal(starts.mock.callCount(), 0); + assert.equal(failures.mock.callCount(), 0); const running = { ...fixture.run, @@ -2868,19 +2863,14 @@ async function testJoinAndOrphanRules() { run: running, registry: { get: () => undefined, - start: () => { - starts += 1; - return existingCoordinator; - }, + start: starts, }, options: { ...defaults, store: { claimCoordinator: async () => ({ run: running, acquired: false }), get: async () => ({ ...running, status: "succeeded" }), - fail: async () => { - failures += 1; - }, + fail: failures, }, }, }); @@ -2890,13 +2880,19 @@ async function testJoinAndOrphanRules() { "succeeded", ); assert.equal( - starts, + starts.mock.callCount(), 0, "another process must not create a duplicate coordinator", ); - assert.equal(failures, 0, "another process must not interrupt the live scan"); + assert.equal( + failures.mock.callCount(), + 0, + "another process must not interrupt the live scan", + ); - let staleClaims = 0; + const staleClaims = mock.fn(async () => { + return { run: staleRunning, acquired: false }; + }); const staleRunning = { ...running, updatedAt: "2026-01-01T00:00:00Z" }; const staleObserver = await startOrJoinDeepScanCoordinator({ run: staleRunning, @@ -2904,10 +2900,7 @@ async function testJoinAndOrphanRules() { options: { ...defaults, store: { - claimCoordinator: async () => { - staleClaims += 1; - return { run: staleRunning, acquired: false }; - }, + claimCoordinator: staleClaims, get: async () => staleRunning, }, }, @@ -2917,7 +2910,7 @@ async function testJoinAndOrphanRules() { undefined, ); assert.equal( - staleClaims, + staleClaims.mock.callCount(), 1, "a confirmed live lease must not be reclaimed on every poll", ); @@ -2927,7 +2920,7 @@ async function testJoinAndOrphanRules() { registry: { get: () => undefined, start: (options) => { - starts += 1; + starts(); assert.equal(options.run.coordinatorGeneration, 2); return existingCoordinator; }, @@ -2939,29 +2932,31 @@ async function testJoinAndOrphanRules() { acquired: true, run: { ...fixture.run, coordinatorGeneration: 2 }, }), - fail: async () => { - failures += 1; - }, + fail: failures, }, }, }); assert.equal(recovered.coordinator, existingCoordinator); assert.equal(recovered.joined, false); - assert.equal(starts, 1, "only an expired coordinator may be adopted"); assert.equal( - failures, + starts.mock.callCount(), + 1, + "only an expired coordinator may be adopted", + ); + assert.equal( + failures.mock.callCount(), 0, "recovering an orphan must not fail the logical scan", ); - const lock = new DeepScanStartLock(); - const firstGate = deferred(); + const lock = new AsyncLock(); + const firstGate = Promise.withResolvers(); let liveCoordinator; - starts = 0; + starts.mock.resetCalls(); const registry = { get: () => liveCoordinator, start: () => { - starts += 1; + starts(); liveCoordinator = { marker: "concurrent" }; return liveCoordinator; }, @@ -2970,9 +2965,7 @@ async function testJoinAndOrphanRules() { ...defaults, store: { claimCoordinator: async () => ({ acquired: true, run: fixture.run }), - fail: async () => { - failures += 1; - }, + fail: failures, }, }; const first = lock.run(async () => { @@ -2993,7 +2986,7 @@ async function testJoinAndOrphanRules() { ); await new Promise((resolve) => setImmediate(resolve)); assert.equal( - starts, + starts.mock.callCount(), 0, "the second caller must wait through begin plus registry start", ); @@ -3002,8 +2995,8 @@ async function testJoinAndOrphanRules() { assert.equal(created.joined, false); assert.equal(joined.joined, true); assert.equal(created.coordinator, joined.coordinator); - assert.equal(starts, 1); - assert.equal(failures, 0); + assert.equal(starts.mock.callCount(), 1); + assert.equal(failures.mock.callCount(), 0); } async function testPausedDiscoverySurvivesCoordinatorRestart() { @@ -3064,7 +3057,7 @@ async function testPausedDiscoverySurvivesCoordinatorRestart() { 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.failureInputs.length, 0); assert.equal(store.run.manifestPath, undefined); await assert.rejects( readFile( @@ -3114,7 +3107,7 @@ async function testPausedDiscoverySurvivesCoordinatorRestart() { assert.equal(continuationClaims.length, 1); assert.equal(continuationClaims[0].handoffClaimToken, handoffClaimToken); assert.equal(terminal?.status, "succeeded"); - assert.equal(store.failCalls, 0); + assert.equal(store.failureInputs.length, 0); assert.equal(replacementExecutor.logicalDiscoveryWorkers.size, 1); assert.equal( replacementExecutor.logicalDiscoveryWorkers.has( @@ -3203,12 +3196,12 @@ async function testResumedDiscoveryDeadlineUsesPersistedCreationTime( clock, }); resumed.start(); - if (!alreadyExpired) await resumedExecutor.discoveryStarted; + if (!alreadyExpired) await resumedExecutor.discoveryStarted.promise; 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(store.failureInputs.length, 0); assert.equal(resumedExecutor.discoveryCalls, alreadyExpired ? 0 : 1); assert.equal(resumedExecutor.dedupCalls, 1); assert.equal(resumedExecutor.runningDiscovery, 0); @@ -3292,13 +3285,7 @@ async function testResumedManifestPreservesCompletedReducer( persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker), ), - persistedDedupInputs: store.dedupClaims.flatMap((claim) => - claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ - dedupWorkerId: claim.id, - discoveryWorkerId, - inputOrder, - })), - ), + persistedDedupInputs: store.dedupClaims.flatMap(claimDedupInputs), }; const replacement = createCoordinator(fixture, store, new FakeExecutor(), { run: store.run, @@ -3319,9 +3306,7 @@ async function testResumedManifestPreservesCompletedReducer( ); } -async function testResumeUsesHistoricalCandidateSnapshotForEachReducer( - legacyLayout = false, -) { +async function testResumeUsesHistoricalCandidateSnapshotForEachReducer() { const fixture = await fixtureRun({ workers: 3, subagents: 0, @@ -3351,16 +3336,12 @@ async function testResumeUsesHistoricalCandidateSnapshotForEachReducer( persistedWorkers: [...store.workers.values()].map((worker) => structuredClone(worker), ), - persistedDedupInputs: store.dedupClaims.flatMap((claim) => - claim.workerIds.map((discoveryWorkerId, inputOrder) => ({ - dedupWorkerId: claim.id, - discoveryWorkerId, - inputOrder, - })), - ), + persistedDedupInputs: store.dedupClaims.flatMap(claimDedupInputs), }; + const completedDrafts = []; const replacement = createCoordinator(fixture, store, new FakeExecutor(), { run: store.run, + onComplete: async (draft) => completedDrafts.push(structuredClone(draft)), }); replacement.start(); @@ -3390,6 +3371,10 @@ async function testResumeUsesHistoricalCandidateSnapshotForEachReducer( latestResult.findings[0]?.rootCause.summary, "final reducer evidence", ); + assert.deepEqual( + completedDrafts.map((draft) => draft.findings), + [latestResult.findings], + ); } async function testPersistedErrorLimitStopsBeforeRescheduling() { @@ -3500,10 +3485,9 @@ async function testPersistedReducerErrorLimitStopsBeforeRescheduling() { } async function fixtureRun(config) { - const root = await realpath( - await mkdtemp(path.join(tmpdir(), "codex-security-deep-coordinator-")), + const root = await temporaryDirectories.create( + "codex-security-deep-coordinator-", ); - temporaryRoots.push(root); const targetPath = path.join(root, "target"); const scanDir = path.join(root, "scan"); const pluginRoot = path.join(root, "plugin"); @@ -3545,7 +3529,6 @@ class FakeStore { failDedupCommitFromCall = undefined; loseEveryDedupCommitResponseAfterCommit = false; progressCalls = 0; - failCalls = 0; failureInputs = []; finishCalls = []; failFinish = false; @@ -3558,15 +3541,15 @@ class FakeStore { loseFirstDedupCommitResponseAfterCommit = false; dedupCommitResponseLosses = 0; blockDedupCommitResponse = false; - dedupCommitPersisted = deferred(); - dedupCommitResponseGate = deferred(); + dedupCommitPersisted = Promise.withResolvers(); + dedupCommitResponseGate = Promise.withResolvers(); blockDiscoverySuccess = false; - discoverySuccessBlocked = deferred(); - discoverySuccessGate = deferred(); + discoverySuccessBlocked = Promise.withResolvers(); + discoverySuccessGate = Promise.withResolvers(); blockDiscoveryFailure = false; - discoveryFailureBlocked = deferred(); - discoveryFailureGate = deferred(); - dedupCommitted = deferred(); + discoveryFailureBlocked = Promise.withResolvers(); + discoveryFailureGate = Promise.withResolvers(); + dedupCommitted = Promise.withResolvers(); failNextTerminalGet = false; publicationFailureMessages = []; @@ -3835,7 +3818,6 @@ class FakeStore { manifestPath, stagedManifestPath, ) { - this.failCalls += 1; this.failureInputs.push({ message, status, manifestPath }); if (this.rejectFailurePersistence) { throw new DeepScanNonRetryableError( @@ -3888,15 +3870,11 @@ class FakeStore { class FakeExecutor { constructor(options = {}) { this.options = options; - this.discoveryStarted = new Promise((resolve) => { - this.resolveDiscoveryStarted = resolve; - }); - this.dedupStarted = new Promise((resolve) => { - this.resolveDedupStarted = resolve; - }); - this.dedupGate = deferred(); - this.dedupArtifactsWritten = deferred(); - this.discoveryArtifactsWritten = deferred(); + this.discoveryStarted = Promise.withResolvers(); + this.dedupStarted = Promise.withResolvers(); + this.dedupGate = Promise.withResolvers(); + this.dedupArtifactsWritten = Promise.withResolvers(); + this.discoveryArtifactsWritten = Promise.withResolvers(); } discoveryCalls = 0; @@ -3952,7 +3930,7 @@ class FakeExecutor { this.maximumDiscoveryConcurrency, this.runningDiscovery, ); - this.resolveDiscoveryStarted(); + this.discoveryStarted.resolve(); try { if (this.options.policyRefusalWorkers?.includes(workerId)) { const message = @@ -4050,7 +4028,7 @@ class FakeExecutor { this.maximumDedupConcurrency, this.runningDedup, ); - this.resolveDedupStarted(); + this.dedupStarted.resolve(); try { if (this.options.blockDedup) await this.dedupGate.promise; if (request.signal.aborted) throw abortError(); @@ -4153,13 +4131,7 @@ async function writeDedupArtifacts(request, consumedOverride, options = {}) { const workerDrafts = await Promise.all( reducer.claimedWorkers.map(async (worker) => { const result = JSON.parse(await readFile(worker.resultPath, "utf8")); - result.findings = result.findings.map((finding, index) => ({ - ...finding, - provenance: { - ...finding.provenance, - sourceFindingIds: [`${worker.id}:${index}`], - }, - })); + result.findings = result.findings.map(sourceReferences(worker)); return result; }), ); @@ -4260,11 +4232,8 @@ function standardScanDraft(scanId, candidateId, workerLabel) { async function waitForAbort(signal) { if (signal.aborted) throw abortError(); - await new Promise((_, reject) => { - signal.addEventListener("abort", () => reject(abortError()), { - once: true, - }); - }); + await once(signal, "abort"); + throw abortError(); } function abortError() { @@ -4322,6 +4291,7 @@ function boundedFixtureErrorText(message, maximum) { } const { + testResumeRequiresHistoricalWorkerPrompt, testRecoverableWorkerErrorsCannotFailScan, testPolicyRefusedReducerPreservesInputsAndCommittedAggregate, testNonRetryableReducerAbortsScanWithoutRetry, @@ -4330,6 +4300,7 @@ const { FakeStore, FakeExecutor, createCoordinator, + DeepScanCoordinator, DeepScanNonRetryableError, classifyCodexWorkerError, deferred, @@ -4433,6 +4404,8 @@ try { await testCoordinatorHeartbeatsContinueDuringBlockedOwnershipRead(); await testRemoteObserverRetriesTransientPersistenceFailures(); await testJoinAndOrphanRules(); + await testResumeRequiresHistoricalWorkerPrompt("failed"); + await testResumeRequiresHistoricalWorkerPrompt("canceled"); await testPausedDiscoverySurvivesCoordinatorRestart(); await testResumedDiscoveryDeadlineUsesPersistedCreationTime(); await testResumedDiscoveryDeadlineUsesPersistedCreationTime(true); @@ -4441,13 +4414,10 @@ try { await testResumedManifestPreservesCompletedReducer(); await testResumedManifestPreservesCompletedReducer(true); await testResumeUsesHistoricalCandidateSnapshotForEachReducer(); - await testResumeUsesHistoricalCandidateSnapshotForEachReducer(true); await testPersistedErrorLimitStopsBeforeRescheduling(); await testPersistedReducerErrorLimitStopsBeforeRescheduling(); } finally { - await Promise.all( - temporaryRoots.map((root) => rm(root, { recursive: true, force: true })), - ); + await temporaryDirectories.cleanup(); } function createCoordinator(fixture, store, executor, options) { diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_stdio_lifecycle.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_stdio_lifecycle.mjs index 89e3f9e567..eb25fd2c05 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_stdio_lifecycle.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_stdio_lifecycle.mjs @@ -1,25 +1,24 @@ import { assertNoError, assertFlagPair } from "./assertions.mjs"; import { readOnlyParentSandboxState } from "./sandbox-state.mjs"; +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; -import { execFile, spawn } from "node:child_process"; +import { execFile } from "node:child_process"; import { randomUUID } from "node:crypto"; import { chmod, mkdir, - mkdtemp, readFile, realpath, rm, writeFile, } from "node:fs/promises"; -import { tmpdir } from "node:os"; import path from "node:path"; import { setTimeout as delay } from "node:timers/promises"; import { promisify } from "node:util"; import { pathToFileURL } from "node:url"; import { applicationRoot as mcpAppRoot, buildServer } from "./build-server.mjs"; -import { consumeLines } from "./consume-lines.mjs"; +import * as streams from "./support/streams.mjs"; const execFileAsync = promisify(execFile); @@ -36,9 +35,7 @@ if (process.platform === "win32") { } async function testDeepScanStdioLifecycle() { - const fixtureRoot = await mkdtemp( - path.join(tmpdir(), "codex-security-deep-stdio-"), - ); + const fixtureRoot = await temporaryDirectory("codex-security-deep-stdio-"); const targetPath = path.join(fixtureRoot, "target"); const failedTargetPath = path.join(fixtureRoot, "failed-target"); const stateDir = path.join(fixtureRoot, "state"); @@ -898,89 +895,15 @@ async function testDeepScanStdioLifecycle() { } function startServer(serverPath, env) { - const child = spawn(process.execPath, [serverPath, "--stdio"], { + return streams.startServer(serverPath, env, { cwd: pluginRoot, - env, - stdio: ["pipe", "pipe", "pipe"], - }); - const responses = new Map(); - const waiters = new Map(); - const stderrEvents = []; - const stderrLines = []; - let stdoutBuffer = ""; - let stderrBuffer = ""; - - child.stdout.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { - stdoutBuffer += chunk; - stdoutBuffer = consumeLines(stdoutBuffer, (line) => { - const response = JSON.parse(line); - responses.set(response.id, response); - waiters.get(response.id)?.(response); - waiters.delete(response.id); - }); + component: "codex_security_deep_scan", + withTimeout, + responseLabel: "JSON-RPC response", + // Non-structured diagnostics remain available in the child process on test failure. + stderrLines: [], + checkSignalCode: true, }); - child.stderr.setEncoding("utf8"); - child.stderr.on("data", (chunk) => { - stderrBuffer += chunk; - stderrBuffer = consumeLines(stderrBuffer, (line) => { - stderrLines.push(line); - try { - const event = JSON.parse(line); - if (event.component === "codex_security_deep_scan") - stderrEvents.push(event); - } catch { - // Non-structured diagnostics remain available in the child process on test failure. - } - }); - }); - - return { - pid: child.pid, - notify(method, params = {}) { - child.stdin.write( - `${JSON.stringify({ jsonrpc: "2.0", method, params })}\n`, - ); - }, - sendRequest(id, method, params = {}) { - child.stdin.write( - `${JSON.stringify({ jsonrpc: "2.0", id, method, params })}\n`, - ); - }, - request(id, method, params = {}) { - this.sendRequest(id, method, params); - return this.waitForResponse(id); - }, - waitForResponse(id, timeoutMs = 15_000) { - const existing = responses.get(id); - if (existing) return Promise.resolve(existing); - return withTimeout( - new Promise((resolve) => waiters.set(id, resolve)), - timeoutMs, - `JSON-RPC response ${id}`, - ); - }, - stderrEvents() { - return [...stderrEvents]; - }, - stderrText() { - return stderrLines.join("\n"); - }, - response(id) { - return responses.get(id); - }, - async stop() { - if (child.exitCode !== null || child.signalCode !== null) return; - child.stdin.end(); - const exited = new Promise((resolve) => child.once("exit", resolve)); - const graceful = await Promise.race([ - exited.then(() => true), - delay(2_000).then(() => false), - ]); - if (!graceful && child.exitCode === null) child.kill("SIGKILL"); - await exited; - }, - }; } function toolCall( diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_store.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_store.mjs index 59fea31dc1..e61cd47e0f 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_store.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_store.mjs @@ -1,4 +1,5 @@ import { deferred } from "./deferred.mjs"; +import { mock } from "node:test"; import assert from "node:assert/strict"; import { randomUUID } from "node:crypto"; import { mkdtemp, readFile, rm } from "node:fs/promises"; @@ -104,7 +105,7 @@ async function testBeginProtocolAndParsing() { async function testWriteSerializationAndRecovery() { const calls = []; - const firstGate = deferred(); + const firstGate = Promise.withResolvers(); let first = true; const runner = async (args) => { calls.push(args); @@ -126,7 +127,10 @@ async function testWriteSerializationAndRecovery() { scanId: randomUUID(), reviewItemsCompleted: 1, }); - const cancellation = store.cancel(randomUUID(), "thread-fixture"); + const finalWrite = store.updateProgress({ + scanId: randomUUID(), + reviewItemsCompleted: 2, + }); await new Promise((resolve) => setImmediate(resolve)); assert.equal( calls.length, @@ -135,7 +139,7 @@ async function testWriteSerializationAndRecovery() { ); firstGate.resolve(); await assert.rejects(firstWrite, /first write failed/); - await Promise.all([secondWrite, cancellation]); + await Promise.all([secondWrite, finalWrite]); assert.equal( calls.length, 3, @@ -144,17 +148,17 @@ async function testWriteSerializationAndRecovery() { assert.equal(calls[0][0], "update-progress"); assert.equal(flagValue(calls[0], "--claim-token"), claimToken); assert.equal(calls[1][0], "update-progress"); - assert.equal( - calls[2][0], - "cancel-scan", - "cancellation must use the same ordered persistence queue", + assert.deepEqual( + calls.slice(1).map((args) => flagValue(args, "--review-items-completed")), + ["1", "2"], + "later writes must retain invocation order through the persistence queue", ); } async function testBeginUsesTheWriteQueue() { const scanId = randomUUID(); const calls = []; - const firstGate = deferred(); + const firstGate = Promise.withResolvers(); const runner = async (args) => { calls.push(args); if (args[0] === "update-progress") { @@ -185,7 +189,7 @@ async function testHeartbeatBypassesBlockedWriteQueue() { const scanId = randomUUID(); const handoffClaimToken = randomUUID(); const scanDir = await mkdtemp(join(tmpdir(), "deep-scan-heartbeat-")); - const blockedWriteGate = deferred(); + const blockedWriteGate = Promise.withResolvers(); const calls = []; const runner = async (args) => { calls.push(args); @@ -378,11 +382,10 @@ async function testCanonicalCommitProtocol() { const reducerId = randomUUID(); const resultManifestPath = "/fixture/scans/run/artifacts/deep_discovery/dedup/result.json"; - let command; - const store = new WorkbenchDeepScanStore(async (args) => { - command = args; + const runWorkbench = mock.fn(async (args) => { return stateResult(scanId, { deepScan: { canonicalArtifacts: canonical } }); }); + const store = new WorkbenchDeepScanStore(runWorkbench); const state = await store.commitDedup({ id: reducerId, scanId, @@ -390,7 +393,7 @@ async function testCanonicalCommitProtocol() { resultManifestPath, }); assert.deepEqual(state.canonicalArtifacts, canonical); - assert.deepEqual(command, [ + assert.deepEqual(runWorkbench.mock.calls.at(-1)?.arguments[0], [ "commit-deep-scan-dedup", "--scan-id", scanId, @@ -534,12 +537,12 @@ async function testPersistenceRetriesRemainInsideTheWriteQueue() { promptPath: "/fixture/reducer/prompt.md", artifactDir: "/fixture/reducer/output", }); - const cancellation = store.cancel(scanId, "thread-fixture"); - await Promise.all([claim, cancellation]); + const progress = store.updateProgress({ scanId, phase: "discovery" }); + await Promise.all([claim, progress]); assert.deepEqual( calls, - ["claim-deep-scan-dedup", "claim-deep-scan-dedup", "cancel-scan"], + ["claim-deep-scan-dedup", "claim-deep-scan-dedup", "update-progress"], "later mutations must not interleave with an idempotent persistence replay", ); } @@ -656,17 +659,14 @@ async function testDeterministicPersistenceFailuresAreNotRetried() { code: "ABORT_ERR", name: "AbortError", }); - let attempts = 0; - const store = new WorkbenchDeepScanStore(async () => { - attempts += 1; - throw canceled; - }); + const attempts = mock.fn(Promise.reject.bind(Promise, canceled)); + const store = new WorkbenchDeepScanStore(attempts); await assert.rejects( idempotentPersistenceScenarios()[1].invoke(store), (error) => error === canceled, ); assert.equal( - attempts, + attempts.mock.callCount(), 1, "explicit cancellation must never be treated as a transient timeout", ); @@ -690,11 +690,8 @@ async function testDeterministicPersistenceFailuresAreNotRetried() { name: "AbortError", }), ]) { - let workerAttempts = 0; - const workerStore = new WorkbenchDeepScanStore(async () => { - workerAttempts += 1; - throw failure; - }); + const workerAttempts = mock.fn(Promise.reject.bind(Promise, failure)); + const workerStore = new WorkbenchDeepScanStore(workerAttempts); await assert.rejects( workerStore.updateWorker({ @@ -709,7 +706,7 @@ async function testDeterministicPersistenceFailuresAreNotRetried() { (error) => error === failure, ); assert.equal( - workerAttempts, + workerAttempts.mock.callCount(), 1, `${status} updates must not replay ${failure.code}`, ); @@ -718,30 +715,25 @@ async function testDeterministicPersistenceFailuresAreNotRetried() { } async function testNonIdempotentMutationsAreNotRetried() { - for (const operation of ["progress", "cancel", "fail", "begin"]) { - let attempts = 0; + for (const operation of ["progress", "fail", "begin"]) { const expected = new Error("sqlite3.OperationalError: database is locked"); - const store = new WorkbenchDeepScanStore(async () => { - attempts += 1; - throw expected; - }); + const attempts = mock.fn(Promise.reject.bind(Promise, expected)); + const store = new WorkbenchDeepScanStore(attempts); const scanId = randomUUID(); const request = operation === "progress" ? store.updateProgress({ scanId, phase: "discovery" }) - : operation === "cancel" - ? store.cancel(scanId, "thread-fixture") - : operation === "fail" - ? store.fail(scanId, "fixture failure") - : store.begin({ - targetPath: "/fixture/repository", - threadId: "thread-fixture", - scanRoot: "/fixture/scans", - }); + : operation === "fail" + ? store.fail(scanId, "fixture failure") + : store.begin({ + targetPath: "/fixture/repository", + threadId: "thread-fixture", + scanRoot: "/fixture/scans", + }); await assert.rejects(request, (error) => error === expected); assert.equal( - attempts, + attempts.mock.callCount(), 1, `${operation} must not gain a new persistence retry policy`, ); diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_store_integration.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_store_integration.mjs index e234d6a818..eae9758e84 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_store_integration.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_store_integration.mjs @@ -1,9 +1,9 @@ +import { temporaryDirectory } from "./support/temporary-directories.mjs"; import assert from "node:assert/strict"; import { execFile, spawn } from "node:child_process"; import { randomUUID } from "node:crypto"; import { once } from "node:events"; -import { mkdir, mkdtemp, readFile, rm, writeFile } from "node:fs/promises"; -import { tmpdir } from "node:os"; +import { mkdir, readFile, rm, writeFile } from "node:fs/promises"; import path from "node:path"; import { promisify } from "node:util"; import { fileURLToPath } from "node:url"; @@ -40,7 +40,7 @@ await testNoopStoppedRefreshRetainsPublicationFailure(); await testConcurrentParentDraftsPreserveBothCheckpoints(); async function createWorkbenchFixture(prefix, homeName = "home") { - const fixtureRoot = await mkdtemp(path.join(tmpdir(), prefix)); + const fixtureRoot = await temporaryDirectory(prefix); const targetPath = path.join(fixtureRoot, "target"); const environment = { ...process.env, @@ -201,14 +201,11 @@ async function testConcurrentParentDraftsPreserveBothCheckpoints() { const python = process.env.PYTHON?.trim() || "python3"; const rawRunWorkbench = createWorkbenchRunner(python, environment); let stagedWrites = 0; - let releaseInitialWrites; - const initialWritesReady = new Promise((resolve) => { - releaseInitialWrites = resolve; - }); + const initialWritesReady = Promise.withResolvers(); const runWorkbench = async (args) => { if (args[0] === "write-scan-draft" && ++stagedWrites <= 2) { - if (stagedWrites === 2) releaseInitialWrites(); - await initialWritesReady; + if (stagedWrites === 2) initialWritesReady.resolve(); + await initialWritesReady.promise; } return rawRunWorkbench(args); }; diff --git a/plugins/codex-security/mcp-app/tests/test_workbench_state_fallback.mjs b/plugins/codex-security/mcp-app/tests/test_workbench_state_fallback.mjs index cbc93c3514..ff202b81c1 100644 --- a/plugins/codex-security/mcp-app/tests/test_workbench_state_fallback.mjs +++ b/plugins/codex-security/mcp-app/tests/test_workbench_state_fallback.mjs @@ -1,5 +1,5 @@ import assert from "node:assert/strict"; -import { execFileSync, spawn } from "node:child_process"; +import { execFileSync } from "node:child_process"; import { randomUUID } from "node:crypto"; import { chmod, @@ -16,7 +16,7 @@ import path from "node:path"; import { setTimeout as delay } from "node:timers/promises"; import { applicationRoot as mcpAppRoot, buildServer } from "./build-server.mjs"; -import { consumeLines } from "./consume-lines.mjs"; +import * as streams from "./support/streams.mjs"; if (process.platform !== "win32") { await testWorkbenchStateFallback(); @@ -530,66 +530,18 @@ async function writeFakePython(executablePath) { } function startServer(serverPath, env) { - const child = spawn(process.execPath, [serverPath, "--stdio"], { + const server = streams.startServer(serverPath, env, { cwd: path.dirname(path.dirname(serverPath)), - env, - stdio: ["pipe", "pipe", "pipe"], - }); - const responses = new Map(); - const waiters = new Map(); - const stderrEvents = []; - let stdout = ""; - let stderr = ""; - child.stdout.setEncoding("utf8"); - child.stdout.on("data", (chunk) => { - stdout += chunk; - stdout = consumeLines(stdout, (line) => { - const response = JSON.parse(line); - responses.set(response.id, response); - waiters.get(response.id)?.(response); - waiters.delete(response.id); - }); - }); - child.stderr.setEncoding("utf8"); - child.stderr.on("data", (chunk) => { - stderr += chunk; - stderr = consumeLines(stderr, (line) => { - try { - const event = JSON.parse(line); - if (event.component === "codex_security_workbench") - stderrEvents.push(event); - } catch { - // Tool errors are asserted from MCP responses; only structured diagnostics matter here. - } - }); + component: "codex_security_workbench", + withTimeout, + responseLabel: "response", + // Tool errors are asserted from MCP responses; only structured diagnostics matter here. + checkSignalCode: false, }); return { - request(id, method, params = {}) { - child.stdin.write( - `${JSON.stringify({ jsonrpc: "2.0", id, method, params })}\n`, - ); - const existing = responses.get(id); - if (existing) return Promise.resolve(existing); - return withTimeout( - new Promise((resolve) => waiters.set(id, resolve)), - 15_000, - `response ${id}`, - ); - }, - stderrEvents() { - return [...stderrEvents]; - }, - async stop() { - if (child.exitCode !== null) return; - child.stdin.end(); - const exited = new Promise((resolve) => child.once("exit", resolve)); - const graceful = await Promise.race([ - exited.then(() => true), - delay(2_000).then(() => false), - ]); - if (!graceful && child.exitCode === null) child.kill("SIGKILL"); - await exited; - }, + request: server.request.bind(server), + stderrEvents: server.stderrEvents, + stop: server.stop, }; } @@ -677,10 +629,7 @@ function assertToolError(response, pattern) { async function readJsonLines(filePath) { const content = await readFile(filePath, "utf8"); - return content - .split(/\r?\n/) - .filter(Boolean) - .map((line) => JSON.parse(line)); + return content.split(/\r?\n/).filter(Boolean).map(JSON.parse); } async function pathExists(filePath) { From 0f907561abfd2933d6eeb2e59cdf262ead87c0a4 Mon Sep 17 00:00:00 2001 From: Codex Date: Sat, 3 Oct 2026 22:51:21 +0000 Subject: [PATCH 4/4] refactor(plugin): use runtime timer for retry cancellation --- .../mcp-app/src/deep-scan/coordinator.ts | 22 +++++-------------- .../tests/test_deep_scan_coordinator.mjs | 9 +++----- 2 files changed, 9 insertions(+), 22 deletions(-) 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 29e492ff85..11d5bb6d46 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts @@ -2,6 +2,7 @@ import { randomUUID } from "node:crypto"; import { promises as fs } from "node:fs"; import { join } from "node:path"; +import { setTimeout as delay } from "node:timers/promises"; import { createDeepScanArtifacts, ensureDeepScanDirectories, @@ -1088,22 +1089,11 @@ export class DeepScanCoordinator { const systemClock: DeepScanClock = { now: () => Date.now(), sleep: async (delayMs, signal) => { - if (signal.aborted) throw abortError(signal.reason); - await new Promise((resolvePromise, rejectPromise) => { - const timeout = setTimeout(() => { - cleanup(); - resolvePromise(); - }, delayMs); - const onAbort = (): void => { - cleanup(); - rejectPromise(abortError(signal.reason)); - }; - const cleanup = (): void => { - clearTimeout(timeout); - signal.removeEventListener("abort", onAbort); - }; - signal.addEventListener("abort", onAbort, { once: true }); - }); + try { + await delay(delayMs, undefined, { signal }); + } catch (error) { + throw signal.aborted ? abortError(signal.reason) : error; + } }, }; 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 7bbda94c3f..23d567d9e5 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 @@ -1398,12 +1398,9 @@ async function testCancellationClearsRetryWait() { }); const sleepStarted = Promise.withResolvers(); const coordinator = createCoordinator(fixture, store, executor, { - clock: { - now: immediateClock.now, - sleep: async (_delayMs, signal) => { - sleepStarted.resolve(); - await waitForAbort(signal); - }, + clock: undefined, + log: (event) => { + if (event.event === "worker_retry_scheduled") sleepStarted.resolve(); }, }); coordinator.start();