From ef9e2ab03da7adb57b35cbd0f85d4690ec1f1b76 Mon Sep 17 00:00:00 2001 From: Michael D'Angelo Date: Sat, 12 Sep 2026 08:13:43 +0000 Subject: [PATCH 1/4] fix(plugin): fence stale Deep publication and preserve replay --- plugins/codex-security/mcp-app/server.ts | 4 +- .../mcp-app/src/artifact-scan-draft.ts | 12 +- .../mcp-app/src/deep-scan/coordinator.ts | 10 +- .../tests/deep_scan_publication_cases.mjs | 12 +- .../tests/test_artifact_scan_draft.mjs | 8 + .../test_deep_scan_store_integration.mjs | 5 +- .../codex-security/scripts/workbench_db.py | 4 +- .../scripts/workbench_saved_results.py | 50 ++++-- .../test_deep_scan_publication_authority.py | 84 ++++++++++ .../test_deep_scan_publication_replay.py | 146 ++++++++++++++++++ .../test_deep_scan_successful_publication.py | 31 +++- 11 files changed, 348 insertions(+), 18 deletions(-) create mode 100644 plugins/codex-security/tests/test_deep_scan_publication_authority.py create mode 100644 plugins/codex-security/tests/test_deep_scan_publication_replay.py diff --git a/plugins/codex-security/mcp-app/server.ts b/plugins/codex-security/mcp-app/server.ts index a21fbb627..2c35f3c7d 100644 --- a/plugins/codex-security/mcp-app/server.ts +++ b/plugins/codex-security/mcp-app/server.ts @@ -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, @@ -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([ 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 4c560822c..125a02ce9 100644 --- a/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts +++ b/plugins/codex-security/mcp-app/src/artifact-scan-draft.ts @@ -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, @@ -165,6 +171,7 @@ export async function recordCodexSecurityScanDraftViaWorkbench( input: ScanDraftInput, runWorkbench: RunArtifactWorkbench, signal?: AbortSignal, + publication?: DeepScanPublication, ): Promise { return recordCodexSecurityScanDraft( context, @@ -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", 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 905c0c9a8..e88b473b4 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts @@ -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"; @@ -58,6 +59,7 @@ interface SchedulerResult { mergedWorkerIds: string[]; reducers: AcceptedReducer[]; result?: DeepReductionInput; + resultPath?: string; } type CoordinatorPhase = "setup" | "discovery" | "terminal"; @@ -86,7 +88,7 @@ export interface CoordinatorOptions { threadId?: string; heartbeatIntervalMs?: number; observeReplacement?: (run: DeepScanRunState) => Promise; - onComplete?: (draft: ScanDraftInput, signal: AbortSignal) => Promise; + onComplete?: (draft: ScanDraftInput, signal: AbortSignal, publication: DeepScanPublication) => Promise; onStopped?: (run: DeepScanRunState) => Promise; } @@ -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; @@ -966,6 +971,7 @@ export class DeepScanCoordinator { mergedWorkerIds: unique(mergedDiscoveries.map((worker) => worker.id)), reducers: reducerOutcomes, result: latestResult, + resultPath: previousReducerResultPath, }; } 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 5e44d1796..4221c95ee 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 @@ -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) => { @@ -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(); 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 72f815fc3..aa5859ad2 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 @@ -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, @@ -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")), []); 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 a93384919..589f7874f 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 @@ -91,7 +91,10 @@ async function testRecoveredPublicationRejectsLateFailure() { paths: ["fixture.py"], }], }, - }, runWorkbench); + }, runWorkbench, undefined, { + coordinatorGeneration: claim.run.coordinatorGeneration, + resultPath: null, + }); await runWorkbench([ "cancel-scan", "--scan-id", run.scanId, "--thread-id", "publication-failure-owner", diff --git a/plugins/codex-security/scripts/workbench_db.py b/plugins/codex-security/scripts/workbench_db.py index 161756ee0..b53a506ec 100644 --- a/plugins/codex-security/scripts/workbench_db.py +++ b/plugins/codex-security/scripts/workbench_db.py @@ -1547,8 +1547,10 @@ def add_warning() -> None: wrote = True manifest, findings, _ = _write_prepared_scan_finalization(prepared) except ContractError as exc: - if wrote or ( + # Replay a validated Deep aggregate after an output write fails. + if (wrote and scan["mode"] != "deep") or ( scan["mode"] == "deep" + and not wrote and not already_sealed and not isinstance(exc, RecoverableContractError) ): diff --git a/plugins/codex-security/scripts/workbench_saved_results.py b/plugins/codex-security/scripts/workbench_saved_results.py index a32b96513..0c72a2e0d 100644 --- a/plugins/codex-security/scripts/workbench_saved_results.py +++ b/plugins/codex-security/scripts/workbench_saved_results.py @@ -1383,6 +1383,41 @@ def save_scan_artifact(db: Any, connection: Any, args: Any) -> dict[str, Any]: return {"scanId": scan_id, "path": str(scan_dir / output)} +def _read_staged_scan_draft(scan_dir: Path, draft_path: str) -> dict[str, Any]: + try: + relative = Path(draft_path).relative_to(scan_dir).as_posix() + except ValueError as exc: + raise SystemExit("Scan draft must be inside the registered scan drafts directory.") from exc + if not re.fullmatch(r"drafts/[0-9a-fA-F-]+\.json", relative): + raise SystemExit("Scan draft must be inside the registered scan drafts directory.") + return _read_scan_local_json(scan_dir, relative, "Staged scan draft") + + +def _require_current_deep_publication( + db: Any, connection: Any, scan_id: str, draft: dict[str, Any] +) -> None: + publication = draft.get("deepScanPublication") + run = db.deep_scan.require_deep_scan_run(connection, scan_id) + db.deep_scan.require_current_coordinator( + run, + argparse.Namespace( + coordinator_generation=publication.get("coordinatorGeneration") if publication else None + ), + ) + # Generation-one runs predate host publication metadata. Keep their existing + # draft path; adopted coordinators must carry their generation and selection. + if publication is None: + return + reducer = _latest_successful_reducer( + connection.execute( + "SELECT * FROM deep_scan_workers WHERE scan_id = ?", (scan_id,) + ).fetchall() + ) + selected_result = reducer["result_manifest_path"] if reducer is not None else None + if publication["resultPath"] != selected_result: + raise SystemExit("Deep Scan aggregate belongs to a superseded publication selection.") + + def write_scan_draft(db: Any, connection: Any, args: Any) -> dict[str, Any]: scan_id = db.require_uuid(args.scan_id, "scan-id") with db.scan_completion_lock(scan_id): @@ -1395,6 +1430,10 @@ def write_scan_draft(db: Any, connection: Any, args: Any) -> dict[str, Any]: "The scan stopped; its saved checkpoint was retained without replacing sealed results." ) scan_dir = db.require_canonical_scan_directory(Path(scan["scan_dir"])) + draft = None + if scan["mode"] == "deep": + draft = _read_staged_scan_draft(scan_dir, args.draft_path) + _require_current_deep_publication(db, connection, scan_id, draft) if args.checkpoint_path is not None: try: checkpoint_relative = Path(args.checkpoint_path).relative_to(scan_dir).as_posix() @@ -1424,15 +1463,8 @@ def write_scan_draft(db: Any, connection: Any, args: Any) -> dict[str, Any]: raise SystemExit( "scan_draft_conflict: canonical scan results changed; reconcile the saved checkpoint again." ) - try: - relative = Path(args.draft_path).relative_to(scan_dir).as_posix() - except ValueError as exc: - raise SystemExit( - "Scan draft must be inside the registered scan drafts directory." - ) from exc - if not re.fullmatch(r"drafts/[0-9a-fA-F-]+\.json", relative): - raise SystemExit("Scan draft must be inside the registered scan drafts directory.") - draft = _read_scan_local_json(scan_dir, relative, "Staged scan draft") + if draft is None: + draft = _read_staged_scan_draft(scan_dir, args.draft_path) manifest, findings, coverage = draft["manifest"], draft["findings"], draft["coverage"] binding = db.workbench_completion_binding(scan, db.now()) # Validate on copies: saved canonical documents remain ordinary unsealed drafts. diff --git a/plugins/codex-security/tests/test_deep_scan_publication_authority.py b/plugins/codex-security/tests/test_deep_scan_publication_authority.py new file mode 100644 index 000000000..b4ff3a0bd --- /dev/null +++ b/plugins/codex-security/tests/test_deep_scan_publication_authority.py @@ -0,0 +1,84 @@ +from __future__ import annotations + +import copy +import json +import uuid +from argparse import Namespace + +import pytest +from test_deep_scan_successful_publication import add_worker +from test_deep_scan_successful_publication import publication_scan as publication_scan + + +def stage_publication(scan, *, generation, result_path, title): + draft_dir = scan.scan_dir / "drafts" + draft_dir.mkdir(exist_ok=True) + draft_path = draft_dir / f"{uuid.uuid4()}.json" + checkpoint_path = draft_dir / f"{uuid.uuid4()}.checkpoint.json" + findings = copy.deepcopy(scan.findings) + findings[0]["title"] = title + draft = { + "manifest": json.loads((scan.scan_dir / "scan-manifest.json").read_text()), + "findings": {"findings": findings}, + "coverage": scan.coverage, + } + if generation is not None: + draft["deepScanPublication"] = { + "coordinatorGeneration": generation, + "resultPath": str(result_path), + } + draft_path.write_text(json.dumps(draft)) + checkpoint_path.write_text( + json.dumps({"scanId": scan.scan_id, "findings": findings, "coverage": scan.coverage}) + ) + return Namespace( + scan_id=scan.scan_id, + claim_token=None, + draft_path=str(draft_path), + checkpoint_path=str(checkpoint_path), + expected_draft_digest=None, + ) + + +@pytest.mark.parametrize("stale", ["generation", "aggregate", "unfenced"]) +def test_stale_coordinator_cannot_replace_newer_canonical_publication( + workbench_api, workbench_db, publication_scan, stale +): + scan = publication_scan() + old_result = add_worker(workbench_db, scan) + new_result = add_worker(workbench_db, scan) + with workbench_db: + workbench_db.execute( + "UPDATE deep_scan_runs SET coordinator_generation = 3 WHERE scan_id = ?", + (scan.scan_id,), + ) + for result, completed_at in ((old_result, "2026-01-01"), (new_result, "2026-01-02")): + workbench_db.execute( + "UPDATE deep_scan_workers SET kind = 'dedup', merge_state = 'none', completed_at = ? " + "WHERE result_manifest_path = ?", + (completed_at, str(result)), + ) + current = stage_publication( + scan, generation=3, result_path=new_result, title="Current accepted aggregate" + ) + workbench_api["write_scan_draft"](workbench_db, current) + saved = { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*.json") + if "drafts" not in path.parts + } + old = stage_publication( + scan, + generation=None if stale == "unfenced" else 2 if stale == "generation" else 3, + result_path=old_result if stale == "aggregate" else new_result, + title="Superseded aggregate", + ) + + with pytest.raises(SystemExit, match="coordinator|aggregate"): + workbench_api["write_scan_draft"](workbench_db, old) + + assert { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*.json") + if "drafts" not in path.parts + } == saved diff --git a/plugins/codex-security/tests/test_deep_scan_publication_replay.py b/plugins/codex-security/tests/test_deep_scan_publication_replay.py new file mode 100644 index 000000000..af7a629be --- /dev/null +++ b/plugins/codex-security/tests/test_deep_scan_publication_replay.py @@ -0,0 +1,146 @@ +from __future__ import annotations + +import json +import sqlite3 +import subprocess +import sys +from argparse import Namespace +from pathlib import Path + +import pytest +from test_deep_scan_publication_authority import stage_publication +from test_deep_scan_successful_publication import add_worker +from test_deep_scan_successful_publication import publication_scan as publication_scan + +_CRASH_PUBLICATION = """ +import json, os, runpy, sqlite3, sys +from argparse import Namespace + +api = runpy.run_path(sys.argv[1], run_name="publication_crash_test") +args = Namespace(**json.loads(sys.argv[3])) +boundary = sys.argv[4] + +class CrashConnection(sqlite3.Connection): + def commit(self): + completing = self.execute( + "SELECT status FROM scans WHERE id = ?", (args.scan_id,) + ).fetchone()[0] == "complete" + if completing and boundary == "sqlite-before": + os._exit(72) + super().commit() + if completing and boundary == "sqlite-after": + os._exit(73) + +connection = sqlite3.connect(sys.argv[2], factory=CrashConnection) +connection.row_factory = sqlite3.Row +connection.execute("PRAGMA foreign_keys = ON") +if boundary.startswith("sqlite-"): + api["complete_scan"](connection, Namespace( + scan_id=args.scan_id, claim_token=None, cost_json=None + )) +else: + saved = api["saved_results"] + original_write = saved.write_scan_local_bytes + def crash_after_write(root, relative, contents): + original_write(root, relative, contents) + if relative == boundary: + os._exit(71) + saved.write_scan_local_bytes = crash_after_write + api["write_scan_draft"](connection, args) +raise AssertionError("publication never reached the requested crash boundary") +""" + + +@pytest.mark.parametrize( + "boundary", + ["findings.json", "coverage.json", "scan-manifest.json", "sqlite-before", "sqlite-after"], +) +def test_publication_crash_replays_selected_input_without_stale_overwrite( + workbench_api, workbench_db, publication_scan, tmp_path, boundary +): + scan = publication_scan() + result_path = add_worker(workbench_db, scan) + with workbench_db: + workbench_db.execute( + "UPDATE deep_scan_runs SET coordinator_generation = 3 WHERE scan_id = ?", + (scan.scan_id,), + ) + workbench_db.execute( + "UPDATE deep_scan_workers SET kind = 'dedup', merge_state = 'none' WHERE scan_id = ?", + (scan.scan_id,), + ) + current = stage_publication( + scan, generation=3, result_path=result_path, title="Selected aggregate" + ) + stale = stage_publication( + scan, generation=2, result_path=result_path, title="Obsolete coordinator draft" + ) + database_path = tmp_path / "publication.sqlite3" + with sqlite3.connect(database_path) as connection: + workbench_db.backup(connection) + connection.row_factory = sqlite3.Row + if boundary.startswith("sqlite-"): + workbench_api["write_scan_draft"](connection, current) + + child = subprocess.run( + [ + sys.executable, + "-c", + _CRASH_PUBLICATION, + str(Path(__file__).resolve().parents[1] / "scripts" / "workbench_db.py"), + str(database_path), + json.dumps(vars(current)), + boundary, + ], + capture_output=True, + text=True, + ) + assert child.returncode == {"sqlite-before": 72, "sqlite-after": 73}.get(boundary, 71), ( + child.stdout, + child.stderr, + ) + interrupted = { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*") + if path.is_file() and "drafts" not in path.parts + } + with sqlite3.connect(database_path) as connection: + connection.row_factory = sqlite3.Row + connection.execute("PRAGMA foreign_keys = ON") + row = connection.execute("SELECT * FROM scans WHERE id = ?", (scan.scan_id,)).fetchone() + assert row["status"] == ("complete" if boundary == "sqlite-after" else "running") + assert bool(row["seal_manifest_digest"]) == (boundary == "sqlite-after") + run_before = dict(connection.execute("SELECT * FROM deep_scan_runs").fetchone()) + workers_before = [ + dict(row) for row in connection.execute("SELECT * FROM deep_scan_workers") + ] + + with pytest.raises(SystemExit, match="coordinator|stopped"): + workbench_api["write_scan_draft"](connection, stale) + assert all(path.read_bytes() == contents for path, contents in interrupted.items()) + + if not boundary.startswith("sqlite-"): + workbench_api["write_scan_draft"](connection, current) + completed = workbench_api["complete_scan"]( + connection, Namespace(scan_id=scan.scan_id, claim_token=None, cost_json=None) + )["scan"] + assert completed["progress"]["status"] == "complete" + assert completed["findingCount"] == 1 + assert dict(connection.execute("SELECT * FROM deep_scan_runs").fetchone()) == run_before + assert [ + dict(row) for row in connection.execute("SELECT * FROM deep_scan_workers") + ] == workers_before + assert connection.execute("SELECT COUNT(*) FROM finding_occurrences").fetchone()[0] == 1 + published = { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*") + if path.is_file() and "drafts" not in path.parts + } + workbench_api["complete_scan"]( + connection, Namespace(scan_id=scan.scan_id, claim_token=None, cost_json=None) + ) + assert all(path.read_bytes() == contents for path, contents in published.items()) + if boundary.startswith("sqlite-"): + assert published == interrupted + findings = json.loads((scan.scan_dir / "findings.json").read_text())["findings"] + assert findings[0]["title"] == "Selected aggregate" diff --git a/plugins/codex-security/tests/test_deep_scan_successful_publication.py b/plugins/codex-security/tests/test_deep_scan_successful_publication.py index 8a144ece4..95c799dc8 100644 --- a/plugins/codex-security/tests/test_deep_scan_successful_publication.py +++ b/plugins/codex-security/tests/test_deep_scan_successful_publication.py @@ -72,7 +72,7 @@ def create(*, mode="deep", scope="."): "INSERT INTO deep_scan_runs (scan_id, schema_version, workflow_version, " "status, phase, workers, subagents, stop_after_no_new, max_discovery_runs, " "manifest_path, terminal_reason, created_at, updated_at, completed_at) " - "VALUES (?, 1, 'publication-test', 'succeeded', 'terminal', 1, 0, 1, 1, " + "VALUES (?, 1, 'deep-security-scan/v1', 'succeeded', 'terminal', 1, 0, 1, 1, " "?, 'saturated', ?, ?, ?)", ( scan_id, @@ -460,3 +460,32 @@ def test_standard_publication_preserves_deliberately_partial_coverage( complete(workbench_api, workbench_db, scan) assert_published_aggregate(scan) + + +def test_deep_publication_write_failure_keeps_original_terminal_cause( + workbench_api, workbench_db, publication_scan, monkeypatch +): + scan = publication_scan() + finalizer_globals = workbench_api["_write_prepared_scan_finalization"].__globals__ + write_bytes = finalizer_globals["write_scan_local_bytes"] + + def fail_report(scan_dir, relative_path, payload, **kwargs): + if relative_path == "report.md": + raise finalizer_globals["ContractError"]("Synthetic report write interruption") + return write_bytes(scan_dir, relative_path, payload, **kwargs) + + with monkeypatch.context() as patch: + patch.setitem(finalizer_globals, "write_scan_local_bytes", fail_report) + with pytest.raises(SystemExit, match="Synthetic report write interruption"): + complete(workbench_api, workbench_db, scan) + + assert ( + workbench_db.execute("SELECT status FROM scans WHERE id = ?", (scan.scan_id,)).fetchone()[0] + == "running" + ) + run = workbench_db.execute( + "SELECT status, terminal_reason FROM deep_scan_runs WHERE scan_id = ?", (scan.scan_id,) + ).fetchone() + assert tuple(run) == ("succeeded", "saturated") + assert complete(workbench_api, workbench_db, scan)["progress"]["status"] == "complete" + assert_published_aggregate(scan) From 0e623ac2ffc7d4d535268235144e6e9346424988 Mon Sep 17 00:00:00 2001 From: Michael D'Angelo Date: Wed, 16 Sep 2026 06:45:50 +0000 Subject: [PATCH 2/4] test(plugin): preserve legacy Deep publication replay --- .../test_deep_scan_publication_authority.py | 42 +++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/plugins/codex-security/tests/test_deep_scan_publication_authority.py b/plugins/codex-security/tests/test_deep_scan_publication_authority.py index b4ff3a0bd..99865d82d 100644 --- a/plugins/codex-security/tests/test_deep_scan_publication_authority.py +++ b/plugins/codex-security/tests/test_deep_scan_publication_authority.py @@ -82,3 +82,45 @@ def test_stale_coordinator_cannot_replace_newer_canonical_publication( for path in scan.scan_dir.rglob("*.json") if "drafts" not in path.parts } == saved + + +@pytest.mark.parametrize("generation", [None, 3], ids=["legacy-generation-one", "current-lease"]) +def test_current_publication_replays_without_changing_checkpoint_or_worker_state( + workbench_api, workbench_db, publication_scan, generation +): + scan = publication_scan() + result = add_worker(workbench_db, scan) + with workbench_db: + workbench_db.execute( + "UPDATE deep_scan_runs SET coordinator_generation = ? WHERE scan_id = ?", + (generation or 1, scan.scan_id), + ) + workbench_db.execute( + "UPDATE deep_scan_workers SET kind = 'dedup', merge_state = 'none' WHERE scan_id = ?", + (scan.scan_id,), + ) + draft = stage_publication( + scan, generation=generation, result_path=result, title="Accepted aggregate" + ) + run_before = dict(workbench_db.execute("SELECT * FROM deep_scan_runs").fetchone()) + worker_before = dict(workbench_db.execute("SELECT * FROM deep_scan_workers").fetchone()) + + workbench_api["write_scan_draft"](workbench_db, draft) + published = { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*.json") + if "drafts" not in path.parts + } + replay = workbench_api["write_scan_draft"](workbench_db, draft) + + assert replay == {"scanId": scan.scan_id, "status": "draft_written"} + assert { + path: path.read_bytes() + for path in scan.scan_dir.rglob("*.json") + if "drafts" not in path.parts + } == published + assert len(list((scan.scan_dir / "checkpoints").glob("*.json"))) == 1 + assert dict(workbench_db.execute("SELECT * FROM deep_scan_runs").fetchone()) == run_before + assert dict(workbench_db.execute("SELECT * FROM deep_scan_workers").fetchone()) == worker_before + findings = json.loads((scan.scan_dir / "findings.json").read_text())["findings"] + assert findings[0]["title"] == "Accepted aggregate" From fed714002416bc2f73bf0cdf77fe72a4c4cba866 Mon Sep 17 00:00:00 2001 From: Michael D'Angelo Date: Wed, 16 Sep 2026 08:29:25 +0000 Subject: [PATCH 3/4] Preserve Deep parent drafts and reducer publication order --- .../mcp-app/tests/clock_workbench.py | 14 + ...test_deep_scan_publication_integration.mjs | 279 ++++++++++++++++++ .../scripts/workbench_saved_results.py | 24 +- .../test_deep_scan_publication_authority.py | 29 +- 4 files changed, 336 insertions(+), 10 deletions(-) create mode 100644 plugins/codex-security/mcp-app/tests/clock_workbench.py create mode 100644 plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs diff --git a/plugins/codex-security/mcp-app/tests/clock_workbench.py b/plugins/codex-security/mcp-app/tests/clock_workbench.py new file mode 100644 index 000000000..22d2c52b6 --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/clock_workbench.py @@ -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"]() diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs new file mode 100644 index 000000000..270b9466b --- /dev/null +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs @@ -0,0 +1,279 @@ +import assert from "node:assert/strict"; +import { execFile } from "node:child_process"; +import { mkdir, mkdtemp, readFile, readdir, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import path from "node:path"; +import { after, test } from "node:test"; +import { fileURLToPath } from "node:url"; +import { promisify } from "node:util"; +import { Client } from "@modelcontextprotocol/sdk/client/index.js"; +import { StdioClientTransport } from "@modelcontextprotocol/sdk/client/stdio.js"; +import { build } from "esbuild"; + +const execFileAsync = promisify(execFile); +const testsRoot = path.dirname(fileURLToPath(import.meta.url)); +const mcpAppRoot = path.resolve(testsRoot, ".."); +const pluginRoot = path.resolve(mcpAppRoot, ".."); +const workbenchPath = path.join(pluginRoot, "scripts", "workbench_db.py"); +const python = process.env.PYTHON?.trim() || "python3"; +const owner = "publication-test-owner"; +const coverage = { completeness: "complete", surfaces: [], explicitExclusions: [], deferred: [] }; +const highId = "ffffffff-ffff-4fff-8fff-ffffffffffff"; +const lowId = "00000000-0000-4000-8000-000000000001"; +const bundleRoot = await mkdtemp(path.join(tmpdir(), "deep-publication-bundle-")); +after(() => rm(bundleRoot, { recursive: true, force: true })); +const serverPath = path.join(bundleRoot, "server.cjs"); +await build({ + bundle: true, + define: { __dirname: JSON.stringify(mcpAppRoot), "import.meta.url": "__filename" }, + entryPoints: [path.join(mcpAppRoot, "main.ts")], + external: ["fsevents"], + format: "cjs", + loader: { ".md": "text" }, + logLevel: "silent", + outfile: serverPath, + platform: "node", + target: "node20", +}); +const hostBundle = await build({ + bundle: true, + stdin: { + contents: [ + 'export { WorkbenchDeepScanStore } from "./src/deep-scan/store.ts";', + 'export { DeepScanCoordinator } from "./src/deep-scan/coordinator.ts";', + 'export { createScanArtifactContext } from "./src/artifact-context.ts";', + 'export { recordCodexSecurityScanDraftViaWorkbench } from "./src/artifact-scan-draft.ts";', + ].join("\n"), + resolveDir: mcpAppRoot, + }, + format: "esm", + loader: { ".md": "text" }, + platform: "node", + write: false, +}); +const { + WorkbenchDeepScanStore, DeepScanCoordinator, createScanArtifactContext, + recordCodexSecurityScanDraftViaWorkbench, +} = await import(`data:text/javascript;base64,${Buffer.from(hostBundle.outputFiles[0].contents).toString("base64")}`); + +test("public partial drafts retain checkpoints after coordinator adoption", async (t) => { + const fixture = await createFixture(t); + const { run, store, call, runWorkbench } = fixture; + assertSuccess(await call("record_codex_security_scan_draft", partial(run, "generation-one"))); + const first = await checkpoints(run); + assert.equal(Object.keys(first).length, 1); + const claim = await store.claimCoordinator({ scanId: run.scanId, threadId: owner }); + assert.equal(claim.acquired, true); + assert.equal(claim.run.coordinatorGeneration, 2); + + assertSuccess(await call("record_codex_security_scan_draft", partial(run, "generation-two"))); + const retained = await checkpoints(run); + assert.equal(Object.keys(retained).length, 2); + for (const [name, bytes] of Object.entries(first)) assert.equal(retained[name], bytes); + assert.deepEqual(new Set(Object.values(retained).flatMap((bytes) => ( + JSON.parse(bytes).coverage.deferred.map((item) => item.candidateId) + ))), new Set(["generation-one", "generation-two"])); + + const before = await snapshot(run); + assertToolError(await call("record_codex_security_scan_draft", { + ...partial(run, "unfenced-final"), complete: true, + }), /current coordinator lease/); + const context = await createScanArtifactContext(run.scanId, runWorkbench, { requireRunning: true }); + await assert.rejects(recordCodexSecurityScanDraftViaWorkbench( + context, partial(run, "stale-generation"), runWorkbench, undefined, + { coordinatorGeneration: 1, resultPath: null }, + ), /newer generation/); + assertToolError(await call("complete_codex_security_scan", { scanId: run.scanId }), /deep/i); + assert.deepEqual(await snapshot(run), before); + assert.deepEqual(await readdir(path.join(run.scanDir, "drafts")), []); +}); + +for (const [label, offsets, ids] of [ + ["increasing timestamps", [1, 2], [highId, lowId]], + ["equal timestamps and descending UUIDs", [1, 1], [highId, lowId]], + ["decreasing timestamps", [2, 1], [lowId, highId]], + ["equal timestamps and ascending UUIDs", [1, 1], [lowId, highId]], +]) { + test(`selected reducer survives recovery and public completion: ${label}`, async (t) => { + const fixture = await createFixture(t); + const { run, store, call, runWorkbench, instant } = fixture; + assertSuccess(await call("record_codex_security_scan_draft", partial(run, "parent-checkpoint-only"))); + const parentCheckpoints = await checkpoints(run); + assert.equal(Object.keys(parentCheckpoints).length, 1); + const claimed = await store.claimCoordinator({ scanId: run.scanId, threadId: owner }); + assert.equal(claimed.run.coordinatorGeneration, 2); + assert.equal(claimed.run.config.stopAfterNoNew, 4); + const results = await commitReducers(fixture, offsets, ids); + const persisted = await store.get(run.scanId, owner); + assert.equal(persisted.noNewStreak, 4); + assert.deepEqual(persisted.persistedWorkers.filter((worker) => worker.kind === "discovery") + .map((worker) => worker.mergeState), Array(4).fill("merged")); + const context = await createScanArtifactContext(run.scanId, runWorkbench, { requireRunning: true }); + const publish = async (index, coordinatorGeneration = 2) => ( + recordCodexSecurityScanDraftViaWorkbench(context, { + ...JSON.parse(await readFile(results[index], "utf8")), coverage, + }, runWorkbench, undefined, { coordinatorGeneration, resultPath: results[index] }) + ); + + await publish(1); + const selected = await snapshot(run); + await assert.rejects(publish(0), /superseded publication selection/); + await assert.rejects(publish(1, 1), /newer generation/); + assert.deepEqual(await snapshot(run), selected); + + const publications = []; + const coordinator = new DeepScanCoordinator({ + run: persisted, store, pluginRoot, + executor: { run() { throw new Error("Recovery must reuse the committed workers."); } }, + clock: { now: () => Date.parse(instant), sleep: async () => {} }, + onComplete: async (draft, signal, publication) => { + publications.push(publication); + await recordCodexSecurityScanDraftViaWorkbench(context, draft, runWorkbench, signal, publication); + }, + }); + coordinator.start(); + const terminal = await coordinator.wait(undefined, 30_000); + assert.equal(terminal?.status, "succeeded", terminal?.error); + assert.deepEqual(publications, [{ coordinatorGeneration: 2, resultPath: results[1] }]); + + const beforeCompletion = await snapshot(run); + const lateDraft = partial(run, "late-parent-checkpoint"); + assertToolError(await call("record_codex_security_scan_draft", lateDraft)); + assert.deepEqual(await snapshot(run), beforeCompletion); + assertSuccess(await call("complete_codex_security_scan", { scanId: run.scanId })); + const sealed = await snapshot(run); + const manifest = JSON.parse(sealed.files["scan-manifest.json"]); + assert.equal(manifest.scan.threatModel.summary, "dedup-0002"); + assert.ok(manifest.scan.sealedAt); + assert.deepEqual(JSON.parse(sealed.files["findings.json"]).findings, []); + assert.deepEqual(JSON.parse(sealed.files["coverage.json"]).deferred, []); + assert.ok(sealed.files["report.md"]); + for (const [name, bytes] of Object.entries(parentCheckpoints)) { + assert.equal(sealed.checkpoints[name], bytes, "completion retains the immutable parent checkpoint"); + } + assert.deepEqual(sealed.checkpoints, beforeCompletion.checkpoints); + assertSuccess(await call("complete_codex_security_scan", { scanId: run.scanId })); + assert.deepEqual(await snapshot(run), sealed, "completion replay is byte-stable"); + assertToolError(await call("record_codex_security_scan_draft", lateDraft)); + assert.deepEqual(await snapshot(run), sealed, "late writes cannot replace sealed results"); + }); +} + +async function createFixture(t) { + const root = await mkdtemp(path.join(tmpdir(), "deep-publication-")); + let client; + t.after(async () => { + await client?.close(); + await rm(root, { recursive: true, force: true }); + }); + const targetPath = path.join(root, "target"); + await mkdir(targetPath); + await writeFile(path.join(targetPath, "fixture.py"), "# Synthetic publication fixture\n"); + const environment = { + ...process.env, + CODEX_HOME: path.join(root, "home"), + CODEX_SECURITY_STATE_DIR: path.join(root, "state"), + CODEX_SECURITY_SCAN_ROOT: path.join(root, "scans"), + }; + delete environment.CODEX_SECURITY_DEEP_SCAN_CONFIG_PATH; + const instant = new Date().toISOString(); + let now = instant; + const runWorkbench = async (args, input) => { + const child = execFileAsync(python, [path.join(testsRoot, "clock_workbench.py"), workbenchPath, ...args], { + cwd: pluginRoot, env: { ...environment, TEST_WORKBENCH_NOW: now }, + timeout: 30_000, maxBuffer: 4 * 1024 * 1024, + }); + child.child.stdin.end(input); + return JSON.parse((await child).stdout); + }; + const store = new WorkbenchDeepScanStore(runWorkbench); + const { run } = await store.begin({ targetPath, scope: ".", threadId: owner, scanRoot: environment.CODEX_SECURITY_SCAN_ROOT }); + client = new Client({ name: "publication-integration", version: "1.0.0" }); + const transport = new StdioClientTransport({ + command: process.execPath, args: [serverPath, "--stdio"], cwd: mcpAppRoot, + env: environment, stderr: "pipe", + }); + transport.stderr?.resume(); + await client.connect(transport); + return { + run, store, runWorkbench, instant, + setTime(offset) { now = new Date(Date.parse(instant) + offset * 1_000).toISOString(); }, + call(name, args) { + return client.callTool({ name, arguments: args, _meta: { "openai/threadId": owner } }); + }, + }; +} + +async function commitReducers(fixture, offsets, ids) { + const { run, store } = fixture; + const results = []; + for (let batch = 0; batch < 2; batch++) { + const workerIds = []; + for (let index = batch * 2 + 1; index <= batch * 2 + 2; index++) { + const id = `10000000-0000-4000-8000-${String(index).padStart(12, "0")}`; + const label = `discovery-${String(index).padStart(4, "0")}`; + const artifact = await workerArtifact(run, "workers", label); + const update = { id, scanId: run.scanId, kind: "discovery", status: "running", attempt: 1, ...artifact }; + await store.updateWorker(update); + const resultManifestPath = path.join(artifact.artifactDir, "result.json"); + await writeFile(resultManifestPath, JSON.stringify({ scanId: run.scanId, findings: [], coverage })); + await store.updateWorker({ ...update, status: "succeeded", resultManifestPath }); + workerIds.push(id); + } + const label = `dedup-${String(batch + 1).padStart(4, "0")}`; + const artifact = await workerArtifact(run, "dedup", label); + const id = ids[batch]; + await store.claimDedup({ id, scanId: run.scanId, workerIds, ...artifact }); + await store.updateWorker({ id, scanId: run.scanId, kind: "dedup", status: "running", attempt: 1, ...artifact }); + const resultManifestPath = path.join(artifact.artifactDir, "result.json"); + await writeFile(resultManifestPath, JSON.stringify({ scanId: run.scanId, findings: [], threatModel: { summary: label } })); + fixture.setTime(offsets[batch]); + await store.commitDedup({ id, scanId: run.scanId, resultManifestPath, newFindings: 0 }); + fixture.setTime(0); + results.push(resultManifestPath); + } + return results; +} + +async function workerArtifact(run, kind, label) { + const artifactDir = path.join(run.scanDir, "artifacts", "deep_discovery", kind, label, "output"); + const promptPath = path.join(path.dirname(artifactDir), "prompt.md"); + await mkdir(artifactDir, { recursive: true }); + await writeFile(promptPath, `Synthetic ${label}\n`); + return { artifactDir, promptPath }; +} + +function partial(run, candidateId) { + return { + scanId: run.scanId, complete: false, findings: [], + coverage: { + ...coverage, completeness: "partial", + deferred: [{ candidateId, reason: "Synthetic review remains pending.", paths: ["fixture.py"] }], + }, + }; +} + +async function checkpoints(run) { + const directory = path.join(run.scanDir, "checkpoints"); + const files = {}; + for (const name of await readdir(directory)) files[name] = await readFile(path.join(directory, name), "utf8"); + return files; +} + +async function snapshot(run) { + const files = {}; + for (const name of ["scan-manifest.json", "findings.json", "coverage.json", "report.md"]) { + try { files[name] = await readFile(path.join(run.scanDir, name), "utf8"); } + catch (error) { if (error.code !== "ENOENT") throw error; } + } + return { files, checkpoints: await checkpoints(run) }; +} + +function assertSuccess(result) { + assert.notEqual(result.isError, true, JSON.stringify(result)); +} + +function assertToolError(result, pattern) { + assert.equal(result.isError, true, JSON.stringify(result)); + if (pattern) assert.match(JSON.stringify(result), pattern); +} diff --git a/plugins/codex-security/scripts/workbench_saved_results.py b/plugins/codex-security/scripts/workbench_saved_results.py index 6cbc7ad54..f51ba74a5 100644 --- a/plugins/codex-security/scripts/workbench_saved_results.py +++ b/plugins/codex-security/scripts/workbench_saved_results.py @@ -1396,6 +1396,8 @@ def _require_current_deep_publication( db: Any, connection: Any, scan_id: str, draft: dict[str, Any] ) -> None: publication = draft.get("deepScanPublication") + if publication is None and draft["manifest"]["scan"].get("complete") is False: + return run = db.deep_scan.require_deep_scan_run(connection, scan_id) db.deep_scan.require_current_coordinator( run, @@ -1407,10 +1409,24 @@ def _require_current_deep_publication( # draft path; adopted coordinators must carry their generation and selection. if publication is None: return - reducer = _latest_successful_reducer( - connection.execute( - "SELECT * FROM deep_scan_workers WHERE scan_id = ?", (scan_id,) - ).fetchall() + + # Match the durable reducer sequence used by coordinator recovery. + def reducer_order(worker: Any) -> tuple[int, str]: + match = re.search(r"dedup-(\d+)", worker["prompt_path"]) + return (int(match[1]) if match else 0, worker["id"]) + + reducer = max( + ( + worker + for worker in connection.execute( + "SELECT * FROM deep_scan_workers WHERE scan_id = ?", (scan_id,) + ) + if worker["kind"] == "dedup" + and worker["status"] == "succeeded" + and worker["result_manifest_path"] + ), + key=reducer_order, + default=None, ) selected_result = reducer["result_manifest_path"] if reducer is not None else None if publication["resultPath"] != selected_result: diff --git a/plugins/codex-security/tests/test_deep_scan_publication_authority.py b/plugins/codex-security/tests/test_deep_scan_publication_authority.py index 99865d82d..1d17ca8be 100644 --- a/plugins/codex-security/tests/test_deep_scan_publication_authority.py +++ b/plugins/codex-security/tests/test_deep_scan_publication_authority.py @@ -10,7 +10,7 @@ from test_deep_scan_successful_publication import publication_scan as publication_scan -def stage_publication(scan, *, generation, result_path, title): +def stage_publication(scan, *, generation, result_path, title, complete=True): draft_dir = scan.scan_dir / "drafts" draft_dir.mkdir(exist_ok=True) draft_path = draft_dir / f"{uuid.uuid4()}.json" @@ -22,6 +22,8 @@ def stage_publication(scan, *, generation, result_path, title): "findings": {"findings": findings}, "coverage": scan.coverage, } + if not complete: + draft["manifest"]["scan"]["complete"] = False if generation is not None: draft["deepScanPublication"] = { "coordinatorGeneration": generation, @@ -40,9 +42,18 @@ def stage_publication(scan, *, generation, result_path, title): ) -@pytest.mark.parametrize("stale", ["generation", "aggregate", "unfenced"]) +@pytest.mark.parametrize( + ("stale", "complete"), + [ + ("generation", True), + ("generation", False), + ("aggregate", True), + ("aggregate", False), + ("unfenced", True), + ], +) def test_stale_coordinator_cannot_replace_newer_canonical_publication( - workbench_api, workbench_db, publication_scan, stale + workbench_api, workbench_db, publication_scan, stale, complete ): scan = publication_scan() old_result = add_worker(workbench_db, scan) @@ -52,11 +63,16 @@ def test_stale_coordinator_cannot_replace_newer_canonical_publication( "UPDATE deep_scan_runs SET coordinator_generation = 3 WHERE scan_id = ?", (scan.scan_id,), ) - for result, completed_at in ((old_result, "2026-01-01"), (new_result, "2026-01-02")): + for sequence, result in enumerate((old_result, new_result), start=1): workbench_db.execute( - "UPDATE deep_scan_workers SET kind = 'dedup', merge_state = 'none', completed_at = ? " + "UPDATE deep_scan_workers SET kind = 'dedup', merge_state = 'none', " + "prompt_path = ?, completed_at = ? " "WHERE result_manifest_path = ?", - (completed_at, str(result)), + ( + str(scan.scan_dir / f"dedup-{sequence:04d}" / "prompt.md"), + f"2026-01-0{sequence}", + str(result), + ), ) current = stage_publication( scan, generation=3, result_path=new_result, title="Current accepted aggregate" @@ -72,6 +88,7 @@ def test_stale_coordinator_cannot_replace_newer_canonical_publication( generation=None if stale == "unfenced" else 2 if stale == "generation" else 3, result_path=old_result if stale == "aggregate" else new_result, title="Superseded aggregate", + complete=complete, ) with pytest.raises(SystemExit, match="coordinator|aggregate"): From 4a5983f7aac613131d2ee8eec12a6b5170857136 Mon Sep 17 00:00:00 2001 From: Michael D'Angelo Date: Wed, 16 Sep 2026 09:23:03 +0000 Subject: [PATCH 4/4] Read reducer sequence from its worker directory --- .../mcp-app/src/deep-scan/coordinator.ts | 2 +- ...test_deep_scan_publication_integration.mjs | 26 ++++++++++++------- .../scripts/workbench_saved_results.py | 2 +- 3 files changed, 18 insertions(+), 12 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 efb8b48da..4e9dd0809 100644 --- a/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts +++ b/plugins/codex-security/mcp-app/src/deep-scan/coordinator.ts @@ -1184,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; } diff --git a/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs b/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs index 270b9466b..93bf5498d 100644 --- a/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs +++ b/plugins/codex-security/mcp-app/tests/test_deep_scan_publication_integration.mjs @@ -88,14 +88,18 @@ test("public partial drafts retain checkpoints after coordinator adoption", asyn assert.deepEqual(await readdir(path.join(run.scanDir, "drafts")), []); }); -for (const [label, offsets, ids] of [ +for (const [label, offsets, ids, paths, recoveryOnly] of [ ["increasing timestamps", [1, 2], [highId, lowId]], ["equal timestamps and descending UUIDs", [1, 1], [highId, lowId]], ["decreasing timestamps", [2, 1], [lowId, highId]], ["equal timestamps and ascending UUIDs", [1, 1], [lowId, highId]], + ["scan root collision, live publication", [1, 1], [highId, lowId], { scanRoot: "dedup-123/scans" }], + ["scan root collision, recovery", [1, 1], [highId, lowId], { scanRoot: "dedup-123/scans" }, true], + ["target name collision, live publication", [1, 1], [highId, lowId], { target: "dedup-456-target" }], + ["target name collision, recovery", [1, 1], [highId, lowId], { target: "dedup-456-target" }, true], ]) { test(`selected reducer survives recovery and public completion: ${label}`, async (t) => { - const fixture = await createFixture(t); + const fixture = await createFixture(t, paths); const { run, store, call, runWorkbench, instant } = fixture; assertSuccess(await call("record_codex_security_scan_draft", partial(run, "parent-checkpoint-only"))); const parentCheckpoints = await checkpoints(run); @@ -115,11 +119,13 @@ for (const [label, offsets, ids] of [ }, runWorkbench, undefined, { coordinatorGeneration, resultPath: results[index] }) ); - await publish(1); - const selected = await snapshot(run); - await assert.rejects(publish(0), /superseded publication selection/); - await assert.rejects(publish(1, 1), /newer generation/); - assert.deepEqual(await snapshot(run), selected); + if (!recoveryOnly) { + await publish(1); + const selected = await snapshot(run); + await assert.rejects(publish(0), /superseded publication selection/); + await assert.rejects(publish(1, 1), /newer generation/); + assert.deepEqual(await snapshot(run), selected); + } const publications = []; const coordinator = new DeepScanCoordinator({ @@ -159,21 +165,21 @@ for (const [label, offsets, ids] of [ }); } -async function createFixture(t) { +async function createFixture(t, paths = {}) { const root = await mkdtemp(path.join(tmpdir(), "deep-publication-")); let client; t.after(async () => { await client?.close(); await rm(root, { recursive: true, force: true }); }); - const targetPath = path.join(root, "target"); + const targetPath = path.join(root, paths.target ?? "target"); await mkdir(targetPath); await writeFile(path.join(targetPath, "fixture.py"), "# Synthetic publication fixture\n"); const environment = { ...process.env, CODEX_HOME: path.join(root, "home"), CODEX_SECURITY_STATE_DIR: path.join(root, "state"), - CODEX_SECURITY_SCAN_ROOT: path.join(root, "scans"), + CODEX_SECURITY_SCAN_ROOT: path.join(root, paths.scanRoot ?? "scans"), }; delete environment.CODEX_SECURITY_DEEP_SCAN_CONFIG_PATH; const instant = new Date().toISOString(); diff --git a/plugins/codex-security/scripts/workbench_saved_results.py b/plugins/codex-security/scripts/workbench_saved_results.py index f51ba74a5..abb428519 100644 --- a/plugins/codex-security/scripts/workbench_saved_results.py +++ b/plugins/codex-security/scripts/workbench_saved_results.py @@ -1412,7 +1412,7 @@ def _require_current_deep_publication( # Match the durable reducer sequence used by coordinator recovery. def reducer_order(worker: Any) -> tuple[int, str]: - match = re.search(r"dedup-(\d+)", worker["prompt_path"]) + match = re.search(r"dedup-(\d+)", Path(worker["prompt_path"]).parent.name) return (int(match[1]) if match else 0, worker["id"]) reducer = max(