diff --git a/plugins/codex-security/mcp-app/helpers-main.ts b/plugins/codex-security/mcp-app/helpers-main.ts index 3beb90680..0d36a5f41 100644 --- a/plugins/codex-security/mcp-app/helpers-main.ts +++ b/plugins/codex-security/mcp-app/helpers-main.ts @@ -6,6 +6,7 @@ import { normalizeCandidatesCommand } from "./src/helpers/normalize-candidates"; import { validatePatchRiskAssessmentCommand } from "./src/helpers/validate-patch-risk-assessment"; import { deepReviewInputCommand } from "./src/helpers/deep-review-input"; import { rankShardsCommand } from "./src/helpers/rank-shards"; +import { rankPoolCommand } from "./src/helpers/rank-pool"; let commandLine = process.argv.slice(2); if (process.platform === "win32") { @@ -48,9 +49,15 @@ if (command === "resolve-security-md") { command === "merge-rank-outputs" ) { process.exitCode = rankShardsCommand(command, args, posixHome); +} else if ( + command === "make-rank-pool-plan" || + command === "validate-rank-worker" || + command === "validate-rank-pool" +) { + process.exitCode = rankPoolCommand(command, args, posixHome); } else { console.error( - "Usage: launch_codex_security_mcp[.cmd] --helper [options]", + "Usage: launch_codex_security_mcp[.cmd] --helper [options]", ); process.exitCode = 2; } diff --git a/plugins/codex-security/mcp-app/src/helpers/normalize-candidates.ts b/plugins/codex-security/mcp-app/src/helpers/normalize-candidates.ts index 523733325..7e006f6e4 100644 --- a/plugins/codex-security/mcp-app/src/helpers/normalize-candidates.ts +++ b/plugins/codex-security/mcp-app/src/helpers/normalize-candidates.ts @@ -1,4 +1,5 @@ import { decodeUtf8 } from "./utf8"; +import { resolvedPath } from "./resolve-path"; import { createHash, randomBytes } from "node:crypto"; import { closeSync, @@ -11,12 +12,7 @@ import { writeFileSync, } from "node:fs"; import { basename, dirname, isAbsolute, join, relative, sep } from "node:path"; -import { - decodePosixBytes, - encodePosixPath, - resolvePosixPath, - SymlinkLoopError, -} from "./posix-path"; +import { encodePosixPath, SymlinkLoopError } from "./posix-path"; import { expandHome, HomeExpansionError, @@ -106,20 +102,6 @@ const stat = (path: string) => const pathKey = (value: string) => process.platform === "win32" ? value.toLowerCase() : value; -function resolvedPath(value: string, strict = true): string { - if (process.platform !== "win32") - return decodePosixBytes( - resolvePosixPath(encodePosixPath(parsedPath(value)), strict), - ); - try { - return pathText(windowsFiles().realpath(widePath(value), strict)); - } catch (error) { - if ((error as NodeJS.ErrnoException).code === "ELOOP") - throw new SymlinkLoopError(`Symlink loop from ${value}`); - throw error; - } -} - function inside(path: string, root: string, allowMissing = false): string { let result: string | undefined; if (windows) { diff --git a/plugins/codex-security/mcp-app/src/helpers/python-json.ts b/plugins/codex-security/mcp-app/src/helpers/python-json.ts index 0796749df..26ec23e8a 100644 --- a/plugins/codex-security/mcp-app/src/helpers/python-json.ts +++ b/plugins/codex-security/mcp-app/src/helpers/python-json.ts @@ -1,3 +1,5 @@ +import { decodeUtf8 } from "./utf8"; + // Preserve Python's integer/float distinction and arbitrary-size JSON integers. export class JsonFloat { constructor(readonly source: string) {} @@ -22,6 +24,54 @@ export function objectEntries(value: Row): [string, unknown][] { export class JsonSyntaxError extends Error {} +// Detect UTF-8/16/32 JSON bytes; UTF-8 input uses the shared strict decoder. +export function parseJsonBytes(bytes: Buffer): unknown { + let width = 1; + let little = true; + let offset = 0; + const prefix = bytes.subarray(0, 4).toString("hex"); + if (prefix === "fffe0000" || prefix === "0000feff") { + width = 4; + little = prefix === "fffe0000"; + offset = 4; + } else if (prefix.startsWith("fffe") || prefix.startsWith("feff")) { + width = 2; + little = prefix.startsWith("fffe"); + offset = 2; + } else if (prefix.startsWith("efbbbf")) { + offset = 3; + } else if (bytes.length >= 4) { + if (bytes[0] === 0) { + width = bytes[1] === 0 ? 4 : 2; + little = false; + } else if (bytes[1] === 0) { + width = bytes[2] || bytes[3] ? 2 : 4; + } + } else if (bytes.length === 2 && (bytes[0] === 0 || bytes[1] === 0)) { + width = 2; + little = bytes[0] !== 0; + } + let text = ""; + if (width === 1) { + text = decodeUtf8(bytes.subarray(offset)); + } else { + if ((bytes.length - offset) % width !== 0) + throw new Error(`Truncated UTF-${width * 8} JSON input`); + for (let index = offset; index < bytes.length; index += width) { + const point = + width === 2 + ? little + ? bytes.readUInt16LE(index) + : bytes.readUInt16BE(index) + : little + ? bytes.readUInt32LE(index) + : bytes.readUInt32BE(index); + text += String.fromCodePoint(point); + } + } + return parseJson(text); +} + export function parseJson(source: string, rejectDuplicates = false): unknown { const tokens = [ ...source.matchAll( diff --git a/plugins/codex-security/mcp-app/src/helpers/rank-pool.ts b/plugins/codex-security/mcp-app/src/helpers/rank-pool.ts new file mode 100644 index 000000000..6b574164e --- /dev/null +++ b/plugins/codex-security/mcp-app/src/helpers/rank-pool.ts @@ -0,0 +1,346 @@ +import { createHash } from "node:crypto"; +import { basename, dirname } from "node:path"; +import { mkdir, readFile, writeFile } from "./helper-files"; +import { + JsonSyntaxError, + object, + parseJsonBytes, + pythonRepr, +} from "./python-json"; +import { resolvedPath } from "./resolve-path"; +import { + childPath, + discoverInputShards, + shardNames, + validateShard, +} from "./rank-shards"; +import { + ArgumentError, + argumentsFor, + compare, + print, + worklistPath, +} from "./rank-worklists"; + +type Command = + | "make-rank-pool-plan" + | "validate-rank-worker" + | "validate-rank-pool"; +interface Worker { + input_shards: string[]; + output_shards: string[]; + slot: number; +} +const outputName = (name: string) => + name.replaceAll(".input.jsonl", ".output.jsonl"); +const same = (left: string[], right: string[]) => + left.length === right.length && + left.every((name, index) => name === right[index]); + +function requirePlanDirectory(plan: string, directory: string): void { + const expected = childPath(dirname(plan), "rank_shards"); + const key = (value: string) => { + const path = resolvedPath(value, false); + return process.platform === "win32" ? path.toLowerCase() : path; + }; + if (key(directory) !== key(expected)) + throw new Error( + `Rank shard directory must be the assignment plan's sibling rank_shards directory; expected=${expected}; actual=${directory}`, + ); +} + +function makePlan(directory: string, slots: bigint, output: string): void { + requirePlanDirectory(output, directory); + const inputs = discoverInputShards(directory).map((path) => basename(path)); + const count = Math.min(inputs.length, Number(slots), 6); + const workers = Array.from({ length: count }, (_, index): Worker => { + const assigned = inputs.filter( + (_, inputIndex) => inputIndex % count === index, + ); + return { + input_shards: assigned, + output_shards: assigned.map(outputName), + slot: index + 1, + }; + }); + const plan = { + ranking_worker_count: count, + schema_version: 1, + shard_count: inputs.length, + strategy: "round_robin", + workers, + }; + mkdir(dirname(output)); + writeFile(output, [ + Buffer.from( + (JSON.stringify(plan, null, 2) + "\n").replaceAll( + "\n", + process.platform === "win32" ? "\r\n" : "\n", + ), + ), + ]); + print( + `Assigned ${inputs.length} rank shards to ${count} ranking workers in ${output}`, + ); +} + +function integer(value: unknown, label: string, minimum: bigint): bigint { + if (typeof value !== "bigint" || value < minimum) + throw new Error(`${label} must be an integer of at least ${minimum}`); + return value; +} + +function strings(value: unknown, label: string): string[] { + if (!Array.isArray(value) || value.length === 0) + throw new Error(`${label} must be a non-empty list`); + if (value.some((item: unknown) => typeof item !== "string" || !item)) + throw new Error(`${label} entries must be non-empty strings`); + return value as string[]; +} + +function fields( + value: Record, + expected: string[], + label: string, +): void { + const actual = Object.keys(value); + const missing = expected + .filter((name) => !Object.hasOwn(value, name)) + .sort(compare); + const unexpected = actual + .filter((name) => !expected.includes(name)) + .sort(compare); + if (missing.length || unexpected.length) + throw new Error( + `${label} fields do not match schema; missing=${pythonRepr(missing)}; unexpected=${pythonRepr(unexpected)}`, + ); +} + +function validatePlan(plan: string, directory: string) { + requirePlanDirectory(plan, directory); + const misplaced = [ + ...shardNames(dirname(plan), "input"), + ...shardNames(dirname(plan), "output"), + ].sort(compare); + if (misplaced.length) + throw new Error( + `Rank shard artifacts must be stored in the assignment plan's sibling rank_shards directory; misplaced=${pythonRepr(misplaced)}`, + ); + const inputs = discoverInputShards(directory); + const inputNames = inputs.map((path) => basename(path)); + const outputNames = inputNames.map(outputName); + const bytes = readFile(plan); + let payload: unknown; + try { + payload = parseJsonBytes(bytes); + } catch (error) { + if (error instanceof JsonSyntaxError) + throw new Error( + `${plan}: invalid JSON: ${error.message.replace(/: line \d+ column \d+ \(char \d+\)$/u, "")}`, + ); + throw error; + } + if (!object(payload)) throw new Error(`${plan}: expected a JSON object`); + fields( + payload, + [ + "schema_version", + "strategy", + "shard_count", + "ranking_worker_count", + "workers", + ], + `${plan}: rank pool plan`, + ); + if (integer(payload.schema_version, `${plan}: schema_version`, 1n) !== 1n) + throw new Error(`${plan}: schema_version must be 1`); + if (payload.strategy !== "round_robin") + throw new Error(`${plan}: strategy must be round_robin`); + const shardCount = integer(payload.shard_count, `${plan}: shard_count`, 0n); + if (shardCount !== BigInt(inputs.length)) + throw new Error( + `${plan}: shard_count does not match input shards; plan=${shardCount}; actual=${inputs.length}`, + ); + const workerCount = integer( + payload.ranking_worker_count, + `${plan}: ranking_worker_count`, + 0n, + ); + if (shardCount > 0n && workerCount === 0n) + throw new Error( + `${plan}: ranking_worker_count must be at least 1 when input shards exist`, + ); + if (workerCount > shardCount) + throw new Error(`${plan}: ranking_worker_count cannot exceed shard_count`); + if (workerCount > 6n) + throw new Error(`${plan}: ranking_worker_count cannot exceed 6`); + if ( + !Array.isArray(payload.workers) || + BigInt(payload.workers.length) !== workerCount + ) + throw new Error( + `${plan}: workers must contain exactly ${workerCount} worker assignments`, + ); + const workers = payload.workers.map((raw: unknown, index): Worker => { + const label = `${plan}: workers[${index}]`; + if (!object(raw)) throw new Error(`${label} must be a JSON object`); + fields(raw, ["slot", "input_shards", "output_shards"], label); + const slot = integer(raw.slot, `${label}.slot`, 1n); + if (slot !== BigInt(index + 1)) + throw new Error(`${label}.slot must be ${index + 1}`); + const assignedInputs = strings(raw.input_shards, `${label}.input_shards`); + const assignedOutputs = strings( + raw.output_shards, + `${label}.output_shards`, + ); + if (assignedInputs.length !== assignedOutputs.length) + throw new Error( + `${label} input_shards and output_shards lengths must match`, + ); + if (!same(assignedOutputs, assignedInputs.map(outputName))) + throw new Error(`${label}.output_shards do not match its input_shards`); + return { + slot: Number(slot), + input_shards: assignedInputs, + output_shards: assignedOutputs, + }; + }); + for (const [index, worker] of workers.entries()) { + const assigned = inputNames.filter( + (_, inputIndex) => inputIndex % workers.length === index, + ); + if ( + !same(worker.input_shards, assigned) || + !same(worker.output_shards, assigned.map(outputName)) + ) + throw new Error( + `${plan}: worker slot ${index + 1} does not match the deterministic round_robin assignment`, + ); + } + return { inputs, outputNames, workers, bytes }; +} + +function validateWorker( + plan: string, + directory: string, + slotArgument: unknown, +): void { + const { workers, bytes } = validatePlan(plan, directory); + const slot = integer(slotArgument, "--slot", 1n); + if (slot > BigInt(workers.length)) + throw new Error(`--slot must be at most ${workers.length}`); + const worker = workers[Number(slot) - 1]!; + let rows = 0; + const digest = createHash("sha256"); + for (const [index, name] of worker.input_shards.entries()) { + const output = childPath(directory, worker.output_shards[index]!); + const [, outputs] = validateShard(childPath(directory, name), output); + const outputBytes = readFile(output); + rows += outputs.length; + digest + .update(worker.output_shards[index]!) + .update("\0") + .update(outputBytes) + .update("\0"); + } + print( + "RANK_WORKER_RECEIPT " + + JSON.stringify({ + output_shards: worker.output_shards.length, + outputs_sha256: digest.digest("hex"), + plan_sha256: createHash("sha256").update(bytes).digest("hex"), + ranking_worker_count: workers.length, + rows, + schema_version: 1, + slot: Number(slot), + status: "complete", + }), + ); +} + +function validatePool(plan: string, directory: string): void { + const { inputs, outputNames, workers } = validatePlan(plan, directory); + const actual = new Set(shardNames(directory, "output")); + const expected = new Set(outputNames); + const missing = outputNames.filter((name) => !actual.has(name)).sort(compare); + const unexpected = [...actual] + .filter((name) => !expected.has(name)) + .sort(compare); + if (missing.length || unexpected.length) + throw new Error( + `Rank pool outputs are incomplete; missing output shards=${pythonRepr(missing)}; unexpected output shards=${pythonRepr(unexpected)}`, + ); + let rows = 0; + for (const input of inputs) { + const [, outputs] = validateShard( + input, + childPath(directory, outputName(basename(input))), + ); + rows += outputs.length; + } + print( + `Validated ${workers.length} ranking workers, ${inputs.length} shards, and ${rows} ranking rows`, + ); +} + +export function rankPoolCommand( + command: Command, + args: string[], + posixHome = process.env.HOME, +): number { + const required = + command === "make-rank-pool-plan" + ? ["shard-dir", "usable-worker-slots", "out"] + : command === "validate-rank-worker" + ? ["plan", "shard-dir", "slot"] + : ["plan", "shard-dir"]; + const integers = + command === "make-rank-pool-plan" + ? ["usable-worker-slots"] + : command === "validate-rank-worker" + ? ["slot"] + : []; + const syntax = required + .map((name) => `--${name} ${integers.includes(name) ? "INT" : "PATH"}`) + .join(" "); + const usage = `usage: launch_codex_security_mcp[.cmd] --helper ${command} [-h] ${syntax}`; + try { + const values = argumentsFor(args, required, integers); + if (values.help) { + const description = + command === "make-rank-pool-plan" + ? "Assign rank shards to a deterministic bounded worker pool. Usable worker slots are capped at 6." + : command === "validate-rank-worker" + ? "Validate one assigned ranking-worker slot and emit its completion receipt." + : "Validate a rank pool plan and every assigned shard output."; + print( + `${usage}\n\n${description}\n\noptions:\n -h, --help show this help message and exit\n${required.map((name) => ` --${name} ${integers.includes(name) ? "INT" : "PATH"}`).join("\n")}`, + ); + return 0; + } + if (command === "make-rank-pool-plan") { + const slots = values["usable-worker-slots"] as bigint; + if (slots < 1n) + throw new Error("--usable-worker-slots must be at least 1"); + const directory = worklistPath(values["shard-dir"] as string, posixHome); + const output = worklistPath(values.out as string, posixHome); + makePlan(directory, slots, output); + } else { + const plan = worklistPath(values.plan as string, posixHome); + const directory = worklistPath(values["shard-dir"] as string, posixHome); + if (command === "validate-rank-worker") + validateWorker(plan, directory, values.slot); + else validatePool(plan, directory); + } + return 0; + } catch (error) { + const message = error instanceof Error ? error.message : String(error); + print( + error instanceof ArgumentError + ? `${usage}\n${command}: error: ${message}` + : message, + true, + ); + return error instanceof ArgumentError ? 2 : 1; + } +} diff --git a/plugins/codex-security/mcp-app/src/helpers/rank-shards.ts b/plugins/codex-security/mcp-app/src/helpers/rank-shards.ts index 1f950a82d..6947440dd 100644 --- a/plugins/codex-security/mcp-app/src/helpers/rank-shards.ts +++ b/plugins/codex-security/mcp-app/src/helpers/rank-shards.ts @@ -36,7 +36,10 @@ export function childPath(directory: string, name: string): string { ); } -function shardNames(directory: string, kind: "input" | "output"): string[] { +export function shardNames( + directory: string, + kind: "input" | "output", +): string[] { const matches = (name: string) => { // Python's Unicode case-insensitive globbing includes these ASCII equivalents. const matched = @@ -68,7 +71,7 @@ function shardNames(directory: string, kind: "input" | "output"): string[] { } } -function discoverInputShards(directory: string): string[] { +export function discoverInputShards(directory: string): string[] { let isDirectory = false; try { isDirectory = ( @@ -102,7 +105,10 @@ function discoverInputShards(directory: string): string[] { return expected.map((name) => childPath(directory, name)); } -function validateShard(input: string, output: string): [RankRow[], RankRow[]] { +export function validateShard( + input: string, + output: string, +): [RankRow[], RankRow[]] { const inputs = loadRankRows(input, false); requireUniquePaths(inputs, `Rank input shard ${basename(input)}`); const outputs = loadRankRows(output, true); diff --git a/plugins/codex-security/mcp-app/src/helpers/rank-worklists.ts b/plugins/codex-security/mcp-app/src/helpers/rank-worklists.ts index b3087035c..72ab587ba 100644 --- a/plugins/codex-security/mcp-app/src/helpers/rank-worklists.ts +++ b/plugins/codex-security/mcp-app/src/helpers/rank-worklists.ts @@ -140,7 +140,7 @@ export function argumentsFor( required: readonly string[], integerOptions: readonly string[] = [], ): Record { - const names = [...required, "help", ...integerOptions]; + const names = [...new Set([...required, "help", ...integerOptions])]; const values: Record = {}; const extra: string[] = []; const looksOptional = (arg: string) => diff --git a/plugins/codex-security/mcp-app/src/helpers/resolve-path.ts b/plugins/codex-security/mcp-app/src/helpers/resolve-path.ts new file mode 100644 index 000000000..017c8ef66 --- /dev/null +++ b/plugins/codex-security/mcp-app/src/helpers/resolve-path.ts @@ -0,0 +1,29 @@ +import { windowsBinding } from "../native"; +import { + pathText, + widePath, + windowsFileSystem, +} from "../../../native/windows-files.mjs"; +import { parsedPath } from "./resolve-security-md"; +import { + decodePosixBytes, + encodePosixPath, + resolvePosixPath, + SymlinkLoopError, +} from "./posix-path"; + +export function resolvedPath(value: string, strict = true): string { + if (process.platform !== "win32") + return decodePosixBytes( + resolvePosixPath(encodePosixPath(parsedPath(value)), strict), + ); + try { + return pathText( + windowsFileSystem(windowsBinding()).realpath(widePath(value), strict), + ); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ELOOP") + throw new SymlinkLoopError(`Symlink loop from ${value}`); + throw error; + } +} diff --git a/plugins/codex-security/native/README.md b/plugins/codex-security/native/README.md index 129de7f9c..855574bd4 100644 --- a/plugins/codex-security/native/README.md +++ b/plugins/codex-security/native/README.md @@ -41,7 +41,7 @@ Windows uses `windows-binding.mts` and the same Rust crate. `WindowsHandle` owns The binding exposes synchronous file and directory creation, attributes and reparse tags, identity and final/opened names, read/write/seek/size/EOF/flush, exact-handle rename and deletion, and exclusive whole-file locking. Rust's `File` supplies ordinary I/O, cursor-preserving truncation, `sync_all` for flush, and locks. Calls return numeric Windows errors, including 6 for closed handles and 33 for nonblocking lock contention. Buffer ranges, path encoding, and 64-bit seek arguments are checked before use. Overlapped handles are unsupported because pending operations could retain native buffers beyond the call. Path authorization, ancestor traversal, and reparse-point policy remain the caller's responsibility. -Five additional operations preserve Windows strings at the Node boundary. `windowsArguments` returns the complete OS argument vector, including the executable and Node options, using Rust's CRT-compatible parser. `windowsEnvironment` reads one wide environment name and distinguishes an absent value (`null`) from an empty buffer. `windowsAbsolutePath` resolves against the native current directory and drive directories without requiring the destination to exist. `windowsDirectoryEntries` uses `std::fs::read_dir` and cached `DirEntry::file_type()` values without opening each child; names remain UTF-16LE, and construction or iteration failures return their numeric Windows error and an empty array. Directory symlinks and junctions have both directory and symbolic-link flags. The typed adapter exposes this enumerator through `entriesWithTypes`, which `resolve-security-md --list` uses on Windows. `windowsReadLink` returns a UTF-16LE link target or its numeric Windows error; candidate normalization uses it to resolve missing paths without losing raw filenames. Assessment validation, deep-review worklists, and rank shards share these typed wide-path operations. +Five additional operations preserve Windows strings at the Node boundary. `windowsArguments` returns the complete OS argument vector, including the executable and Node options, using Rust's CRT-compatible parser. `windowsEnvironment` reads one wide environment name and distinguishes an absent value (`null`) from an empty buffer. `windowsAbsolutePath` resolves against the native current directory and drive directories without requiring the destination to exist. `windowsDirectoryEntries` uses `std::fs::read_dir` and cached `DirEntry::file_type()` values without opening each child; names remain UTF-16LE, and construction or iteration failures return their numeric Windows error and an empty array. Directory symlinks and junctions have both directory and symbolic-link flags. The typed adapter exposes this enumerator through `entriesWithTypes`, which `resolve-security-md --list` uses on Windows. `windowsReadLink` returns a UTF-16LE link target or its numeric Windows error; candidate normalization uses it to resolve missing paths without losing raw filenames. Assessment validation, deep-review worklists, and rank shard/pool helpers share these typed wide-path operations. `windows-files.mts` leaves ordinary absolute-path resolution and canonicalization to `GetFullPathNameW` and `GetFinalPathNameByHandleW`, trimming trailing separators below the root. Its small verbatim-path normalizer preserves drive and UNC share roots when resolving dot segments, including literal trailing dots and spaces. Non-strict `realpath` can retain unresolved components; callers must check containment independently. It also supports missing output paths. `stat(path, false)` retains exact symbolic-link and reparse-point metadata so callers can reject junction traversal independently of the enumerator's link label. The SDK's public runtime floor remains Node 22.13.0. Node 20.0.0 is an additional native-foundation compatibility proof; it does not change the SDK engine requirement. diff --git a/plugins/codex-security/native/examples/windows-wide-launcher.rs b/plugins/codex-security/native/examples/windows-wide-launcher.rs index 8ec31aa5e..09fe93b58 100644 --- a/plugins/codex-security/native/examples/windows-wide-launcher.rs +++ b/plugins/codex-security/native/examples/windows-wide-launcher.rs @@ -529,7 +529,71 @@ fn main() -> std::io::Result<()> { "Wide shard discovery or replacement paths changed", )); } - println!("{{\"policyHelperRawPaths\":true,\"candidateHelperRawPaths\":true,\"assessmentHelperRawPaths\":true,\"deepReviewHelperRawPaths\":true,\"rankShardHelperRawPaths\":true,\"directoryIdentity\":true}}"); + fs::remove_file(shards.join(&malformed_name))?; + let pool_name = raw("pool-", 0xd800); + let plan_name = raw("plan-", 0xdfff); + let pool = repo.join(&pool_name); + let pool_shards = pool.join("rank_shards"); + fs::create_dir(&pool)?; + fs::rename(&shards, &pool_shards)?; + let plan = pool.join(&plan_name); + let plan_sentinel = pool.join(raw("plan-", 0xfffd)); + fs::write(&plan_sentinel, "replacement plan sentinel")?; + let made_plan = shard_command( + "make-rank-pool-plan", + &[ + "--shard-dir".into(), + Path::new("~").join(&pool_name).join("rank_shards"), + "--usable-worker-slots".into(), + "2".into(), + "--out".into(), + Path::new(&pool_name).join(&plan_name), + ], + )?; + if !made_plan.status.success() || !made_plan.stderr.is_empty() { + return Err(io::Error::other(format!( + "Wide pool plan creation failed: {}", + String::from_utf8_lossy(&made_plan.stderr) + ))); + } + let plan_bytes = fs::read(&plan)?; + let expected_plan = concat!( + "{\r\n \"ranking_worker_count\": 2,\r\n \"schema_version\": 1,\r\n", + " \"shard_count\": 2,\r\n \"strategy\": \"round_robin\",\r\n \"workers\": [\r\n", + " {\r\n \"input_shards\": [\r\n \"rank-shard-0001.input.jsonl\"\r\n ],\r\n", + " \"output_shards\": [\r\n \"rank-shard-0001.output.jsonl\"\r\n ],\r\n \"slot\": 1\r\n },\r\n", + " {\r\n \"input_shards\": [\r\n \"rank-shard-0002.input.jsonl\"\r\n ],\r\n", + " \"output_shards\": [\r\n \"rank-shard-0002.output.jsonl\"\r\n ],\r\n \"slot\": 2\r\n }\r\n ]\r\n}\r\n", + ); + if plan_bytes != expected_plan.as_bytes() { + return Err(io::Error::other("Wide pool assignment bytes changed")); + } + let pool_args = [ + "--plan".into(), + plan.clone(), + "--shard-dir".into(), + Path::new(".").join(&pool_name).join("rank_shards"), + ]; + let mut worker_args = pool_args.to_vec(); + worker_args.extend(["--slot".into(), "1".into()]); + let worker = shard_command("validate-rank-worker", &worker_args)?; + let complete = shard_command("validate-rank-pool", &pool_args)?; + if !worker.status.success() + || !worker.stderr.is_empty() + || !worker.stdout.starts_with(b"RANK_WORKER_RECEIPT ") + || !complete.status.success() + || !complete.stderr.is_empty() + || complete.stdout != b"Validated 2 ranking workers, 2 shards, and 2 ranking rows\r\n" + || fs::read(&plan)? != plan_bytes + || fs::read(&plan_sentinel)? != b"replacement plan sentinel" + { + return Err(io::Error::other(format!( + "Wide pool validation failed: {}{}", + String::from_utf8_lossy(&worker.stderr), + String::from_utf8_lossy(&complete.stderr) + ))); + } + println!("{{\"policyHelperRawPaths\":true,\"candidateHelperRawPaths\":true,\"assessmentHelperRawPaths\":true,\"deepReviewHelperRawPaths\":true,\"rankShardHelperRawPaths\":true,\"rankPoolHelperRawPaths\":true,\"directoryIdentity\":true}}"); Ok(()) } diff --git a/plugins/codex-security/scripts/generate_rank_input.py b/plugins/codex-security/scripts/generate_rank_input.py index 21cbb1296..569432efd 100644 --- a/plugins/codex-security/scripts/generate_rank_input.py +++ b/plugins/codex-security/scripts/generate_rank_input.py @@ -8,24 +8,15 @@ - `make-diff-rank-input` creates the deterministic diff-scoped JSONL candidate worklist from Git changed paths. It supports committed revision diffs and local working-tree patches. -- `make-rank-pool-plan` assigns those shards to a deterministic bounded worker - pool. -- `validate-rank-worker` validates one worker slot and emits a content-bound - completion receipt. -- `validate-rank-pool` validates the pool plan and every assigned shard output. """ from __future__ import annotations import argparse -import hashlib import json import os -import re import subprocess import sys -from collections import Counter -from collections.abc import Callable from pathlib import Path # Some plugin hosts launch Python with safe-path isolation enabled. @@ -114,16 +105,8 @@ "yarn.lock", } -SHARD_INPUT_GLOB = "rank-shard-*.input.jsonl" -SHARD_OUTPUT_GLOB = "rank-shard-*.output.jsonl" -SHARD_INPUT_PATTERN = re.compile(r"^rank-shard-([0-9]{4,})\.input\.jsonl$") DIRECT_SCOPE_PREVIEW_READ_BYTES = 64 * 1024 -RANK_POOL_PLAN_SCHEMA_VERSION = 1 -RANK_POOL_STRATEGY = "round_robin" -RANK_POOL_WORKER_CAP = 6 JsonRow = dict[str, object] -RowValidator = Callable[[JsonRow, Path, int], None] -RankWorkerAssignment = tuple[int, list[str], list[str]] def parse_args() -> argparse.Namespace: @@ -195,39 +178,6 @@ def parse_args() -> argparse.Namespace: help=f"Maximum UTF-8 bytes in each preview. Defaults to {DEFAULT_PREVIEW_BYTES}.", ) - pool_plan = subparsers.add_parser( - "make-rank-pool-plan", - help="Assign rank shards to a deterministic bounded worker pool.", - ) - pool_plan.add_argument("--shard-dir", required=True, help="Directory of rank shards.") - pool_plan.add_argument( - "--usable-worker-slots", - required=True, - type=int, - help="Usable ranking-worker slots reported by capability preflight; capped at 6.", - ) - pool_plan.add_argument("--out", required=True, help="Output rank_worker_assignments.json path.") - - validate_worker = subparsers.add_parser( - "validate-rank-worker", - help="Validate one assigned ranking-worker slot and emit its completion receipt.", - ) - validate_worker.add_argument("--plan", required=True, help="Rank pool plan JSON path.") - validate_worker.add_argument("--shard-dir", required=True, help="Directory of rank shards.") - validate_worker.add_argument( - "--slot", - required=True, - type=int, - help="One-based ranking-worker slot from the rank pool plan.", - ) - - validate_pool = subparsers.add_parser( - "validate-rank-pool", - help="Validate a rank pool plan and every assigned shard output.", - ) - validate_pool.add_argument("--plan", required=True, help="Rank pool plan JSON path.") - validate_pool.add_argument("--shard-dir", required=True, help="Directory of rank shards.") - return parser.parse_args() @@ -302,14 +252,6 @@ def write_jsonl(output: Path, rows: list[JsonRow]) -> None: handle.write("\n") -def write_json(output: Path, payload: dict[str, object]) -> None: - output.parent.mkdir(parents=True, exist_ok=True) - output.write_text( - json.dumps(payload, ensure_ascii=False, indent=2, sort_keys=True) + "\n", - encoding="utf-8", - ) - - def load_scopes_file(scopes_file: Path) -> list[str]: try: loaded: object = json.loads(scopes_file.read_text(encoding="utf-8")) @@ -324,82 +266,6 @@ def load_scopes_file(scopes_file: Path) -> list[str]: return loaded -def load_jsonl(path: Path, label: str, validator: RowValidator) -> list[JsonRow]: - if not path.exists(): - raise SystemExit(f"{label} missing: {path}") - - rows: list[JsonRow] = [] - with path.open(encoding="utf-8") as handle: - for line_number, raw_line in enumerate(handle, start=1): - if not raw_line.strip(): - raise SystemExit(f"{path}:{line_number}: blank JSONL rows are not allowed") - try: - parsed: object = json.loads(raw_line) - except json.JSONDecodeError as exc: - raise SystemExit(f"{path}:{line_number}: invalid JSON: {exc.msg}") from exc - if not isinstance(parsed, dict): - raise SystemExit(f"{path}:{line_number}: expected a JSON object") - row = {str(key): value for key, value in parsed.items()} - validator(row, path, line_number) - rows.append(row) - return rows - - -def require_exact_fields(row: JsonRow, expected: set[str], path: Path, line_number: int) -> None: - actual = set(row) - if actual != expected: - missing = sorted(expected - actual) - unexpected = sorted(actual - expected) - details: list[str] = [] - if missing: - details.append(f"missing fields {missing}") - if unexpected: - details.append(f"unexpected fields {unexpected}") - raise SystemExit(f"{path}:{line_number}: {'; '.join(details)}") - - -def require_string( - row: JsonRow, field: str, path: Path, line_number: int, *, allow_empty: bool -) -> None: - value = row[field] - if not isinstance(value, str) or (not allow_empty and not value.strip()): - requirement = "a string" if allow_empty else "a non-empty string" - raise SystemExit(f"{path}:{line_number}: {field} must be {requirement}") - - -def validate_rank_input_row(row: JsonRow, path: Path, line_number: int) -> None: - require_exact_fields(row, {"path", "area", "preview"}, path, line_number) - require_string(row, "path", path, line_number, allow_empty=False) - require_string(row, "area", path, line_number, allow_empty=True) - require_string(row, "preview", path, line_number, allow_empty=True) - - -def validate_rank_output_row(row: JsonRow, path: Path, line_number: int) -> None: - require_exact_fields(row, {"path", "area", "score", "include", "reason"}, path, line_number) - require_string(row, "path", path, line_number, allow_empty=False) - require_string(row, "area", path, line_number, allow_empty=True) - score = row["score"] - if isinstance(score, bool) or not isinstance(score, int): - raise SystemExit(f"{path}:{line_number}: score must be an integer from 1 through 10") - if not 1 <= score <= 10: - raise SystemExit(f"{path}:{line_number}: score must be from 1 through 10") - if not isinstance(row["include"], bool): - raise SystemExit(f"{path}:{line_number}: include must be a boolean") - require_string(row, "reason", path, line_number, allow_empty=False) - - -def require_unique_paths(rows: list[JsonRow], label: str) -> None: - seen: set[str] = set() - duplicates: set[str] = set() - for row in rows: - path = str(row["path"]) - if path in seen: - duplicates.add(path) - seen.add(path) - if duplicates: - raise SystemExit(f"{label} contains duplicate paths: {sorted(duplicates)}") - - def make_repo_rank_input(args: argparse.Namespace) -> None: repo = Path(args.repo).expanduser().resolve() if not repo.is_dir(): @@ -683,326 +549,6 @@ def make_diff_rank_input(args: argparse.Namespace) -> None: print(f"Wrote {len(rows)} rows to {output}") -def discover_input_shards(shard_dir: Path) -> list[Path]: - if not shard_dir.is_dir(): - raise SystemExit(f"Rank shard directory missing: {shard_dir}") - - numbered_shards: list[tuple[int, Path]] = [] - for path in shard_dir.glob(SHARD_INPUT_GLOB): - match = SHARD_INPUT_PATTERN.fullmatch(path.name) - if match is None: - raise SystemExit(f"Rank input shard has invalid name: {path.name}") - numbered_shards.append((int(match.group(1)), path)) - numbered_shards.sort(key=lambda item: (item[0], item[1].name)) - input_shards = [path for _, path in numbered_shards] - expected_names = [ - f"rank-shard-{index:04d}.input.jsonl" for index in range(1, len(input_shards) + 1) - ] - actual_names = [path.name for path in input_shards] - if actual_names != expected_names: - raise SystemExit( - "Rank input shards must use contiguous canonical names; " - f"expected={expected_names}; actual={actual_names}" - ) - return input_shards - - -def output_name_for(input_name: str) -> str: - return input_name.replace(".input.jsonl", ".output.jsonl") - - -def require_plan_shard_dir(plan_path: Path, shard_dir: Path) -> None: - expected = plan_path.parent / "rank_shards" - if shard_dir.resolve() != expected.resolve(): - raise SystemExit( - "Rank shard directory must be the assignment plan's sibling rank_shards " - f"directory; expected={expected}; actual={shard_dir}" - ) - - -def require_no_misplaced_rank_shards(plan_path: Path) -> None: - misplaced = sorted( - ( - *plan_path.parent.glob(SHARD_INPUT_GLOB), - *plan_path.parent.glob(SHARD_OUTPUT_GLOB), - ), - key=lambda path: path.name, - ) - if misplaced: - raise SystemExit( - "Rank shard artifacts must be stored in the assignment plan's sibling " - f"rank_shards directory; misplaced={[path.name for path in misplaced]}" - ) - - -def make_rank_pool_plan(args: argparse.Namespace) -> None: - if args.usable_worker_slots < 1: - raise SystemExit("--usable-worker-slots must be at least 1") - - shard_dir = Path(args.shard_dir).expanduser() - output = Path(args.out).expanduser() - require_plan_shard_dir(output, shard_dir) - input_shards = discover_input_shards(shard_dir) - worker_count = min(len(input_shards), args.usable_worker_slots, RANK_POOL_WORKER_CAP) - workers: list[dict[str, object]] = [] - for worker_index in range(worker_count): - assigned_inputs = [path.name for path in input_shards[worker_index::worker_count]] - workers.append( - { - "slot": worker_index + 1, - "input_shards": assigned_inputs, - "output_shards": [output_name_for(name) for name in assigned_inputs], - } - ) - - plan: dict[str, object] = { - "schema_version": RANK_POOL_PLAN_SCHEMA_VERSION, - "strategy": RANK_POOL_STRATEGY, - "shard_count": len(input_shards), - "ranking_worker_count": worker_count, - "workers": workers, - } - write_json(output, plan) - print(f"Assigned {len(input_shards)} rank shards to {worker_count} ranking workers in {output}") - - -def load_rank_pool_plan(plan_path: Path) -> tuple[dict[str, object], bytes]: - if not plan_path.exists(): - raise SystemExit(f"Rank pool plan missing: {plan_path}") - plan_bytes = plan_path.read_bytes() - try: - payload: object = json.loads(plan_bytes) - except json.JSONDecodeError as exc: - raise SystemExit(f"{plan_path}: invalid JSON: {exc.msg}") from exc - if not isinstance(payload, dict): - raise SystemExit(f"{plan_path}: expected a JSON object") - return {str(key): value for key, value in payload.items()}, plan_bytes - - -def require_integer(value: object, label: str, *, minimum: int) -> int: - if isinstance(value, bool) or not isinstance(value, int) or value < minimum: - raise SystemExit(f"{label} must be an integer of at least {minimum}") - return value - - -def require_string_list(value: object, label: str) -> list[str]: - if not isinstance(value, list) or not value: - raise SystemExit(f"{label} must be a non-empty list") - if any(not isinstance(item, str) or not item for item in value): - raise SystemExit(f"{label} entries must be non-empty strings") - return [item for item in value if isinstance(item, str)] - - -def assignment_differences( - assigned_names: list[str], expected_names: list[str] -) -> tuple[list[str], list[str], list[str]]: - counts = Counter(assigned_names) - duplicates = sorted(name for name, count in counts.items() if count > 1) - assigned = set(assigned_names) - expected = set(expected_names) - return sorted(expected - assigned), duplicates, sorted(assigned - expected) - - -def validate_rank_pool_plan( - plan_path: Path, shard_dir: Path -) -> tuple[list[Path], list[str], list[RankWorkerAssignment], bytes]: - require_plan_shard_dir(plan_path, shard_dir) - require_no_misplaced_rank_shards(plan_path) - input_shards = discover_input_shards(shard_dir) - input_names = [path.name for path in input_shards] - output_names = [output_name_for(name) for name in input_names] - plan, plan_bytes = load_rank_pool_plan(plan_path) - expected_fields = { - "schema_version", - "strategy", - "shard_count", - "ranking_worker_count", - "workers", - } - actual_fields = set(plan) - if actual_fields != expected_fields: - raise SystemExit( - f"{plan_path}: rank pool plan fields do not match schema; " - f"missing={sorted(expected_fields - actual_fields)}; " - f"unexpected={sorted(actual_fields - expected_fields)}" - ) - schema_version = require_integer( - plan["schema_version"], f"{plan_path}: schema_version", minimum=1 - ) - if schema_version != RANK_POOL_PLAN_SCHEMA_VERSION: - raise SystemExit(f"{plan_path}: schema_version must be {RANK_POOL_PLAN_SCHEMA_VERSION}") - if plan["strategy"] != RANK_POOL_STRATEGY: - raise SystemExit(f"{plan_path}: strategy must be {RANK_POOL_STRATEGY}") - - shard_count = require_integer(plan["shard_count"], f"{plan_path}: shard_count", minimum=0) - if shard_count != len(input_shards): - raise SystemExit( - f"{plan_path}: shard_count does not match input shards; " - f"plan={shard_count}; actual={len(input_shards)}" - ) - worker_count = require_integer( - plan["ranking_worker_count"], f"{plan_path}: ranking_worker_count", minimum=0 - ) - if shard_count > 0 and worker_count == 0: - raise SystemExit( - f"{plan_path}: ranking_worker_count must be at least 1 when input shards exist" - ) - if worker_count > shard_count: - raise SystemExit(f"{plan_path}: ranking_worker_count cannot exceed shard_count") - if worker_count > RANK_POOL_WORKER_CAP: - raise SystemExit(f"{plan_path}: ranking_worker_count cannot exceed {RANK_POOL_WORKER_CAP}") - - workers = plan["workers"] - if not isinstance(workers, list) or len(workers) != worker_count: - raise SystemExit( - f"{plan_path}: workers must contain exactly {worker_count} worker assignments" - ) - - assigned_inputs: list[str] = [] - assigned_outputs: list[str] = [] - parsed_workers: list[RankWorkerAssignment] = [] - worker_fields = {"slot", "input_shards", "output_shards"} - for worker_index, raw_worker in enumerate(workers): - label = f"{plan_path}: workers[{worker_index}]" - if not isinstance(raw_worker, dict): - raise SystemExit(f"{label} must be a JSON object") - worker = {str(key): value for key, value in raw_worker.items()} - if set(worker) != worker_fields: - raise SystemExit( - f"{label} fields do not match schema; " - f"missing={sorted(worker_fields - set(worker))}; " - f"unexpected={sorted(set(worker) - worker_fields)}" - ) - slot = require_integer(worker["slot"], f"{label}.slot", minimum=1) - if slot != worker_index + 1: - raise SystemExit(f"{label}.slot must be {worker_index + 1}") - worker_inputs = require_string_list(worker["input_shards"], f"{label}.input_shards") - worker_outputs = require_string_list(worker["output_shards"], f"{label}.output_shards") - if len(worker_inputs) != len(worker_outputs): - raise SystemExit(f"{label} input_shards and output_shards lengths must match") - expected_worker_outputs = [output_name_for(name) for name in worker_inputs] - if worker_outputs != expected_worker_outputs: - raise SystemExit(f"{label}.output_shards do not match its input_shards") - assigned_inputs.extend(worker_inputs) - assigned_outputs.extend(worker_outputs) - parsed_workers.append((slot, worker_inputs, worker_outputs)) - - missing, duplicates, unexpected = assignment_differences(assigned_inputs, input_names) - if missing or duplicates or unexpected: - raise SystemExit( - f"{plan_path}: pool plan must assign each input shard exactly once; " - f"missing={missing}; duplicates={duplicates}; unexpected={unexpected}" - ) - missing, duplicates, unexpected = assignment_differences(assigned_outputs, output_names) - if missing or duplicates or unexpected: - raise SystemExit( - f"{plan_path}: pool plan must assign each output shard exactly once; " - f"missing={missing}; duplicates={duplicates}; unexpected={unexpected}" - ) - - for worker_index, (_, worker_inputs, worker_outputs) in enumerate(parsed_workers): - expected_inputs = input_names[worker_index::worker_count] - expected_outputs = output_names[worker_index::worker_count] - if worker_inputs != expected_inputs or worker_outputs != expected_outputs: - raise SystemExit( - f"{plan_path}: worker slot {worker_index + 1} does not match the deterministic " - f"{RANK_POOL_STRATEGY} assignment" - ) - return input_shards, output_names, parsed_workers, plan_bytes - - -def validate_rank_worker_command(args: argparse.Namespace) -> None: - plan_path = Path(args.plan).expanduser() - shard_dir = Path(args.shard_dir).expanduser() - _, _, workers, plan_bytes = validate_rank_pool_plan(plan_path, shard_dir) - - slot = require_integer(args.slot, "--slot", minimum=1) - worker_count = len(workers) - if slot > worker_count: - raise SystemExit(f"--slot must be at most {worker_count}") - - assigned_slot, input_names, output_names = workers[slot - 1] - if assigned_slot != slot: - raise SystemExit(f"{plan_path}: worker assignment for slot {slot} is inconsistent") - - row_count = 0 - outputs_digest = hashlib.sha256() - for input_name, output_name in zip(input_names, output_names, strict=True): - input_shard = shard_dir / input_name - output_shard = shard_dir / output_name - _, output_rows = validate_rank_shard(input_shard, output_shard) - output_bytes = output_shard.read_bytes() - row_count += len(output_rows) - outputs_digest.update(output_name.encode("utf-8")) - outputs_digest.update(b"\0") - outputs_digest.update(output_bytes) - outputs_digest.update(b"\0") - - receipt: dict[str, object] = { - "schema_version": 1, - "plan_sha256": hashlib.sha256(plan_bytes).hexdigest(), - "slot": slot, - "ranking_worker_count": worker_count, - "output_shards": len(output_names), - "rows": row_count, - "outputs_sha256": outputs_digest.hexdigest(), - "status": "complete", - } - print("RANK_WORKER_RECEIPT " + json.dumps(receipt, sort_keys=True, separators=(",", ":"))) - - -def validate_rank_pool_command(args: argparse.Namespace) -> None: - plan_path = Path(args.plan).expanduser() - shard_dir = Path(args.shard_dir).expanduser() - input_shards, expected_output_names, workers, _ = validate_rank_pool_plan(plan_path, shard_dir) - - actual_output_names = {path.name for path in shard_dir.glob(SHARD_OUTPUT_GLOB)} - expected_outputs = set(expected_output_names) - if actual_output_names != expected_outputs: - missing = sorted(expected_outputs - actual_output_names) - unexpected = sorted(actual_output_names - expected_outputs) - raise SystemExit( - "Rank pool outputs are incomplete; " - f"missing output shards={missing}; unexpected output shards={unexpected}" - ) - - row_count = 0 - for input_shard in input_shards: - output_shard = input_shard.with_name(output_name_for(input_shard.name)) - _, shard_outputs = validate_rank_shard(input_shard, output_shard) - row_count += len(shard_outputs) - print( - f"Validated {len(workers)} ranking workers, " - f"{len(input_shards)} shards, and {row_count} ranking rows" - ) - - -def validate_rank_shard( - input_shard: Path, output_shard: Path -) -> tuple[list[JsonRow], list[JsonRow]]: - shard_inputs = load_jsonl(input_shard, "Rank input shard", validate_rank_input_row) - require_unique_paths(shard_inputs, f"Rank input shard {input_shard.name}") - shard_outputs = load_jsonl(output_shard, "Rank output shard", validate_rank_output_row) - require_unique_paths(shard_outputs, f"Rank output shard {output_shard.name}") - - expected_paths = {str(row["path"]) for row in shard_inputs} - actual_paths = {str(row["path"]) for row in shard_outputs} - if expected_paths != actual_paths: - missing = sorted(expected_paths - actual_paths) - unknown = sorted(actual_paths - expected_paths) - raise SystemExit( - f"{output_shard}: paths do not match its input shard; " - f"missing={missing}; unknown={unknown}" - ) - - area_by_path = {str(row["path"]): row["area"] for row in shard_inputs} - for row in shard_outputs: - row_path = str(row["path"]) - if row["area"] != area_by_path[row_path]: - raise SystemExit(f"{output_shard}: area does not match rank input for {row_path}") - return shard_inputs, shard_outputs - - def main() -> None: args = parse_args() if args.command == "make-repo-rank-input": @@ -1013,12 +559,6 @@ def main() -> None: bind_repo_scopes(args) elif args.command == "make-diff-rank-input": make_diff_rank_input(args) - elif args.command == "make-rank-pool-plan": - make_rank_pool_plan(args) - elif args.command == "validate-rank-worker": - validate_rank_worker_command(args) - elif args.command == "validate-rank-pool": - validate_rank_pool_command(args) else: raise SystemExit(f"Unknown command: {args.command}") diff --git a/plugins/codex-security/tests/test_generate_rank_input.py b/plugins/codex-security/tests/test_generate_rank_input.py index 2d05495cb..29eb2e329 100644 --- a/plugins/codex-security/tests/test_generate_rank_input.py +++ b/plugins/codex-security/tests/test_generate_rank_input.py @@ -1,6 +1,5 @@ from __future__ import annotations -import hashlib import json import subprocess import sys @@ -33,24 +32,10 @@ def test_cli_loads_preview_helper_with_safe_path() -> None: assert "Codex Security scan worklist helper" in result.stdout -def write_jsonl(path: Path, rows: list[dict[str, object]]) -> None: - path.parent.mkdir(parents=True, exist_ok=True) - path.write_text( - "".join(f"{json.dumps(row, separators=(',', ':'))}\n" for row in rows), - encoding="utf-8", - ) - - def read_jsonl(path: Path) -> list[dict[str, object]]: return [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines()] -def read_json(path: Path) -> dict[str, object]: - payload = json.loads(path.read_text(encoding="utf-8")) - assert isinstance(payload, dict) - return payload - - def git(repo: Path, *args: str, input: str | None = None) -> str: result = subprocess.run( ["git", "-C", str(repo), *args], @@ -68,78 +53,6 @@ def initialize_repo(repo: Path) -> None: git(repo, "config", "user.name", "Codex Security Tests") -def make_rank_rows(count: int) -> list[dict[str, object]]: - return [ - {"path": f"src/file_{index:02d}.py", "area": "src", "preview": f"value = {index}"} - for index in range(count) - ] - - -def rank_result(row: dict[str, object], *, score: int = 5) -> dict[str, object]: - return { - "path": row["path"], - "area": row["area"], - "score": score, - "include": True, - "reason": "runtime surface", - } - - -def make_shards_and_pool_plan( - tmp_path: Path, *, shard_count: int = 5, usable_worker_slots: int = 2 -) -> tuple[Path, Path, Path]: - rank_input = tmp_path / "rank_input.jsonl" - rows = make_rank_rows(shard_count) - write_jsonl(rank_input, rows) - shard_dir = tmp_path / "rank_shards" - shard_dir.mkdir() - for index, row in enumerate(rows, start=1): - write_jsonl(shard_dir / f"rank-shard-{index:04d}.input.jsonl", [row]) - plan = tmp_path / "rank_worker_assignments.json" - run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - str(usable_worker_slots), - "--out", - str(plan), - ) - return rank_input, shard_dir, plan - - -def write_valid_shard_outputs(shard_dir: Path) -> None: - for input_shard in sorted(shard_dir.glob("*.input.jsonl")): - output_shard = input_shard.with_name(input_shard.name.replace(".input.", ".output.")) - write_jsonl(output_shard, [rank_result(row) for row in read_jsonl(input_shard)]) - - -def worker_shard_names(plan: Path, slot: int) -> tuple[list[str], list[str]]: - payload = read_json(plan) - workers = payload["workers"] - assert isinstance(workers, list) - worker = workers[slot - 1] - assert isinstance(worker, dict) - input_shards = worker["input_shards"] - output_shards = worker["output_shards"] - assert isinstance(input_shards, list) - assert isinstance(output_shards, list) - assert all(isinstance(name, str) for name in input_shards) - assert all(isinstance(name, str) for name in output_shards) - return input_shards, output_shards - - -def write_worker_shard_outputs(shard_dir: Path, plan: Path, slot: int) -> list[str]: - input_names, output_names = worker_shard_names(plan, slot) - for input_name, output_name in zip(input_names, output_names, strict=True): - input_shard = shard_dir / input_name - write_jsonl( - shard_dir / output_name, - [rank_result(row) for row in read_jsonl(input_shard)], - ) - return output_names - - def test_make_repo_rank_input_matches_golden_and_filters_noise(tmp_path: Path) -> None: repo = tmp_path / "repo" (repo / "src").mkdir(parents=True) @@ -853,479 +766,3 @@ def test_make_rank_input_decodes_bom_marked_utf16_source(tmp_path: Path, mode: s run_cli(*arguments, "--out", str(output)) assert {row["path"]: row["preview"] for row in read_jsonl(output)} == expected - - -def test_make_rank_pool_plan_is_deterministic_round_robin_and_exact_once( - tmp_path: Path, -) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - second_plan = tmp_path / "second_rank_worker_assignments.json" - - run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - "2", - "--out", - str(second_plan), - ) - - assert plan.read_bytes() == second_plan.read_bytes() - assert read_json(plan) == { - "schema_version": 1, - "strategy": "round_robin", - "shard_count": 5, - "ranking_worker_count": 2, - "workers": [ - { - "slot": 1, - "input_shards": [ - "rank-shard-0001.input.jsonl", - "rank-shard-0003.input.jsonl", - "rank-shard-0005.input.jsonl", - ], - "output_shards": [ - "rank-shard-0001.output.jsonl", - "rank-shard-0003.output.jsonl", - "rank-shard-0005.output.jsonl", - ], - }, - { - "slot": 2, - "input_shards": [ - "rank-shard-0002.input.jsonl", - "rank-shard-0004.input.jsonl", - ], - "output_shards": [ - "rank-shard-0002.output.jsonl", - "rank-shard-0004.output.jsonl", - ], - }, - ], - } - - -def test_make_rank_pool_plan_caps_workers_at_shard_count(tmp_path: Path) -> None: - _, _, plan = make_shards_and_pool_plan(tmp_path, shard_count=2, usable_worker_slots=8) - - payload = read_json(plan) - assert payload["shard_count"] == 2 - assert payload["ranking_worker_count"] == 2 - assert len(payload["workers"]) == 2 - - -def test_make_rank_pool_plan_caps_workers_at_six(tmp_path: Path) -> None: - _, _, plan = make_shards_and_pool_plan(tmp_path, shard_count=8, usable_worker_slots=12) - - payload = read_json(plan) - assert payload["shard_count"] == 8 - assert payload["ranking_worker_count"] == 6 - workers = payload["workers"] - assert isinstance(workers, list) - assert len(workers) == 6 - assert worker_shard_names(plan, 1)[0] == [ - "rank-shard-0001.input.jsonl", - "rank-shard-0007.input.jsonl", - ] - assert worker_shard_names(plan, 2)[0] == [ - "rank-shard-0002.input.jsonl", - "rank-shard-0008.input.jsonl", - ] - - -def test_empty_rank_pool_closes_with_zero_shards_and_workers(tmp_path: Path) -> None: - shard_dir = tmp_path / "rank_shards" - shard_dir.mkdir() - - plan = tmp_path / "rank_worker_assignments.json" - plan_result = run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - "6", - "--out", - str(plan), - ) - assert plan_result.stdout == f"Assigned 0 rank shards to 0 ranking workers in {plan}\n" - assert read_json(plan) == { - "schema_version": 1, - "strategy": "round_robin", - "shard_count": 0, - "ranking_worker_count": 0, - "workers": [], - } - - pool_result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - ) - assert pool_result.stdout == "Validated 0 ranking workers, 0 shards, and 0 ranking rows\n" - - -def test_make_rank_pool_plan_requires_sibling_rank_shards_directory(tmp_path: Path) -> None: - shard_dir = tmp_path / "other_shards" - for index, row in enumerate(make_rank_rows(2), start=1): - write_jsonl(shard_dir / f"rank-shard-{index:04d}.input.jsonl", [row]) - plan = tmp_path / "rank_worker_assignments.json" - - result = run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - "2", - "--out", - str(plan), - check=False, - ) - - assert result.returncode != 0 - assert "must be the assignment plan's sibling rank_shards directory" in result.stderr - - -def test_validate_rank_pool_requires_sibling_rank_shards_directory(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - other_shard_dir = tmp_path / "other_shards" - shard_dir.rename(other_shard_dir) - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(other_shard_dir), - check=False, - ) - - assert result.returncode != 0 - assert "must be the assignment plan's sibling rank_shards directory" in result.stderr - - -@pytest.mark.parametrize( - "misplaced_name", - [ - "rank-shard-0001.input.jsonl", - "rank-shard-0001.output.jsonl", - ], -) -def test_validate_rank_pool_rejects_shards_beside_assignment_plan( - tmp_path: Path, misplaced_name: str -) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - write_valid_shard_outputs(shard_dir) - (tmp_path / misplaced_name).write_text("{}\n", encoding="utf-8") - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - - assert result.returncode != 0 - assert "must be stored in the assignment plan's sibling rank_shards directory" in result.stderr - assert f"misplaced=['{misplaced_name}']" in result.stderr - - -def test_make_rank_pool_plan_rejects_invalid_slots_and_shard_names(tmp_path: Path) -> None: - _, shard_dir, _ = make_shards_and_pool_plan(tmp_path) - - result = run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - "0", - "--out", - str(tmp_path / "invalid.json"), - check=False, - ) - assert result.returncode != 0 - assert "--usable-worker-slots must be at least 1" in result.stderr - - (shard_dir / "rank-shard-0005.input.jsonl").rename(shard_dir / "rank-shard-0006.input.jsonl") - result = run_cli( - "make-rank-pool-plan", - "--shard-dir", - str(shard_dir), - "--usable-worker-slots", - "2", - "--out", - str(tmp_path / "invalid.json"), - check=False, - ) - assert result.returncode != 0 - assert "contiguous canonical names" in result.stderr - - -def test_validate_rank_pool_rejects_malformed_and_tampered_plan(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - original_plan = read_json(plan) - - plan.write_text("{not json}\n", encoding="utf-8") - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - assert result.returncode != 0 - assert "invalid JSON" in result.stderr - - workers = original_plan["workers"] - assert isinstance(workers, list) - first_worker = workers[0] - assert isinstance(first_worker, dict) - input_shards = first_worker["input_shards"] - output_shards = first_worker["output_shards"] - assert isinstance(input_shards, list) - assert isinstance(output_shards, list) - input_shards[1] = input_shards[0] - output_shards[1] = output_shards[0] - plan.write_text(json.dumps(original_plan) + "\n", encoding="utf-8") - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - assert result.returncode != 0 - assert "assign each input shard exactly once" in result.stderr - assert "duplicates=['rank-shard-0001.input.jsonl']" in result.stderr - - -def test_validate_rank_worker_emits_content_bound_receipt_for_only_its_slot( - tmp_path: Path, -) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - output_names = write_worker_shard_outputs(shard_dir, plan, slot=1) - - result = run_cli( - "validate-rank-worker", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - "--slot", - "1", - ) - - outputs_digest = hashlib.sha256() - for output_name in output_names: - outputs_digest.update(output_name.encode("utf-8")) - outputs_digest.update(b"\0") - outputs_digest.update((shard_dir / output_name).read_bytes()) - outputs_digest.update(b"\0") - expected = { - "schema_version": 1, - "plan_sha256": hashlib.sha256(plan.read_bytes()).hexdigest(), - "slot": 1, - "ranking_worker_count": 2, - "output_shards": 3, - "rows": 3, - "outputs_sha256": outputs_digest.hexdigest(), - "status": "complete", - } - expected_line = "RANK_WORKER_RECEIPT " + json.dumps( - expected, sort_keys=True, separators=(",", ":") - ) - assert result.stdout == expected_line + "\n" - assert result.stderr == "" - assert not (shard_dir / "rank-shard-0002.output.jsonl").exists() - assert not (shard_dir / "rank-shard-0004.output.jsonl").exists() - - -def test_validate_rank_worker_rejects_missing_assigned_output(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - output_names = write_worker_shard_outputs(shard_dir, plan, slot=1) - missing_output = shard_dir / output_names[1] - missing_output.unlink() - - result = run_cli( - "validate-rank-worker", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - "--slot", - "1", - check=False, - ) - - assert result.returncode != 0 - assert f"Rank output shard missing: {missing_output}" in result.stderr - assert result.stdout == "" - - -def test_validate_rank_worker_rejects_invalid_assigned_output(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - output_names = write_worker_shard_outputs(shard_dir, plan, slot=2) - invalid_output = shard_dir / output_names[0] - invalid_output.write_text("{not json}\n", encoding="utf-8") - - result = run_cli( - "validate-rank-worker", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - "--slot", - "2", - check=False, - ) - - assert result.returncode != 0 - assert f"{invalid_output}:1: invalid JSON" in result.stderr - assert result.stdout == "" - - -def test_validate_rank_worker_reuses_strict_plan_and_path_validation(tmp_path: Path) -> None: - _, _, plan = make_shards_and_pool_plan(tmp_path) - other_shard_dir = tmp_path / "other_shards" - other_shard_dir.mkdir() - - result = run_cli( - "validate-rank-worker", - "--plan", - str(plan), - "--shard-dir", - str(other_shard_dir), - "--slot", - "1", - check=False, - ) - - assert result.returncode != 0 - assert "must be the assignment plan's sibling rank_shards directory" in result.stderr - - -def test_validate_rank_pool_accepts_complete_multi_shard_workers(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - write_valid_shard_outputs(shard_dir) - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - ) - - assert "Validated 2 ranking workers, 5 shards, and 5 ranking rows" in result.stdout - - -def test_rank_pool_accepts_parent_completion_for_an_unstarted_worker_slot( - tmp_path: Path, -) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path, shard_count=5, usable_worker_slots=2) - write_worker_shard_outputs(shard_dir, plan, slot=1) - write_worker_shard_outputs(shard_dir, plan, slot=2) - - for slot in (1, 2): - receipt = run_cli( - "validate-rank-worker", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - "--slot", - str(slot), - ) - assert receipt.stdout.startswith("RANK_WORKER_RECEIPT ") - - pool = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - ) - assert "Validated 2 ranking workers, 5 shards, and 5 ranking rows" in pool.stdout - - -def test_validate_rank_pool_rejects_missing_and_unexpected_outputs(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - write_valid_shard_outputs(shard_dir) - missing_output = shard_dir / "rank-shard-0005.output.jsonl" - missing_output.unlink() - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - assert result.returncode != 0 - assert "missing output shards=['rank-shard-0005.output.jsonl']" in result.stderr - - input_row = read_jsonl(shard_dir / "rank-shard-0005.input.jsonl")[0] - write_jsonl(missing_output, [rank_result(input_row)]) - write_jsonl( - shard_dir / "rank-shard-9999.output.jsonl", - [rank_result(input_row)], - ) - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - assert result.returncode != 0 - assert "unexpected output shards=['rank-shard-9999.output.jsonl']" in result.stderr - - -def test_validate_rank_pool_rejects_one_bad_shard_from_multi_shard_worker( - tmp_path: Path, -) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - write_valid_shard_outputs(shard_dir) - bad_output = shard_dir / "rank-shard-0003.output.jsonl" - bad_output.write_text("{not json}\n", encoding="utf-8") - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - - assert result.returncode != 0 - assert "rank-shard-0003.output.jsonl:1: invalid JSON" in result.stderr - - -def test_validate_rank_pool_rejects_duplicate_output_rows(tmp_path: Path) -> None: - _, shard_dir, plan = make_shards_and_pool_plan(tmp_path) - write_valid_shard_outputs(shard_dir) - output = shard_dir / "rank-shard-0003.output.jsonl" - row = read_jsonl(output)[0] - write_jsonl(output, [row, row]) - - result = run_cli( - "validate-rank-pool", - "--plan", - str(plan), - "--shard-dir", - str(shard_dir), - check=False, - ) - - assert result.returncode != 0 - assert "contains duplicate paths" in result.stderr diff --git a/sdk/typescript/tests-ts/rank-pool.test.ts b/sdk/typescript/tests-ts/rank-pool.test.ts new file mode 100644 index 000000000..29e76ae5b --- /dev/null +++ b/sdk/typescript/tests-ts/rank-pool.test.ts @@ -0,0 +1,657 @@ +import { spawnSync } from "node:child_process"; +import { createHash } from "node:crypto"; +import { + existsSync, + mkdirSync, + mkdtempSync, + readFileSync, + readdirSync, + realpathSync, + rmSync, + statSync, + symlinkSync, + writeFileSync, +} from "node:fs"; +import { tmpdir } from "node:os"; +import { join, sep } from "node:path"; +import { afterEach, describe, expect, test } from "bun:test"; +import { PLUGIN_ROOT } from "./plugin-root.js"; + +const node = Bun.which("node")!; +const helper = join(PLUGIN_ROOT, "mcp", "helpers.mjs"); +const newline = process.platform === "win32" ? "\r\n" : "\n"; +const roots: string[] = []; +interface Worker { + slot: number; + input_shards: string[]; + output_shards: string[]; +} +interface Plan { + schema_version: number; + strategy: string; + shard_count: number; + ranking_worker_count: number; + workers: Worker[]; +} +const shardName = (index: number, output = false) => + `rank-shard-${String(index).padStart(4, "0")}.${output ? "output" : "input"}.jsonl`; +const row = (index: number) => ({ + path: `src/file_${index}.py`, + area: "src", + preview: `value = ${index}`, +}); +const ranked = (index: number) => ({ + path: row(index).path, + area: "src", + score: 5, + include: true, + reason: "runtime surface", +}); +function fixture(count = 5) { + const root = realpathSync(mkdtempSync(join(tmpdir(), "rank-pool-"))); + roots.push(root); + const directory = join(root, "rank_shards"); + mkdirSync(directory); + for (let index = 1; index <= count; index++) + writeFileSync( + join(directory, shardName(index)), + JSON.stringify(row(index)) + "\n", + ); + return { root, directory, plan: join(root, "assignments.json") }; +} +type Fixture = ReturnType; +function run(f: Fixture, command: string, args: string[], env = process.env) { + return spawnSync(node, [helper, command, ...args], { + cwd: f.root, + env, + encoding: "utf8", + input: "stdin is not a plan", + }); +} +function make(f: Fixture, slots = "2", extra: string[] = []) { + return run(f, "make-rank-pool-plan", [ + "--shard-dir", + f.directory, + "--usable-worker-slots", + slots, + "--out", + f.plan, + ...extra, + ]); +} +function validate(f: Fixture, slot?: string, extra: string[] = []) { + return run( + f, + slot === undefined ? "validate-rank-pool" : "validate-rank-worker", + [ + "--plan", + f.plan, + "--shard-dir", + f.directory, + ...(slot === undefined ? [] : ["--slot", slot]), + ...extra, + ], + ); +} +function plan(f: Fixture): Plan { + return JSON.parse(readFileSync(f.plan, "utf8")) as Plan; +} +function change(f: Fixture, edit: (value: Plan) => void) { + const value = plan(f); + edit(value); + writeFileSync(f.plan, JSON.stringify(value)); +} +function complete(f: Fixture, slots?: number[]) { + for (const worker of plan(f).workers) { + if (slots && !slots.includes(worker.slot)) continue; + for (const name of worker.output_shards) { + const index = Number(name.match(/([0-9]+)\.output/u)![1]); + writeFileSync( + join(f.directory, name), + JSON.stringify(ranked(index)) + "\n", + ); + } + } +} +function encoded(text: string, encoding: string): Buffer { + const little = encoding.endsWith("le"); + if (encoding === "utf8") return Buffer.from(text); + if (encoding === "utf8-bom") + return Buffer.concat([Buffer.from([0xef, 0xbb, 0xbf]), Buffer.from(text)]); + if (encoding.startsWith("utf16")) { + const bytes = Buffer.from(text, "utf16le"); + return little ? bytes : bytes.swap16(); + } + const points = Array.from(text, (character) => character.codePointAt(0)!); + const bytes = Buffer.alloc(points.length * 4); + points.forEach((point, index) => + little + ? bytes.writeUInt32LE(point, index * 4) + : bytes.writeUInt32BE(point, index * 4), + ); + return bytes; +} +afterEach(() => { + for (const root of roots.splice(0)) + rmSync(root, { recursive: true, force: true }); +}); + +describe("rank pool helpers", () => { + test("writes stable sorted JSON and assigns every shard round-robin", () => { + const f = fixture(); + expect(make(f).stdout).toBe( + `Assigned 5 rank shards to 2 ranking workers in ${f.plan}${newline}`, + ); + const expected: Plan = { + ranking_worker_count: 2, + schema_version: 1, + shard_count: 5, + strategy: "round_robin", + workers: [ + { + input_shards: [1, 3, 5].map((index) => shardName(index)), + output_shards: [1, 3, 5].map((index) => shardName(index, true)), + slot: 1, + }, + { + input_shards: [2, 4].map((index) => shardName(index)), + output_shards: [2, 4].map((index) => shardName(index, true)), + slot: 2, + }, + ], + }; + expect(readFileSync(f.plan, "utf8")).toBe( + (JSON.stringify(expected, null, 2) + "\n").replaceAll("\n", newline), + ); + const first = readFileSync(f.plan); + expect(make(f).status).toBe(0); + expect(readFileSync(f.plan)).toEqual(first); + complete(f); + expect(validate(f).stdout).toBe( + `Validated 2 ranking workers, 5 shards, and 5 ranking rows${newline}`, + ); + }); + + test.each([ + [0, "6", 0], + [2, "8", 2], + [8, "12", 6], + [8, "9".repeat(400), 6], + ] as const)( + "caps %i shards with %s usable slots at %i workers", + (count, slots, workers) => { + const f = fixture(count); + expect(make(f, slots).status).toBe(0); + expect(plan(f).ranking_worker_count).toBe(workers); + expect(plan(f).workers).toHaveLength(workers); + complete(f); + expect(validate(f).status).toBe(0); + if (count === 0) + expect(validate(f, "1").stderr).toBe( + `--slot must be at most 0${newline}`, + ); + }, + ); + + test("validates one worker independently and binds its receipt to raw plan and output bytes", () => { + const f = fixture(); + make(f); + complete(f, [1]); + const names = [1, 3, 5].map((index) => shardName(index, true)); + const output = join(f.directory, names[1]!); + writeFileSync(output, JSON.stringify(ranked(3), null, 0) + "\r\n"); + const bytes = Buffer.concat([ + Buffer.from([0xef, 0xbb, 0xbf]), + readFileSync(f.plan), + ]); + writeFileSync(f.plan, bytes); + const digest = createHash("sha256"); + for (const name of names) + digest + .update(name) + .update("\0") + .update(readFileSync(join(f.directory, name))) + .update("\0"); + const result = validate(f, "1"); + expect(result.status).toBe(0); + expect(result.stderr).toBe(""); + expect(result.stdout).toBe( + "RANK_WORKER_RECEIPT " + + JSON.stringify({ + output_shards: 3, + outputs_sha256: digest.digest("hex"), + plan_sha256: createHash("sha256").update(bytes).digest("hex"), + ranking_worker_count: 2, + rows: 3, + schema_version: 1, + slot: 1, + status: "complete", + }) + + newline, + ); + expect(existsSync(join(f.directory, shardName(2, true)))).toBe(false); + expect(validate(f).stderr).toContain("missing output shards"); + const second = run(f, "validate-rank-worker", [ + "--plan", + f.plan, + "--shard-dir", + f.directory, + "--slot", + "2", + ]); + expect(second.status).toBe(1); + expect(second.stdout).toBe(""); + }); + + test("accepts completed slots without requiring worker receipts or unrelated outputs", () => { + const f = fixture(); + make(f); + complete(f); + writeFileSync(join(f.directory, "unrelated.txt"), "untouched"); + expect(validate(f, "1").status).toBe(0); + expect(validate(f, "2").status).toBe(0); + expect(validate(f).status).toBe(0); + writeFileSync( + join(f.directory, shardName(2, true)), + "invalid unassigned output", + ); + expect(validate(f, "1").status).toBe(0); + expect(validate(f).status).toBe(1); + }); + + test.each(["missing", "invalid JSON", "duplicate", "area", "paths"])( + "rejects %s assigned output", + (kind) => { + const f = fixture(); + make(f); + complete(f); + const output = join(f.directory, shardName(3, true)); + if (kind === "missing") rmSync(output); + else if (kind === "invalid JSON") writeFileSync(output, "{bad json}\n"); + else if (kind === "duplicate") + writeFileSync(output, (JSON.stringify(ranked(3)) + "\n").repeat(2)); + else + writeFileSync( + output, + JSON.stringify({ + ...ranked(3), + [kind === "area" ? "area" : "path"]: "different", + }) + "\n", + ); + expect(validate(f, "1").status).toBe(1); + expect(validate(f, "1").stdout).toBe(""); + expect(validate(f).status).toBe(1); + }, + ); + + test("reports both missing and unexpected output shards", () => { + const f = fixture(); + make(f); + complete(f); + rmSync(join(f.directory, shardName(5, true))); + writeFileSync(join(f.directory, shardName(9999, true)), ""); + expect(validate(f).stderr).toBe( + `Rank pool outputs are incomplete; missing output shards=['${shardName(5, true)}']; unexpected output shards=['${shardName(9999, true)}']${newline}`, + ); + }); + + test.each(["input", "output"])( + "rejects a misplaced %s shard before loading the plan", + (kind) => { + const f = fixture(); + const name = shardName(1, kind === "output"); + writeFileSync(join(f.root, name), ""); + expect(validate(f).stderr).toContain(`misplaced=['${name}']`); + expect(validate(f, "1").status).toBe(1); + // Plan creation has no misplaced-artifact check in the original helper. + expect(make(f).status).toBe(0); + }, + ); + + test("requires the sibling shard directory and resolves aliases with parent traversal", () => { + const f = fixture(); + mkdirSync(join(f.root, "other")); + const bad = { ...f, directory: join(f.root, "other") }; + expect(make(bad).stderr).toContain( + "must be the assignment plan's sibling rank_shards directory", + ); + expect(validate(bad, "1").stderr).toContain( + "must be the assignment plan's sibling rank_shards directory", + ); + symlinkSync( + f.directory, + join(f.root, "alias"), + process.platform === "win32" ? "junction" : "dir", + ); + expect(make({ ...f, directory: join(f.root, "alias") }).status).toBe(0); + complete(f); + expect( + validate({ + ...f, + plan: f.directory + sep + ".." + sep + "assignments.json", + }).status, + ).toBe(0); + expect( + validate({ + ...f, + directory: f.root + sep + "missing" + sep + ".." + sep + "rank_shards", + }).status, + ).toBe(process.platform === "win32" ? 0 : 1); + }); + + test("checks shard existence and names before loading a missing or malformed plan", () => { + const f = fixture(0); + expect(validate(f).stderr).toContain(f.plan); + writeFileSync(f.plan, "bad"); + writeFileSync(join(f.directory, "rank-shard-001.input.jsonl"), ""); + expect(validate(f).stderr).toContain("contiguous canonical names"); + expect(make(f).status).toBe(1); + expect(readFileSync(f.plan, "utf8")).toBe("bad"); + rmSync(f.directory, { recursive: true }); + expect(make(f).stderr).toContain("Rank shard directory missing"); + }); + + test.each([ + [ + "root fields", + (p: Plan) => { + Object.assign(p, { extra: true }); + }, + "unexpected=['extra']", + ], + [ + "version type", + (p: Plan) => { + Object.assign(p, { schema_version: true }); + }, + "schema_version must be an integer", + ], + [ + "version", + (p: Plan) => { + p.schema_version = 2; + }, + "schema_version must be 1", + ], + [ + "strategy", + (p: Plan) => { + p.strategy = "other"; + }, + "strategy must be round_robin", + ], + [ + "shard count", + (p: Plan) => { + p.shard_count = 4; + }, + "shard_count does not match", + ], + [ + "worker count", + (p: Plan) => { + p.ranking_worker_count = 0; + }, + "must be at least 1", + ], + [ + "workers type", + (p: Plan) => { + Object.assign(p, { workers: {} }); + }, + "exactly 2 worker assignments", + ], + [ + "worker object", + (p: Plan) => { + p.workers[0] = null as unknown as Worker; + }, + "must be a JSON object", + ], + [ + "worker fields", + (p: Plan) => { + Object.assign(p.workers[0]!, { extra: true }); + }, + "unexpected=['extra']", + ], + [ + "slot", + (p: Plan) => { + p.workers[0]!.slot = 2; + }, + ".slot must be 1", + ], + [ + "empty assignments", + (p: Plan) => { + p.workers[0]!.input_shards = []; + }, + "must be a non-empty list", + ], + [ + "empty name", + (p: Plan) => { + p.workers[0]!.input_shards[0] = ""; + }, + "entries must be non-empty strings", + ], + [ + "unequal assignments", + (p: Plan) => { + p.workers[0]!.input_shards.pop(); + }, + "lengths must match", + ], + [ + "output name", + (p: Plan) => { + p.workers[0]!.output_shards[0] = "wrong"; + }, + "do not match its input_shards", + ], + [ + "duplicate assignments", + (p: Plan) => { + p.workers[0]!.input_shards[1] = shardName(1); + p.workers[0]!.output_shards[1] = shardName(1, true); + }, + "does not match the deterministic round_robin assignment", + ], + [ + "round robin", + (p: Plan) => { + p.workers[0]!.input_shards.reverse(); + p.workers[0]!.output_shards.reverse(); + }, + "does not match the deterministic round_robin assignment", + ], + ] as const)("rejects tampered %s", (_name, edit, message) => { + const f = fixture(); + make(f); + change(f, edit); + const result = validate(f); + expect(result.status).toBe(1); + expect(result.stderr).toContain(message); + expect(validate(f, "1").stdout).toBe(""); + }); + + test.each(["utf8", "utf8-bom", "utf16le", "utf16be", "utf32le", "utf32be"])( + "accepts %s plan bytes and hashes the original encoding", + (encoding) => { + const f = fixture(1); + make(f); + complete(f); + for (const bom of encoding.startsWith("utf16") || + encoding.startsWith("utf32") + ? [false, true] + : [false]) { + const text = JSON.stringify({ + ranking_worker_count: 1, + schema_version: 1, + shard_count: 1, + strategy: "round_robin", + workers: [ + { + input_shards: [shardName(1)], + output_shards: [shardName(1, true)], + slot: 1, + }, + ], + }); + const bytes = encoded((bom ? "\ufeff" : "") + text, encoding); + writeFileSync(f.plan, bytes); + const result = validate(f, "1"); + expect(result.status).toBe(0); + expect( + JSON.parse(result.stdout.slice("RANK_WORKER_RECEIPT ".length)) + .plan_sha256, + ).toBe(createHash("sha256").update(bytes).digest("hex")); + expect(validate(f).status).toBe(0); + } + }, + ); + + test.each(["{bad}", "[]", "null", "1", "true", "", "{} trailing"])( + "rejects malformed or nonobject plan %s", + (text) => { + const f = fixture(); + writeFileSync(f.plan, text); + expect(validate(f).status).toBe(1); + expect(validate(f).stdout).toBe(""); + }, + ); + + test("rejects malformed UTF-8 even in an overwritten plan value", () => { + const f = fixture(0); + make(f); + const rest = readFileSync(f.plan, "utf8").slice(1); + writeFileSync(f.plan, '{"strategy":"\\ud800",' + rest); + expect(validate(f).status).toBe(0); + writeFileSync( + f.plan, + Buffer.concat([ + Buffer.from('{"strategy":"'), + Buffer.from([0xed, 0xa0, 0x80]), + Buffer.from('",' + rest), + ]), + ); + const result = validate(f); + expect(result.status).toBe(1); + expect(result.stdout).toBe(""); + }); + + test("preserves duplicate-key last value and integer-versus-float validation", () => { + const f = fixture(0); + make(f); + const text = readFileSync(f.plan, "utf8"); + writeFileSync( + f.plan, + text.replace( + '"schema_version": 1', + '"schema_version": 2, "schema_version": 1', + ), + ); + expect(validate(f).status).toBe(0); + for (const value of ["1.0", "1e0", "true", "NaN", "Infinity", "-1"]) { + writeFileSync( + f.plan, + text.replace('"schema_version": 1', `"schema_version": ${value}`), + ); + expect(validate(f).stderr).toContain( + "schema_version must be an integer of at least 1", + ); + } + }); + + test.each(["+2_0", "20", "٢٠", " 20 "])( + "accepts required integer argument %s with existing abbreviations", + (value) => { + const f = fixture(1); + expect(make(f, value).status).toBe(0); + complete(f); + expect(validate(f, "1", ["--sl=+1"]).status).toBe(0); + }, + ); + + test("preserves argument failure statuses and validation order", () => { + const f = fixture(); + for (const value of ["0", "-1"]) + expect(make({ ...f, directory: "missing" }, value).stderr).toBe( + `--usable-worker-slots must be at least 1${newline}`, + ); + for (const value of ["1.0", "1__0", "\u001c20"]) + expect(make(f, value).status).toBe(2); + for (const command of [ + "make-rank-pool-plan", + "validate-rank-worker", + "validate-rank-pool", + ]) { + expect(run(f, command, []).status).toBe(2); + expect(run(f, command, ["--help"]).status).toBe(0); + } + expect(validate(f, "0").stderr).toContain(f.plan); + make(f); + expect(validate(f, "0").stderr).toBe( + `--slot must be an integer of at least 1${newline}`, + ); + expect(validate(f, "9".repeat(400)).stderr).toBe( + `--slot must be at most 2${newline}`, + ); + expect(validate(f, "1", ["--s", "1"]).status).toBe(2); + }); + + test("preserves literal dash paths, output aliases, and existing file permissions", () => { + const f = fixture(0); + const dash = { ...f, directory: "rank_shards", plan: "-" }; + writeFileSync(join(f.root, "-"), "previous", { mode: 0o640 }); + expect(make(dash).status).toBe(0); + expect(validate(dash).status).toBe(0); + if (process.platform !== "win32") + expect(statSync(join(f.root, "-")).mode & 0o777).toBe(0o640); + symlinkSync( + join(f.root, "-"), + f.plan, + process.platform === "win32" ? "file" : undefined, + ); + expect(make(f).status).toBe(0); + expect(readFileSync(f.plan)).toEqual(readFileSync(join(f.root, "-"))); + expect(readdirSync(f.directory)).toHaveLength(0); + }); + + test.skipIf(process.platform === "win32")( + "uses raw POSIX bytes for plan and shard paths", + () => { + const f = fixture(0); + const pathBytes = + process.platform === "darwin" ? Buffer.from("é") : Buffer.from([255]); + const suffix = Array.from( + pathBytes, + (byte) => `\\${byte.toString(8).padStart(3, "0")}`, + ).join(""); + const raw = Buffer.concat([Buffer.from(f.root + "/"), pathBytes]); + mkdirSync(raw); + mkdirSync(Buffer.concat([raw, Buffer.from("/rank_shards")])); + const launcher = join( + PLUGIN_ROOT, + "scripts", + "launch_codex_security_mcp", + ); + const result = spawnSync( + "bash", + [ + "-c", + `r="$1"/$(printf '${suffix}'); "$2" --helper make-rank-pool-plan --shard-dir "$r/rank_shards" --usable-worker-slots 1 --out "$r/plan.json"`, + "rank-pool", + f.root, + launcher, + ], + { env: { ...process.env, CODEX_MCP_NODE_PATH: node } }, + ); + expect(result.status).toBe(0); + const output = Buffer.concat([raw, Buffer.from("/plan.json")]); + expect(JSON.parse(readFileSync(output, "utf8")).shard_count).toBe(0); + expect(result.stdout.includes(output)).toBe(true); + }, + ); +});