Skip to content
Open
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
4 changes: 2 additions & 2 deletions plugins/codex-security/mcp-app/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -765,7 +765,7 @@ export function createCodexSecurityServer(): McpServer {
log: logDeepScanEvent,
handoffClaimToken,
threadId,
onComplete: async (draft, signal) => {
onComplete: async (draft, signal, publication) => {
const context = await createScanArtifactContext(
begun.run.scanId,
runWorkbench,
Expand All @@ -779,7 +779,7 @@ export function createCodexSecurityServer(): McpServer {
await recordCodexSecurityScanDraftViaWorkbench(context, {
...draft,
...(handoffClaimToken === undefined ? {} : { handoffClaimToken })
}, runWorkbench, signal);
}, runWorkbench, signal, publication);
},
onStopped: async (run) => {
await runWorkbench([
Expand Down
12 changes: 11 additions & 1 deletion plugins/codex-security/mcp-app/src/artifact-scan-draft.ts
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ interface PreparedScanDraft {
coverage: JsonObject;
}

/** Host-selected Deep aggregate, separate from model-authored draft fields. */
export interface DeepScanPublication {
coordinatorGeneration?: number;
resultPath: string | null;
}

type PublishScanDraft = (
draft: PreparedScanDraft,
expectedDigest: string | undefined,
Expand Down Expand Up @@ -165,6 +171,7 @@ export async function recordCodexSecurityScanDraftViaWorkbench(
input: ScanDraftInput,
runWorkbench: RunArtifactWorkbench,
signal?: AbortSignal,
publication?: DeepScanPublication,
): Promise<ScanDraftResult> {
return recordCodexSecurityScanDraft(
context,
Expand All @@ -184,7 +191,10 @@ export async function recordCodexSecurityScanDraftViaWorkbench(
const { handoffClaimToken: _claim, ...snapshot } = checkpoint;
await Promise.all([
replaceArtifactJson(checkpointPath, snapshot),
replaceArtifactJson(draftPath, draft),
replaceArtifactJson(draftPath, {
...draft,
...(publication === undefined ? {} : { deepScanPublication: publication }),
}),
]);
const arguments_ = [
"write-scan-draft",
Expand Down
12 changes: 9 additions & 3 deletions plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
import { validateDiscoveryArtifacts, validateReducerArtifacts, type DeepReductionInput } from "./artifact-validation.js";
import {
scanDraftInputSchema,
type DeepScanPublication,
type ScanDraftInput
} from "../artifact-scan-draft.js";
import type { DeepScanArtifacts } from "./artifacts.js";
Expand Down Expand Up @@ -58,6 +59,7 @@ interface SchedulerResult {
mergedWorkerIds: string[];
reducers: AcceptedReducer[];
result?: DeepReductionInput;
resultPath?: string;
}

type CoordinatorPhase = "setup" | "discovery" | "terminal";
Expand Down Expand Up @@ -86,7 +88,7 @@ export interface CoordinatorOptions {
threadId?: string;
heartbeatIntervalMs?: number;
observeReplacement?: (run: DeepScanRunState) => Promise<DeepScanRunState>;
onComplete?: (draft: ScanDraftInput, signal: AbortSignal) => Promise<void>;
onComplete?: (draft: ScanDraftInput, signal: AbortSignal, publication: DeepScanPublication) => Promise<void>;
onStopped?: (run: DeepScanRunState) => Promise<void>;
}

Expand Down Expand Up @@ -319,7 +321,10 @@ export class DeepScanCoordinator {
if (draft.scanId !== this.state.scanId) {
throw new Error("Deep Scan aggregate does not match its authoritative scan identity.");
}
await this.options.onComplete?.(draft, this.publicationAbortController.signal);
await this.options.onComplete?.(draft, this.publicationAbortController.signal, {
coordinatorGeneration: this.state.coordinatorGeneration,
resultPath: schedulerResult.resultPath ?? null,
});
if (this.canceled || this.externallyFailed) return;
this.state = await this.finishWithReplay(schedulerResult);
if (this.canceled || this.externallyFailed) return;
Expand Down Expand Up @@ -966,6 +971,7 @@ export class DeepScanCoordinator {
mergedWorkerIds: unique(mergedDiscoveries.map((worker) => worker.id)),
reducers: reducerOutcomes,
result: latestResult,
resultPath: previousReducerResultPath,
};
}

Expand Down Expand Up @@ -1178,7 +1184,7 @@ function compareCompletionSequence(left: AcceptedDiscovery, right: AcceptedDisco
}

function workerLabelSequence(worker: PersistedDeepScanWorker, kind: "discovery" | "dedup"): number {
const match = worker.promptPath.match(new RegExp(`${kind}-(\\d+)`));
const match = basename(dirname(worker.promptPath)).match(new RegExp(`${kind}-(\\d+)`));
return match ? Number(match[1]) : 0;
}

Expand Down
14 changes: 14 additions & 0 deletions plugins/codex-security/mcp-app/tests/clock_workbench.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
"""Control the workbench clock while retaining its real commands and database writes."""

import os
import runpy
import sys
from pathlib import Path

source = Path(sys.argv[1])
sys.argv = [str(source), *sys.argv[2:]]
sys.path.insert(0, str(source.parent))
api = runpy.run_path(str(source), run_name="test_workbench")
instant = os.environ["TEST_WORKBENCH_NOW"]
api["main"].__globals__["now"] = lambda: instant
api["main"]()
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,7 @@ export async function testDeepScanPublication({

async function testPublicationUsesAcceptedReducerSnapshot() {
const fixture = await fixtureRun({ workers: 1, subagents: 0, stopAfterNoNew: 1, maxDiscoveryRuns: 1 });
fixture.run.coordinatorGeneration = 3;
const store = new FakeStore(fixture.run);
const commitDedup = store.commitDedup.bind(store);
store.commitDedup = async (commit) => {
Expand All @@ -163,17 +164,26 @@ export async function testDeepScanPublication({
return structuredClone(store.run);
};
const completed = [];
const published = [];
const coordinator = new DeepScanCoordinator({
run: fixture.run, store,
executor: new FakeExecutor({ discoveryCandidateId: "accepted-finding" }),
pluginRoot: fixture.pluginRoot, clock: immediateClock,
onComplete: async (draft) => completed.push(structuredClone(draft)),
onComplete: async (draft, _signal, publication) => {
completed.push(structuredClone(draft));
published.push(publication);
},
});
coordinator.start();
const terminal = await coordinator.wait(undefined, 5_000);
assert.equal(terminal?.status, "succeeded", terminal?.error);
assert.equal(completed[0].findings[0].provenance.candidateId, "accepted-finding");
assert.equal(completed[0].coverage.completeness, "complete");
const reducer = [...store.workers.values()].find((worker) => worker.kind === "dedup");
assert.deepEqual(published, [{
coordinatorGeneration: 3,
resultPath: reducer.resultManifestPath,
}]);
}

await testSaturationOmitsWorkerAcceptedDuringCancellation();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -591,6 +591,10 @@ try {
const obsoleteCheckpointPath = path.join(deepParentRoot, "checkpoints", "obsolete.json");
await writeFile(obsoleteCheckpointPath, "{malformed obsolete checkpoint\n");
let deepWorkbenchWrites = 0;
const deepPublication = {
coordinatorGeneration: 3,
resultPath: path.join(deepParentRoot, "workers", "reducer", "result.json"),
};
await recordCodexSecurityScanDraftViaWorkbench(
deepParentContext,
acceptedDeepDraft,
Expand All @@ -603,11 +607,15 @@ try {
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.deepScanPublication, deepPublication);
assert.equal(stagedCheckpoint.deepScanPublication, undefined);
assert.deepEqual(staged.findings, acceptedDeepFindings);
assert.deepEqual(staged.coverage, acceptedDeepCoverage);
assert.deepEqual(stagedCheckpoint.findings, acceptedDeepDraft.findings);
assert.equal(stagedCheckpoint.handoffClaimToken, undefined);
},
undefined,
deepPublication,
);
assert.equal(deepWorkbenchWrites, 1, "terminal Deep drafts still publish through the workbench lock despite obsolete malformed checkpoints");
assert.deepEqual(await readdir(path.join(deepParentRoot, "drafts")), []);
Expand Down
Loading
Loading