Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
72 changes: 22 additions & 50 deletions plugins/codex-security/mcp-app/server.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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";
Expand All @@ -48,12 +50,10 @@ const WORKBENCH_COMMANDS_WITHOUT_DATABASE = new Set([
"read-artifact",
]);

type JsonObject = Record<string, unknown>;

let fallbackWorkbenchStateDir: Promise<string> | undefined;
let fallbackWorkbenchStateLogged = false;
let persistentWorkbenchStateSucceeded = false;
let workbenchStateSelectionTail: Promise<void> = Promise.resolve();
const workbenchStateSelectionLock = new AsyncLock();

const userContextSchema = z.string().trim().min(1);
const editableUserContextSchema = z.string().trim();
Expand Down Expand Up @@ -126,14 +126,6 @@ async function scanRoot(): Promise<string> {
return result.scanRoot;
}

interface WorkspaceState extends JsonObject {
id: string;
results?: JsonObject;
setup: {
submitted: boolean;
};
}

const diffTargetSchema = z.discriminatedUnion("kind", [
z
.object({
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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.",
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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),
};
}

Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -2567,22 +2558,6 @@ async function executeWorkbenchWithStateSelection(
});
}

async function withWorkbenchStateSelectionLock<T>(
operation: () => Promise<T>,
): Promise<T> {
const predecessor = workbenchStateSelectionTail;
let release!: () => void;
workbenchStateSelectionTail = new Promise<void>((resolvePromise) => {
release = resolvePromise;
});
await predecessor;
try {
return await operation();
} finally {
release();
}
}

async function executeWorkbench(
pythonCommand: string,
args: string[],
Expand Down Expand Up @@ -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.",
Expand All @@ -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
: "",
);
}
67 changes: 21 additions & 46 deletions plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
/// <reference lib="es2023.array" />
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,
Expand Down Expand Up @@ -162,7 +164,7 @@ export class DeepScanCoordinator {
scanId: this.state.scanId,
reason: errorKind(error),
});
this.failLocally(error);
if (this.stopLocally()) this.rejectTerminal(error);
});
}

Expand Down Expand Up @@ -325,7 +327,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;
Expand Down Expand Up @@ -447,29 +449,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;
}

Expand Down Expand Up @@ -601,7 +591,7 @@ export class DeepScanCoordinator {
} else {
this.state = current;
}
this.finishLocally(this.state);
if (this.stopLocally()) this.resolveTerminal(cloneState(this.state));
return true;
}

Expand Down Expand Up @@ -667,11 +657,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,
Expand Down Expand Up @@ -1046,10 +1034,8 @@ export class DeepScanCoordinator {

private trackSchedulerWork<T>(promise: Promise<T>): Promise<T> {
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;
}

Expand Down Expand Up @@ -1103,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<void>((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;
}
},
};

Expand Down
Loading
Loading