From f7cab9ccd44113757e229ef6861faaca20ab3d76 Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 03:42:43 -0600 Subject: [PATCH 1/6] Restore reviewed bounded-launch candidate for collision hardening --- conductor/README.md | 15 +- conductor/borg-conductor.mjs | 12 +- conductor/docs/LAUNCH-RELIABILITY.md | 86 +++++++ conductor/router/receipt-reader-worker.mjs | 72 ++++++ conductor/router/receipt-reader.mjs | 70 ++++++ conductor/router/router.mjs | 237 ++++++++++++-------- conductor/tests/receipt-reader.test.mjs | 56 +++++ conductor/tests/router-reliability.test.mjs | 205 +++++++++++++++++ 8 files changed, 652 insertions(+), 101 deletions(-) create mode 100644 conductor/docs/LAUNCH-RELIABILITY.md create mode 100644 conductor/router/receipt-reader-worker.mjs create mode 100644 conductor/router/receipt-reader.mjs create mode 100644 conductor/tests/receipt-reader.test.mjs create mode 100644 conductor/tests/router-reliability.test.mjs diff --git a/conductor/README.md b/conductor/README.md index b811aeb..8e3f4cc 100644 --- a/conductor/README.md +++ b/conductor/README.md @@ -57,10 +57,17 @@ remaining usable native allowance; earliest reset is only the tie break. Actual provider exhaustion and provider spend controls have separate evidence codes. No discretionary reservation or allowance floor is created. -Dispatch holds a private lock, rechecks the whole configured fleet immediately -before selection, persists an intent and receipt before native lifecycle calls, -and refuses duplicate `{workId,cwd}` intents across process restarts. An -ambiguous thread or turn response is recorded as `DO_NOT_RETRY` evidence. +Dispatch holds a private lock, persists an intent before admission scans, +rechecks the configured fleet before selection, and refuses duplicate +`{workId,cwd}` intents across process restarts. Active workspace and work-ID +claims are checked before native lifecycle calls. An ambiguous thread or turn +response remains `DO_NOT_RETRY` evidence. + +Use `route-status --cwd ABS --work-id ID` to inspect saved phase history and +native IDs without contacting providers or launching work. Admission/provider +stages and read-only receipt scans have bounded deadlines; unknown outcomes +never permit automatic replay. See [launch reliability and recovery](docs/LAUNCH-RELIABILITY.md) +for states, scan bounds, timeout controls and remaining filesystem/release limits. Capacity and claims support local OS observation or an explicitly configured absolute JSON-producing command. Remote machines therefore require an diff --git a/conductor/borg-conductor.mjs b/conductor/borg-conductor.mjs index 92e901c..fca9b23 100644 --- a/conductor/borg-conductor.mjs +++ b/conductor/borg-conductor.mjs @@ -12,7 +12,7 @@ import { validateInstallConfig, } from './config.mjs'; import { startConductor } from './conductor.mjs'; -import { accountIdentityDigest, dispatch, nativeConductorProvider, rank } from './router/router.mjs'; +import { accountIdentityDigest, dispatch, inspectDispatch, nativeConductorProvider, rank } from './router/router.mjs'; const SEAT_POLICY = `# New owner seat policy @@ -171,7 +171,8 @@ function usage() { borg-conductor status [--config ABS|conductors/config.json] borg-conductor auth status|pin|login [--config ABS|conductors/config.json] [--lane ID] borg-conductor rank [--config ABS|conductors/config.json] [--capability tools|reasoning] [--model MODEL] - borg-conductor route --config ABS|conductors/config.json --cwd ABS --prompt-file ABS --work-id ID [--role leaf|lead] [--model MODEL] [--effort EFFORT] + borg-conductor route-status --config ABS|conductors/config.json --cwd ABS --work-id ID + borg-conductor route --config ABS|conductors/config.json --cwd ABS --prompt-file ABS --work-id ID [--role leaf|lead] [--model MODEL] [--effort EFFORT] [--stage-timeout-ms N] `; } @@ -281,6 +282,12 @@ export async function main(argv = process.argv.slice(2)) { }), null, 2)}\n`); return; } + if (args.command === 'route-status') { + process.stdout.write(`${JSON.stringify(await inspectDispatch(config, { + cwd: canonicalAbsolute(args.cwd, 'cwd'), workId: args.workId, + }), null, 2)}\n`); + return; + } if (args.command === 'route') { const promptPath = canonicalAbsolute(args.promptFile, 'prompt-file'); const prompt = fs.readFileSync(promptPath, 'utf8'); @@ -292,6 +299,7 @@ export async function main(argv = process.argv.slice(2)) { capability: args.capability, model: args.model, effort: args.effort, + stageTimeoutMs: args.stageTimeoutMs === undefined ? undefined : Number(args.stageTimeoutMs), }); process.stdout.write(`${JSON.stringify({ ...result.receipt, receiptPath: result.receiptPath }, null, 2)}\n`); return; diff --git a/conductor/docs/LAUNCH-RELIABILITY.md b/conductor/docs/LAUNCH-RELIABILITY.md new file mode 100644 index 0000000..8bf7772 --- /dev/null +++ b/conductor/docs/LAUNCH-RELIABILITY.md @@ -0,0 +1,86 @@ +# Worker launch reliability and recovery + +The router records an intent before admission scans and provider calls. Inspect it +without another dispatch or provider probe: + +```sh +borg-conductor route-status --config /absolute/BORG_HOME/conductors/config.json \ + --cwd /absolute/workspace --work-id stable-work-id +``` + +Use the same canonical workspace and stable work ID supplied to `route`. This is +read-only reconciliation: it does not create, resume, cancel or complete a worker. +The returned receipt contains a bounded phase history, attempt ID, selected lane, +native thread/turn IDs when proved, and error classification. Prompt text is not +persisted, only its SHA-256. `completionVerified` is always false: a start receipt +is not evidence that a task finished, passed review or was deployed. + +## States and decisions + +| State | Meaning | Recovery | +| --- | --- | --- | +| `ATTEMPTING` | Intent exists; inspect phase (`PREPARING`, `ADMISSION`, `RECHECK`, `THREAD_START_PENDING`). | Do not start another copy. | +| `PRE_START_FAILED` | The recorded attempt failed before a native lifecycle request. | Diagnose the recorded phase. `noStartProven` applies only to this attempt, not all possible workers. | +| `THREAD_STARTED` | Native thread ID exists; turn start is pending. | Reconcile that exact thread. | +| `UNKNOWN_DO_NOT_RETRY` | Thread request may have executed; response was lost, late or invalid. | Verify native state before recovery; no automatic replay. | +| `STARTED_TURN_UNKNOWN` | Thread exists; turn outcome is unknown. | Reconcile exact thread/turn history before any continuation. | +| `DISPATCHED` | Native thread and turn IDs were returned. | Monitor actual worker outcome; do not call the task complete. | +| `NOT_FOUND` (status lookup) | No matching local intent was found. | **Not** proof that no worker started elsewhere. `noStartProven` remains false. | + +Identical intents remain blocked on replay. No new retry, fallback, or no-start +claim is inferred from a timeout or a missing receipt. Existing admission selects +another eligible candidate before native mutation when a candidate is unusable; +uncertain lifecycle calls are not retried on a different account or machine. + +## Admission and ownership + +Missing, Boolean, empty, string and non-finite numeric telemetry are unknown, not +safe zeroes. The same rule applies to native usage measurements. Completed +transport dispatches and uncertain lifecycle states remain active claims until +explicitly reconciled. A different work ID cannot evade an existing workspace +claim; a different workspace cannot evade an active work ID. + +A receipt scan is all-or-error. Malformed JSON, unsafe file permissions, symlinks, +unknown receipt states, oversize files or exceeded bounds never become an empty +or partial successful claims list. The fixed read-only reader runs in a separate +Node process, outside the router's filesystem worker pool. It accepts no shell +commands and never starts agents or writes state. + +Reader bounds: 5 seconds by default, 10,000 directory entries, 65,536 bytes per +receipt, 4 MiB total receipt bytes, and 8 MiB captured output. It scans only the +requested directory, not a recursive tree. A timeout reports +`RECEIPT_READ_TIMEOUT` with its stage and target. Signalling the reader is not +reported as proof that the OS process exited. No incomplete scan permits launch. + +## Deadlines and limits + +`route --stage-timeout-ms N` bounds each admission/provider operation (1–60,000 +milliseconds); the default is `routing.timeoutMs`. Independent machine probes +run concurrently. Late thread responses cannot advance to turn creation after +the caller records an unknown outcome. + +These are **stage** deadlines, not an end-to-end service-level guarantee. Workspace +validation, dispatch-lock I/O and durable state writes can still be delayed by an +unhealthy filesystem. State writes are not raced against a timeout, since a late +write could corrupt the evidence used by a subsequent attempt. Existing lock +contention controls remain. The router never removes an ownership check or +claims it cancelled an external lifecycle request just because its wait ended. + +## Verification + +```sh +node --test conductor/tests/*.test.mjs conductor/providers/*.test.mjs +``` + +The regression suite covers missing telemetry/usage, active dispatched claims, +workspace/work-ID conflicts, unsafe/corrupt/oversize receipts, native reader +timeout, pre-start diagnostic persistence, lifecycle timeout with late response, +duplicate suppression, status lookup and the real command-line status path. +Optional native Codex initialization uses `BORG_TEST_CODEX_BIN` with an isolated +unauthenticated profile; it is not a paid inference or production dispatch test. + +This portable-router change does not silently replace an estate's separately +installed fleet router, publish new MCP schemas, or install a goal-driven team +controller. An integrator must adopt the exact reviewed source, preserve its +existing machine/account/ownership policy, verify the deployed artifact, and +refresh client tool catalogs separately where needed. diff --git a/conductor/router/receipt-reader-worker.mjs b/conductor/router/receipt-reader-worker.mjs new file mode 100644 index 0000000..d4adcc3 --- /dev/null +++ b/conductor/router/receipt-reader-worker.mjs @@ -0,0 +1,72 @@ +// Fixed read-only child for receipt-reader.mjs. Never dispatches or mutates state. +import fs from 'node:fs'; +import path from 'node:path'; + +const MAX_ENTRIES = 10000; +const MAX_FILE_BYTES = 65536; +const MAX_TOTAL_BYTES = 4 * 1024 * 1024; +function reject(code) { const error = new Error(code); error.code = code; throw error; } +function privateStat(target, directory = false) { + const stat = fs.lstatSync(target); + if (stat.isSymbolicLink() || (directory ? !stat.isDirectory() : !stat.isFile()) || (stat.mode & 0o077) !== 0) reject('UNSAFE_RECEIPT'); + return stat; +} +function record(target) { + const stat = privateStat(target); + if (stat.size > MAX_FILE_BYTES) reject('RECEIPT_FILE_LIMIT'); + const fd = fs.openSync(target, fs.constants.O_RDONLY | (fs.constants.O_NOFOLLOW ?? 0) | (fs.constants.O_NONBLOCK ?? 0)); + try { + const actual = fs.fstatSync(fd); + if (!actual.isFile() || (actual.mode & 0o077) !== 0 || actual.ino !== stat.ino || actual.dev !== stat.dev) reject('UNSAFE_RECEIPT'); + const buffer = Buffer.alloc(MAX_FILE_BYTES + 1); + let bytes = 0; + while (bytes < buffer.length) { + const count = fs.readSync(fd, buffer, bytes, buffer.length - bytes, null); + if (!count) break; + bytes += count; + } + if (bytes > MAX_FILE_BYTES) reject('RECEIPT_FILE_LIMIT'); + let value; + try { value = JSON.parse(buffer.subarray(0, bytes).toString('utf8')); } + catch { reject('RECEIPT_JSON_INVALID'); } + if (!value || typeof value !== 'object' || Array.isArray(value)) reject('RECEIPT_JSON_INVALID'); + return { value, bytes }; + } finally { fs.closeSync(fd); } +} +function directory(target) { + privateStat(target, true); + const handle = fs.opendirSync(target); + const values = []; + let seen = 0; let total = 0; + try { + let entry; + while ((entry = handle.readSync()) !== null) { + if (++seen > MAX_ENTRIES) reject('RECEIPT_ENTRY_LIMIT'); + if (!entry.name.endsWith('.json')) continue; + if (!entry.isFile() || entry.isSymbolicLink()) reject('UNSAFE_RECEIPT'); + const read = record(path.join(target, entry.name)); + total += read.bytes; + if (total > MAX_TOTAL_BYTES) reject('RECEIPT_TOTAL_LIMIT'); + values.push(read.value); + } + } finally { handle.closeSync(); } + return values; +} +const [mode, target] = process.argv.slice(2); +try { + if (!['record','directory'].includes(mode) || !path.isAbsolute(target ?? '')) reject('RECEIPT_REQUEST_INVALID'); + let value; + // Only a missing requested root is empty/not-found. A receipt disappearing + // during enumeration makes the WHOLE scan unknown, not partially successful. + try { fs.lstatSync(target); } + catch (error) { + if (error.code !== 'ENOENT') throw error; + process.stdout.write(JSON.stringify({ ok: true, value: mode === 'directory' ? [] : null })); + process.exit(0); + } + value = mode === 'directory' ? directory(target) : record(target).value; + process.stdout.write(JSON.stringify({ ok: true, value })); +} catch (error) { + process.stdout.write(JSON.stringify({ ok: false, code: /^[A-Z_]+$/.test(error.code ?? '') ? error.code : 'RECEIPT_READ_FAILED' })); + process.exitCode = 1; +} diff --git a/conductor/router/receipt-reader.mjs b/conductor/router/receipt-reader.mjs new file mode 100644 index 0000000..d0be981 --- /dev/null +++ b/conductor/router/receipt-reader.mjs @@ -0,0 +1,70 @@ +import { spawn } from 'node:child_process'; +import { fileURLToPath } from 'node:url'; + +const MAX_OUTPUT_BYTES = 8 * 1024 * 1024; +const WORKER = fileURLToPath(new URL('./receipt-reader-worker.mjs', import.meta.url)); + +// Keep potentially stuck OS directory/file reads out of the router's libuv pool. +// This helper is strictly read-only. A failed or partial scan NEVER means no claims. +function readInChild(mode, target, options = {}) { + const timeoutMs = options.timeoutMs ?? 5000; + if (!Number.isInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 60000) { + throw new Error('RECEIPT_READ_TIMEOUT_INVALID'); + } + return new Promise((resolve, reject) => { + const child = spawn(process.execPath, [WORKER, mode, target], { + stdio: ['ignore', 'pipe', 'pipe'], windowsHide: true, + }); + let settled = false; + let bytes = 0; + const chunks = []; + let timer; + function finish(error, result) { + if (settled) return; + settled = true; + clearTimeout(timer); + if (error) { + // Sending SIGKILL is not proof of exit. The result makes no exit claim; + // unref/destroy prevents an uninterruptible read retaining this caller. + if (child.exitCode === null && child.signalCode === null) child.kill('SIGKILL'); + child.stdout.destroy(); child.stderr.destroy(); child.unref(); + error.targetPath = target; + error.stage = 'RECEIPT_READ'; + reject(error); + } else resolve(result); + } + function fail(code) { + const error = new Error(`${code}: ${target}`); + error.code = code; + finish(error); + } + timer = setTimeout(() => fail('RECEIPT_READ_TIMEOUT'), timeoutMs); + child.on('error', () => fail('RECEIPT_READER_UNAVAILABLE')); + child.stdout.on('data', chunk => { + bytes += chunk.length; + if (bytes > MAX_OUTPUT_BYTES) fail('RECEIPT_OUTPUT_LIMIT'); + else chunks.push(chunk); + }); + // Never echo diagnostic stderr, which may contain local/private content. + child.stderr.on('data', () => {}); + child.on('close', (code) => { + if (settled) return; + let envelope; + try { envelope = JSON.parse(Buffer.concat(chunks).toString('utf8')); } + catch { fail('RECEIPT_READER_INVALID_OUTPUT'); return; } + if (code !== 0 || envelope.ok !== true) { + const errorCode = /^[A-Z_]+$/.test(envelope.code ?? '') ? envelope.code : 'RECEIPT_READ_FAILED'; + fail(errorCode); return; + } + finish(null, envelope.value); + }); + }); +} + +export async function readPrivateJsonDirectory(directory, options = {}) { + return readInChild('directory', directory, options); +} + +export async function readPrivateJsonRecord(target, options = {}) { + return readInChild('record', target, options); +} diff --git a/conductor/router/router.mjs b/conductor/router/router.mjs index 498669f..3b6a53d 100644 --- a/conductor/router/router.mjs +++ b/conductor/router/router.mjs @@ -6,6 +6,7 @@ import path from 'node:path'; import { promisify } from 'node:util'; import { validateInstallConfig } from '../config.mjs'; +import { readPrivateJsonDirectory, readPrivateJsonRecord } from './receipt-reader.mjs'; // Portable extraction of the supplied conductor-usage-router and // fleet-placement invariants. Estate-specific roster, SSH shipping, and @@ -14,6 +15,7 @@ import { validateInstallConfig } from '../config.mjs'; const execFile = promisify(execFileCallback); const ACTIVE_RECEIPT_STATES = new Set([ + 'DISPATCHED', 'ATTEMPTING', 'THREAD_STARTED', 'UNKNOWN_DO_NOT_RETRY', @@ -45,7 +47,7 @@ function timestampMillis(value) { } function finite(value, minimum, maximum) { - const number = Number(value); + const number = typeof value === 'number' ? value : NaN; return Number.isFinite(number) && number >= minimum && number <= maximum ? number : null; } @@ -129,36 +131,30 @@ export async function observeMachineCapacity(machine, nowMs = Date.now()) { } async function readJsonFiles(directory) { - let entries; - try { - entries = await fs.readdir(directory, { withFileTypes: true }); - } catch (error) { - if (error.code === 'ENOENT') return []; - throw error; - } - const values = []; - for (const entry of entries) { - if (!entry.isFile() || !entry.name.endsWith('.json')) continue; - const target = path.join(directory, entry.name); - const stat = await fs.lstat(target); - if (!stat.isFile() || stat.isSymbolicLink() || (stat.mode & 0o077) !== 0) { - throw new Error(`unsafe private receipt: ${target}`); - } - values.push(JSON.parse(await fs.readFile(target, 'utf8'))); - } - return values; + return readPrivateJsonDirectory(directory); } export async function observeClaims(machine, statePath, nowMs = Date.now()) { if (machine.claims.kind === 'command') return commandJson(machine.claims); if (machine.claims.kind !== 'router-state') throw new Error('CLAIMS_KIND_UNSUPPORTED'); const receipts = await readJsonFiles(path.join(statePath, 'dispatch-receipts')); + const terminal = new Set(['PRE_START_FAILED', 'COMPLETED', 'FAILED', 'CANCELLED', 'RECONCILED_NO_START']); + for (const receipt of receipts) { + if (!ACTIVE_RECEIPT_STATES.has(receipt.state) && !terminal.has(receipt.state)) { + throw new Error('RECEIPT_STATE_UNKNOWN'); + } + if (ACTIVE_RECEIPT_STATES.has(receipt.state) && (typeof receipt.workId !== 'string' || !receipt.workId.trim())) { + throw new Error('RECEIPT_WORK_ID_INVALID'); + } + } return { schemaVersion: 1, observedAt: new Date(nowMs).toISOString(), source: 'router-state', active: receipts.filter((receipt) => ACTIVE_RECEIPT_STATES.has(receipt.state)).map((receipt) => ({ workId: receipt.workId, + attemptId: receipt.attemptId, + cwd: receipt.cwd, laneId: receipt.laneId, state: receipt.state, attemptedAt: receipt.attemptedAt, @@ -337,17 +333,47 @@ async function probeLane(config, lane, machineAdmission, options) { }; } +function stageDeadline(operation, stage, timeoutMs) { + if (!Number.isInteger(timeoutMs) || timeoutMs < 1 || timeoutMs > 60000) throw new Error('STAGE_TIMEOUT_INVALID'); + let timer; + return Promise.race([ + Promise.resolve().then(operation), + new Promise((_, reject) => { + timer = setTimeout(() => { + const error = new Error(`${stage}_TIMEOUT`); + error.code = `${stage}_TIMEOUT`; + reject(error); + }, timeoutMs); + }), + ]).finally(() => clearTimeout(timer)); +} + async function snapshot(config, options) { - const admissions = new Map(); - for (const machine of config.machines) { - let capacity = null; - let claims = null; - try { capacity = await options.capacityProvider(machine, options.nowMs); } catch { /* fail closed below */ } - try { claims = await options.claimsProvider(machine, config.statePath, options.nowMs); } catch { /* fail closed below */ } - admissions.set(machine.id, evaluateMachineAdmission(machine, capacity, claims, options.nowMs)); - } + const rows = await Promise.all(config.machines.map(async (machine) => { + const errors = []; + const read = async (operation, stage) => { + try { return await stageDeadline(operation, stage, options.stageTimeoutMs); } + catch (error) { + errors.push(error.code || `${stage}_PROBE_FAILED`); + return null; + } + }; + const [capacity, claims] = await Promise.all([ + read(() => options.capacityProvider(machine, options.nowMs), 'CAPACITY'), + read(() => options.claimsProvider(machine, config.statePath, options.nowMs), 'CLAIMS'), + ]); + const admission = evaluateMachineAdmission(machine, capacity, claims, options.nowMs); + admission.issues = [...new Set([...admission.issues, ...errors])]; + admission.eligible = admission.issues.length === 0; + return [machine.id, admission]; + })); + const admissions = new Map(rows); + const boundedProvider = (lane, operation) => stageDeadline( + () => options.conductorProvider(lane, operation), + operation.kind.toUpperCase().replaceAll('-', '_'), options.stageTimeoutMs, + ); const candidates = await Promise.all(config.conductors.map((lane) => probeLane( - config, lane, admissions.get(lane.machineId), options, + config, lane, admissions.get(lane.machineId), { ...options, conductorProvider: boundedProvider }, ))); return { schemaVersion: 1, @@ -445,6 +471,7 @@ export async function rank(configInput, options = {}) { return snapshot(config, { nowMs, usageMaxAgeMs: options.usageMaxAgeMs ?? 15_000, + stageTimeoutMs: options.stageTimeoutMs ?? config.routing.timeoutMs, capability: options.capability ?? 'tools', role: options.role ?? 'leaf', model: options.model, @@ -455,108 +482,128 @@ export async function rank(configInput, options = {}) { }); } +function assertClaimsAvailable(snapshot, receipt) { + for (const admission of Object.values(snapshot.admissions)) { + for (const claim of admission.claims?.active ?? []) { + if (claim.attemptId === receipt.attemptId) continue; + if (claim.workId === receipt.workId) throw new Error('WORK_ID_CLAIM_CONFLICT'); + if (typeof claim.cwd === 'string' && path.resolve(claim.cwd) === receipt.cwd) { + throw new Error('WORKSPACE_CLAIM_CONFLICT'); + } + } + } +} + +export async function inspectDispatch(configInput, options = {}) { + const config = validateInstallConfig(configInput, { expectedBorgHome: configInput.borgHome }); + if (typeof options.workId !== 'string' || !options.workId.trim() || !path.isAbsolute(options.cwd ?? '')) { + throw new Error('workId and absolute cwd are required'); + } + const cwd = path.resolve(options.cwd); + const intentPath = path.join(config.statePath, 'intents', `${intentDigest(options.workId.trim(), cwd)}.json`); + const intent = await readPrivateJsonRecord(intentPath); + if (intent === null) return { found: false, state: 'NOT_FOUND', noStartProven: false, completionVerified: false }; + const receiptPath = intent.receiptPath; + if (typeof receiptPath !== 'string' || path.dirname(receiptPath) !== path.join(config.statePath, 'dispatch-receipts') + || !path.basename(receiptPath).endsWith('.json')) throw new Error('UNSAFE_RECEIPT_REFERENCE'); + const receipt = await readPrivateJsonRecord(receiptPath); + if (!receipt || receipt.workId !== options.workId.trim() || receipt.cwd !== cwd) throw new Error('DISPATCH_RECEIPT_MISMATCH'); + return { found: true, receipt, receiptPath, + noStartProven: receipt.state === 'PRE_START_FAILED' && receipt.nativeStartAttempted === false, + completionVerified: false }; +} + export async function dispatch(configInput, options = {}) { const config = validateInstallConfig(configInput, { expectedBorgHome: configInput.borgHome }); if (typeof options.workId !== 'string' || !options.workId.trim()) throw new Error('workId is required for duplicate-safe dispatch'); if (typeof options.prompt !== 'string' || !options.prompt.trim()) throw new Error('prompt is required'); - const cwd = path.resolve(options.cwd ?? ''); + if (!path.isAbsolute(options.cwd ?? '')) throw new Error('absolute cwd is required'); + const cwd = path.resolve(options.cwd); const stat = await fs.stat(cwd); if (!stat.isDirectory()) throw new Error('cwd must be a directory'); + const stageTimeoutMs = options.stageTimeoutMs ?? config.routing.timeoutMs; + if (!Number.isInteger(stageTimeoutMs) || stageTimeoutMs < 1 || stageTimeoutMs > 60000) throw new Error('STAGE_TIMEOUT_INVALID'); const release = await acquireLock(config.statePath, config.routing.lockTimeoutMs); + let receipt; let paths; let nativeAttempted = false; + const clock = () => new Date(options.nowMs ?? Date.now()).toISOString(); + const progress = async (phase, changes = {}) => { + receipt = { ...receipt, ...changes, phase, + events: [...receipt.events, { phase, at: clock() }].slice(-16) }; + await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + }; try { const common = { - ...options, - nowMs: options.nowMs ?? Date.now(), - capability: options.capability ?? 'tools', - role: options.role ?? 'leaf', + ...options, stageTimeoutMs, + capability: options.capability ?? 'tools', role: options.role ?? 'leaf', capacityProvider: options.capacityProvider ?? observeMachineCapacity, claimsProvider: options.claimsProvider ?? observeClaims, conductorProvider: options.conductorProvider ?? ((lane, operation) => nativeConductorProvider(lane, operation, config.routing.timeoutMs)), }; + receipt = { + schemaVersion: 1, attemptId: crypto.randomUUID(), workId: options.workId.trim(), cwd, + laneId: null, machineId: null, state: 'ATTEMPTING', phase: 'PREPARING', + attemptedAt: clock(), dispatchedAt: null, threadId: null, turnId: null, + errorClass: null, nativeStartAttempted: false, stageTimeoutMs, + promptSha256: sha256(options.prompt), events: [{ phase: 'PREPARING', at: clock() }], + }; + // Persist the work intent BEFORE any admission scans or native lifecycle call. + paths = await createReceipt(config, receipt); + await progress('ADMISSION'); const preliminary = await rank(config, common); - const final = await rank(config, { ...common, nowMs: options.recheckNowMs ?? common.nowMs }); + await progress('RECHECK'); + const final = await rank(config, { ...common, nowMs: options.recheckNowMs ?? options.nowMs ?? Date.now() }); const selected = final.candidates.find((candidate) => candidate.eligible); if (!selected) throw new Error(noEligibleMessage(final)); + assertClaimsAvailable(final, receipt); const lane = config.conductors.find((candidate) => candidate.id === selected.laneId); - let receipt = { - schemaVersion: 1, - attemptId: crypto.randomUUID(), - workId: options.workId.trim(), - cwd, - laneId: lane.id, - machineId: lane.machineId, - state: 'ATTEMPTING', - phase: 'THREAD_START_PENDING', - attemptedAt: new Date(common.nowMs).toISOString(), - dispatchedAt: null, - remainingPercent: selected.remainingPercent, - resetAt: selected.resetAt, + await progress('THREAD_START_PENDING', { + laneId: lane.id, machineId: lane.machineId, nativeStartAttempted: true, + remainingPercent: selected.remainingPercent, resetAt: selected.resetAt, preliminaryLeader: preliminary.candidates.find((candidate) => candidate.eligible)?.laneId ?? null, finalLeader: selected.laneId, - threadId: null, - turnId: null, - errorClass: null, - }; - const paths = await createReceipt(config, receipt); + }); + nativeAttempted = true; let thread; try { - thread = await common.conductorProvider(lane, { - kind: 'thread-start', - role: common.role, - body: { - cwd, - role: common.role, - approvalPolicy: 'never', - sandbox: options.sandbox ?? 'danger-full-access', + thread = await stageDeadline(() => common.conductorProvider(lane, { + kind: 'thread-start', role: common.role, + body: { cwd, role: common.role, approvalPolicy: 'never', sandbox: options.sandbox ?? 'danger-full-access', ...(options.model ? { model: options.model } : {}), - ...(options.instructions ? { instructions: options.instructions } : {}), - }, - }); + ...(options.instructions ? { instructions: options.instructions } : {}) }, + }), 'THREAD_START', stageTimeoutMs); } catch (error) { - receipt = { ...receipt, state: 'UNKNOWN_DO_NOT_RETRY', phase: 'THREAD_START', errorClass: error.name || 'Error' }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + await progress('THREAD_START', { state: 'UNKNOWN_DO_NOT_RETRY', errorClass: error.code || error.name || 'Error' }); throw new Error(`thread start outcome unknown; do not retry: ${paths.receiptPath}`); } - if (!thread?.threadId) { - receipt = { ...receipt, state: 'UNKNOWN_DO_NOT_RETRY', phase: 'THREAD_START_RESPONSE', errorClass: 'MISSING_THREAD_ID' }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + if (!thread?.threadId || typeof thread.threadId !== 'string') { + await progress('THREAD_START_RESPONSE', { state: 'UNKNOWN_DO_NOT_RETRY', errorClass: 'MISSING_THREAD_ID' }); throw new Error(`thread start returned no ID; do not retry: ${paths.receiptPath}`); } - receipt = { ...receipt, state: 'THREAD_STARTED', phase: 'TURN_START_PENDING', threadId: thread.threadId }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + await progress('TURN_START_PENDING', { state: 'THREAD_STARTED', threadId: thread.threadId }); let turn; try { - turn = await common.conductorProvider(lane, { - kind: 'turn-start', - body: { - threadId: thread.threadId, - text: options.prompt, + turn = await stageDeadline(() => common.conductorProvider(lane, { + kind: 'turn-start', body: { threadId: thread.threadId, text: options.prompt, ...(options.model ? { model: options.model } : {}), - ...(options.effort ? { effort: options.effort } : {}), - }, - }); + ...(options.effort ? { effort: options.effort } : {}) }, + }), 'TURN_START', stageTimeoutMs); } catch (error) { - receipt = { ...receipt, state: 'STARTED_TURN_UNKNOWN', phase: 'TURN_START', errorClass: error.name || 'Error' }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + await progress('TURN_START', { state: 'STARTED_TURN_UNKNOWN', errorClass: error.code || error.name || 'Error' }); throw new Error(`thread exists but turn outcome unknown; do not retry: ${paths.receiptPath}`); } - if (!turn?.turnId) { - receipt = { ...receipt, state: 'STARTED_TURN_UNKNOWN', phase: 'TURN_START_RESPONSE', errorClass: 'MISSING_TURN_ID' }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + if (!turn?.turnId || typeof turn.turnId !== 'string') { + await progress('TURN_START_RESPONSE', { state: 'STARTED_TURN_UNKNOWN', errorClass: 'MISSING_TURN_ID' }); throw new Error(`turn start returned no ID; do not retry: ${paths.receiptPath}`); } - receipt = { - ...receipt, - state: 'DISPATCHED', - phase: 'TURN_STARTED', - dispatchedAt: new Date(options.recheckNowMs ?? common.nowMs).toISOString(), - threadId: thread.threadId, - turnId: turn.turnId, - }; - await updateReceipt(paths.receiptPath, paths.intentPath, receipt); + await progress('TURN_STARTED', { state: 'DISPATCHED', dispatchedAt: clock(), turnId: turn.turnId }); return { receipt, receiptPath: paths.receiptPath, snapshot: final }; - } finally { - await release(); - } + } catch (error) { + if (paths && !nativeAttempted) { + await progress(receipt.phase, { state: 'PRE_START_FAILED', nativeStartAttempted: false, errorClass: error.code || error.name || 'Error' }); + error.receiptPath = paths.receiptPath; + error.message += `; receipt: ${paths.receiptPath}`; + } + throw error; + } finally { await release(); } } diff --git a/conductor/tests/receipt-reader.test.mjs b/conductor/tests/receipt-reader.test.mjs new file mode 100644 index 0000000..65e8255 --- /dev/null +++ b/conductor/tests/receipt-reader.test.mjs @@ -0,0 +1,56 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import test from 'node:test'; +import { readPrivateJsonDirectory, readPrivateJsonRecord } from '../router/receipt-reader.mjs'; + +function fixture(t) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), 'borg-receipt-reader-')); + fs.chmodSync(root, 0o700); + t.after(() => fs.rmSync(root, { recursive: true, force: true })); + return root; +} +test('real child returns complete private JSON records and explicit missing values', async (t) => { + const root = fixture(t); + fs.writeFileSync(path.join(root,'one.json'), JSON.stringify({ state: 'DISPATCHED', workId: 'one' }), { mode: 0o600 }); + fs.writeFileSync(path.join(root,'ignored.txt'), 'unrelated', { mode: 0o600 }); + assert.deepEqual(await readPrivateJsonDirectory(root), [{ state: 'DISPATCHED', workId: 'one' }]); + assert.deepEqual(await readPrivateJsonRecord(path.join(root,'one.json')), { state: 'DISPATCHED', workId: 'one' }); + assert.deepEqual(await readPrivateJsonDirectory(path.join(root,'missing')), []); + assert.equal(await readPrivateJsonRecord(path.join(root,'missing.json')), null); +}); +test('one malformed receipt makes the entire scan unknown, not partial success', async (t) => { + const root = fixture(t); + fs.writeFileSync(path.join(root,'a.json'), '{"workId":"a"}', { mode: 0o600 }); + fs.writeFileSync(path.join(root,'z.json'), '{broken', { mode: 0o600 }); + await assert.rejects(readPrivateJsonDirectory(root), { code: 'RECEIPT_JSON_INVALID' }); +}); +test('oversized receipts are rejected before unbounded buffering', async (t) => { + const root = fixture(t); + fs.writeFileSync(path.join(root,'large.json'), JSON.stringify({ text: 'x'.repeat(65537) }), { mode: 0o600 }); + await assert.rejects(readPrivateJsonDirectory(root), { code: 'RECEIPT_FILE_LIMIT' }); +}); +test('non-private files and directory symlinks are refused', async (t) => { + const root = fixture(t); + const target = path.join(root,'open.json'); + fs.writeFileSync(target, '{}'); fs.chmodSync(target, 0o644); + await assert.rejects(readPrivateJsonRecord(target), { code: 'UNSAFE_RECEIPT' }); + const alias = path.join(root,'alias'); fs.symlinkSync(root, alias); + await assert.rejects(readPrivateJsonDirectory(alias), { code: 'UNSAFE_RECEIPT' }); +}); +test('native child timeout returns a bounded failure rather than claiming empty receipts', async (t) => { + const root = fixture(t); + // One millisecond is deliberately shorter than a fresh Node process can load + // this worker. This tests actual child deadline handling, not mocked claims. + const start = performance.now(); + await assert.rejects(readPrivateJsonDirectory(root, { timeoutMs: 1 }), { code: 'RECEIPT_READ_TIMEOUT', stage: 'RECEIPT_READ' }); + assert.ok(performance.now() - start < 5000); +}); +test('invalid timeout configuration is rejected without starting a reader', async (t) => { + const root = fixture(t); + for (const timeoutMs of [null, 0, -1, NaN, 60001, '1']) { + if (timeoutMs === null) continue; // null deliberately selects the default. + await assert.rejects(readPrivateJsonDirectory(root, { timeoutMs }), /RECEIPT_READ_TIMEOUT_INVALID/); + } +}); diff --git a/conductor/tests/router-reliability.test.mjs b/conductor/tests/router-reliability.test.mjs new file mode 100644 index 0000000..ddff9dc --- /dev/null +++ b/conductor/tests/router-reliability.test.mjs @@ -0,0 +1,205 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import test from 'node:test'; +import * as router from '../router/router.mjs'; +import { buildDefaultConfig } from '../config.mjs'; + +const NOW = Date.parse('2026-09-16T12:00:00Z'); +const account = { type: 'chatgpt', email: 'owner@example.test', planType: 'pro' }; +const observation = () => ({ observedAt: new Date(NOW).toISOString(), reachable: true, loadPerCore: 0.2, memoryUsePercent: 35 }); +const claims = (active = []) => ({ observedAt: new Date(NOW).toISOString(), active }); +function fixture(t) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), 'borg-reliability-')); + fs.chmodSync(root, 0o700); + t.after(() => fs.rmSync(root, { recursive: true, force: true })); + const config = buildDefaultConfig(root, process.execPath, { nodeBin: process.execPath }); + config.statePath = path.join(root, 'private/router'); + config.conductors[0].accountPin = router.accountIdentityDigest(account); + const calls = []; + const native = async (lane, op) => { + calls.push(op.kind); + if (op.kind === 'status') return { ok: true, port: lane.port, codexHome: lane.codexHome, supportedRoles: ['leaf','lead'], threads: {} }; + if (op.method === 'account/read') return { account }; + if (op.method === 'account/rateLimits/read') return { rateLimits: { limitId: 'codex', primary: { usedPercent: 5, windowDurationMins: 300, resetsAt: '2026-09-17T12:00:00Z' } } }; + if (op.kind === 'thread-start') return { threadId: 'thread-test' }; + if (op.kind === 'turn-start') return { turnId: 'turn-test' }; + throw new Error('unexpected operation'); + }; + const options = { cwd: root, workId: 'repair', prompt: 'Synthetic acceptance only', nowMs: NOW, stageTimeoutMs: 2000, + capacityProvider: async () => observation(), claimsProvider: async () => claims(), conductorProvider: native }; + return { root, config, options, native, calls }; +} +function putReceipt(config, name, value) { + const dir = path.join(config.statePath, 'dispatch-receipts'); + fs.mkdirSync(dir, { recursive: true, mode: 0o700 }); + fs.writeFileSync(path.join(dir, `${name}.json`), JSON.stringify(value), { mode: 0o600 }); +} +async function bounded(promise, ms = 1500) { + let timer; + try { return await Promise.race([promise, new Promise((_, reject) => { timer = setTimeout(() => reject(new Error('TEST_WATCHDOG_EXPIRED')), ms); })]); } + finally { clearTimeout(timer); } +} + +for (const value of [null, false, true, '', ' ', [], {}, '0']) { + test(`unknown telemetry ${JSON.stringify(value)} is not safe zero`, () => { + const machine = { id: 'local', capacity: { maxAgeMs: 15000, loadPerCoreLimit: 1.5, memoryUseLimitPercent: 92 }, claims: { maxAgeMs: 15000 } }; + for (const field of ['loadPerCore', 'memoryUsePercent']) { + const result = router.evaluateMachineAdmission(machine, { ...observation(), [field]: value }, claims(), NOW); + assert.equal(result.eligible, false, `${field}=${JSON.stringify(value)} must refuse`); + assert.ok(result.issues.includes(field === 'loadPerCore' ? 'LOAD_UNKNOWN' : 'MEMORY_UNKNOWN')); + } + }); +} +test('missing native usage never creates a worker', async (t) => { + const f = fixture(t); + await assert.rejects(router.dispatch(f.config, { ...f.options, conductorProvider: async (lane, op) => { + if (op.method === 'account/rateLimits/read') return { rateLimits: { limitId: 'codex', primary: { usedPercent: null, windowDurationMins: 300, resetsAt: '2026-09-17T12:00:00Z' } } }; + return f.native(lane, op); + } }), /PRIMARY_USED_INVALID/); + assert.ok(!f.calls.includes('thread-start')); +}); +test('DISPATCHED and uncertain receipts remain visible as active claims', async (t) => { + const f = fixture(t); + for (const state of ['DISPATCHED','ATTEMPTING','THREAD_STARTED','UNKNOWN_DO_NOT_RETRY','STARTED_TURN_UNKNOWN']) { + putReceipt(f.config, state, { state, workId: state, cwd: f.root, attemptId: state }); + } + const result = await router.observeClaims(f.config.machines[0], f.config.statePath, NOW); + assert.equal(result.active.length, 5); + assert.equal(result.active.find(row => row.state === 'DISPATCHED').cwd, f.root); +}); +test('a differently named work item cannot bypass an existing workspace claim', async (t) => { + const f = fixture(t); + await assert.rejects(router.dispatch(f.config, { ...f.options, + claimsProvider: async () => claims([{ workId: 'other-work', cwd: f.root, state: 'DISPATCHED', attemptId: 'other-attempt' }]) + }), /WORKSPACE_CLAIM_CONFLICT/); + assert.ok(!f.calls.includes('thread-start')); +}); +test('work identity stays reserved across different workspace paths', async (t) => { + const f = fixture(t); + await assert.rejects(router.dispatch(f.config, { ...f.options, + claimsProvider: async () => claims([{ workId: 'repair', cwd: path.join(f.root, 'another'), state: 'UNKNOWN_DO_NOT_RETRY', attemptId: 'other-attempt' }]) + }), /WORK_ID_CLAIM_CONFLICT/); + assert.ok(!f.calls.includes('thread-start')); +}); +test('unsafe JSON symlink does not silently turn into empty claims', async (t) => { + const f = fixture(t); putReceipt(f.config, 'real', { state: 'DISPATCHED', workId: 'owned' }); + fs.symlinkSync(path.join(f.config.statePath,'dispatch-receipts','real.json'), path.join(f.config.statePath,'dispatch-receipts','alias.json')); + await assert.rejects(router.observeClaims(f.config.machines[0], f.config.statePath, NOW), /UNSAFE_RECEIPT/); +}); +test('unknown receipt state fails closed rather than disappearing', async (t) => { + const f = fixture(t); putReceipt(f.config, 'future', { state: 'NEW_UNKNOWN_STATE', workId: 'owned' }); + await assert.rejects(router.observeClaims(f.config.machines[0], f.config.statePath, NOW), /RECEIPT_STATE_UNKNOWN/); +}); +test('admission failure leaves a status receipt with no native start', async (t) => { + const f = fixture(t); + await assert.rejects(router.dispatch(f.config, { ...f.options, capacityProvider: async () => null }), /CAPACITY_UNKNOWN/); + assert.equal(typeof router.inspectDispatch, 'function', 'restart-safe status lookup must exist'); + const result = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'repair' }); + assert.equal(result.receipt.state, 'PRE_START_FAILED'); + assert.equal(result.receipt.threadId, null); + assert.equal(result.noStartProven, true); + assert.equal(result.completionVerified, false); + assert.ok(result.receipt.events.some(event => event.phase === 'ADMISSION')); +}); +test('stalled admission is bounded and records its exact stage', async (t) => { + const f = fixture(t); + await assert.rejects(bounded(router.dispatch(f.config, { ...f.options, stageTimeoutMs: 40, + capacityProvider: async () => new Promise(() => {}) + })), /CAPACITY_TIMEOUT/); + const result = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'repair' }); + assert.equal(result.receipt.state, 'PRE_START_FAILED'); + assert.ok(!f.calls.includes('thread-start')); +}); +test('thread timeout blocks replay and a late response cannot start a turn', async (t) => { + const f = fixture(t); let finishStart; + await assert.rejects(bounded(router.dispatch(f.config, { ...f.options, stageTimeoutMs: 40, + conductorProvider: async (lane, op) => { + if (op.kind === 'thread-start') { f.calls.push(op.kind); return new Promise(resolve => { finishStart = resolve; }); } + return f.native(lane, op); + } + })), /outcome unknown.*do not retry/i); + const result = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'repair' }); + assert.equal(result.receipt.state, 'UNKNOWN_DO_NOT_RETRY'); + assert.equal(result.receipt.errorClass, 'THREAD_START_TIMEOUT'); + assert.equal(result.noStartProven, false); + finishStart({ threadId: 'late-thread' }); + await new Promise(resolve => setImmediate(resolve)); + assert.ok(!f.calls.includes('turn-start')); + await assert.rejects(router.dispatch(f.config, f.options), /duplicate intent.*do not retry/i); + assert.equal(f.calls.filter(kind => kind === 'thread-start').length, 1); +}); +test('successful dispatch exposes IDs but never claims completed work', async (t) => { + const f = fixture(t); const sent = await router.dispatch(f.config, f.options); + const result = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'repair' }); + assert.equal(result.receipt.attemptId, sent.receipt.attemptId); + assert.equal(result.receipt.state, 'DISPATCHED'); + assert.equal(result.receipt.threadId, 'thread-test'); + assert.equal(result.receipt.turnId, 'turn-test'); + assert.equal(result.completionVerified, false); + assert.equal(result.noStartProven, false); + assert.ok(result.receipt.events.some(event => event.phase === 'RECHECK')); + assert.ok(!JSON.stringify(result).includes(f.options.prompt)); +}); +test('missing status receipt is not evidence of a safe retry', async (t) => { + const f = fixture(t); + assert.equal(typeof router.inspectDispatch, 'function'); + const result = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'never-recorded' }); + assert.equal(result.found, false); + assert.equal(result.noStartProven, false); +}); + +import { spawnSync } from 'node:child_process'; +import { fileURLToPath } from 'node:url'; +test('CLI route-status reads persisted IDs without dispatching or probing providers', async (t) => { + const f = fixture(t); await router.dispatch(f.config, f.options); + const configPath = path.join(f.root,'conductors','config.json'); + fs.mkdirSync(path.dirname(configPath), { recursive: true, mode: 0o700 }); + fs.writeFileSync(configPath, JSON.stringify(f.config), { mode: 0o600 }); + const result = spawnSync(process.execPath, [fileURLToPath(new URL('../borg-conductor.mjs', import.meta.url)), + 'route-status', '--cwd', f.root, '--work-id', 'repair'], + { encoding: 'utf8', timeout: 20000, env: { ...process.env, BORG_HOME: f.root } }); + assert.equal(result.status, 0, result.stderr); + const status = JSON.parse(result.stdout); + assert.equal(status.receipt.threadId, 'thread-test'); + assert.equal(status.receipt.turnId, 'turn-test'); + assert.equal(status.completionVerified, false); + assert.equal(f.calls.filter(kind => kind === 'thread-start').length, 1); +}); + +test('a timed-out candidate does not block an eligible fallback before native start', async (t) => { + const f = fixture(t); + f.config.conductors.push({ ...f.config.conductors[0], id: 'backup', accountProfile: 'backup', port: 4748, + codexHome: path.join(f.root,'conductors/backup/profile'), logsPath: path.join(f.root,'conductors/backup/logs') }); + const result = await router.dispatch(f.config, { ...f.options, stageTimeoutMs: 40, + conductorProvider: async (lane, op) => { + if (lane.id === 'primary' && op.kind === 'status') return new Promise(() => {}); + return f.native(lane, op); + } }); + assert.equal(result.receipt.laneId, 'backup'); + assert.equal(f.calls.filter(kind => kind === 'thread-start').length, 1); + assert.equal(f.calls.filter(kind => kind === 'turn-start').length, 1); +}); +test('turn timeout keeps native thread identity and blocks a replacement launch', async (t) => { + const f = fixture(t); + await assert.rejects(router.dispatch(f.config, { ...f.options, stageTimeoutMs: 40, + conductorProvider: async (lane, op) => { + if (op.kind === 'turn-start') { f.calls.push(op.kind); return new Promise(() => {}); } + return f.native(lane, op); + } }), /turn outcome unknown.*do not retry/i); + const status = await router.inspectDispatch(f.config, { cwd: f.root, workId: 'repair' }); + assert.equal(status.receipt.state, 'STARTED_TURN_UNKNOWN'); + assert.equal(status.receipt.threadId, 'thread-test'); + assert.equal(status.receipt.errorClass, 'TURN_START_TIMEOUT'); + assert.equal(status.noStartProven, false); + await assert.rejects(router.dispatch(f.config, f.options), /duplicate intent.*do not retry/i); + assert.equal(f.calls.filter(kind => kind === 'thread-start').length, 1); +}); +test('concurrent same-intent calls start exactly one native thread', async (t) => { + const f = fixture(t); + const results = await Promise.allSettled([router.dispatch(f.config, f.options), router.dispatch(f.config, f.options)]); + assert.equal(results.filter(result => result.status === 'fulfilled').length, 1); + assert.match(results.find(result => result.status === 'rejected').reason.message, /duplicate intent.*do not retry/i); + assert.equal(f.calls.filter(kind => kind === 'thread-start').length, 1); +}); From 8ef8fe3a414c4e7826929b8fff4493d1f0d903b8 Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 04:10:09 -0600 Subject: [PATCH 2/6] Prevent local receipt and workspace alias collisions --- adapters/README.md | 10 +- conductor/README.md | 6 +- conductor/docs/LAUNCH-RELIABILITY.md | 40 +++++- conductor/router/router.mjs | 127 ++++++++++++++--- conductor/tests/router-ownership.test.mjs | 158 ++++++++++++++++++++++ 5 files changed, 310 insertions(+), 31 deletions(-) create mode 100644 conductor/tests/router-ownership.test.mjs diff --git a/adapters/README.md b/adapters/README.md index 894d4f6..5f366fc 100644 --- a/adapters/README.md +++ b/adapters/README.md @@ -19,9 +19,13 @@ capture. Those shapes are what the scrub inserted: synthetic stand-ins look like identifiers because they were built to. No screened registry term appeared in any generation. -**Compatibility is exact-match by construction.** A LoRA adapter is a low-rank delta on one -specific frozen base model. Each adapter below loads only with the exact Hugging Face model -id listed — not other sizes, not other revisions, not other quantizations. Load with: +**Compatibility depends on the base weights.** A LoRA adapter is a low-rank delta on one +specific frozen base model. Use the model ID, quantization and immutable release pin in +[`MANIFEST.json`](MANIFEST.json); do not substitute another size or conversion. +The training records preserve the model IDs but not the training-time commit hashes. +The manifest pins verified release downloads, not reconstructed training revisions. +Matching those pins does not establish compatibility or promotion readiness: validate +each adapter with its base model and run the required canary before activation. Load with: ```bash pip install mlx-lm diff --git a/conductor/README.md b/conductor/README.md index 8e3f4cc..66d5bee 100644 --- a/conductor/README.md +++ b/conductor/README.md @@ -59,8 +59,10 @@ codes. No discretionary reservation or allowance floor is created. Dispatch holds a private lock, persists an intent before admission scans, rechecks the configured fleet before selection, and refuses duplicate -`{workId,cwd}` intents across process restarts. Active workspace and work-ID -claims are checked before native lifecycle calls. An ambiguous thread or turn +`{workId,cwd}` intents across process restarts. Workspaces are compared by +filesystem-canonical identity, so aliases and parent/child paths overlap. Active +workspace and work-ID claims, including this router's own receipts under any +claims collector, are checked before native lifecycle calls. An ambiguous thread or turn response remains `DO_NOT_RETRY` evidence. Use `route-status --cwd ABS --work-id ID` to inspect saved phase history and diff --git a/conductor/docs/LAUNCH-RELIABILITY.md b/conductor/docs/LAUNCH-RELIABILITY.md index 8bf7772..5e69b93 100644 --- a/conductor/docs/LAUNCH-RELIABILITY.md +++ b/conductor/docs/LAUNCH-RELIABILITY.md @@ -8,7 +8,8 @@ borg-conductor route-status --config /absolute/BORG_HOME/conductors/config.json --cwd /absolute/workspace --work-id stable-work-id ``` -Use the same canonical workspace and stable work ID supplied to `route`. This is +Use the workspace (any spelling of it resolves to the same filesystem-canonical +path) and stable work ID supplied to `route`. This is read-only reconciliation: it does not create, resume, cancel or complete a worker. The returned receipt contains a bounded phase history, attempt ID, selected lane, native thread/turn IDs when proved, and error classification. Prompt text is not @@ -40,8 +41,33 @@ transport dispatches and uncertain lifecycle states remain active claims until explicitly reconciled. A different work ID cannot evade an existing workspace claim; a different workspace cannot evade an active work ID. +Under the dispatch lock the router always rescans its own receipt ledger, +whichever claims collector each machine uses, and checks it together with every +machine's observed claims; only the current attempt is exempt. A command +collector therefore cannot hide this router's submitted or uncertain launches, +and a failed ledger scan refuses dispatch instead of trusting an empty external +list. + +Workspace identity is filesystem-canonical. Native realpath resolves symlinks +and, on case-insensitive volumes, letter case; device/inode ancestry also matches +spellings realpath keeps distinct, such as macOS firmlinks and bind mounts. +Receipts, intents, `route-status` and the native thread use the canonical path. +Two workspaces conflict when one is the same directory as, or an ancestor of, +the other. The comparison is segment-aware, so sibling worktrees such as `repo` +and `repo-other` remain independently admissible. + +Claimed paths from every machine are resolved on the router host's filesystem, +where `route` validates the workspace; paths reported for different machines are +never assumed disjoint. A claimed workspace that no longer exists still blocks +its ancestors. A claim without `cwd` reserves only its work ID; a malformed or +unresolvable claimed `cwd` refuses dispatch (`WORKSPACE_CLAIM_UNRESOLVED`). +Intents recorded under a lexical path by earlier versions still block replay and +remain visible to `route-status`. Symlinks inside a workspace that point +elsewhere are not traced. + A receipt scan is all-or-error. Malformed JSON, unsafe file permissions, symlinks, -unknown receipt states, oversize files or exceeded bounds never become an empty +unknown receipt states, active receipts without a work ID or absolute workspace, +oversize files or exceeded bounds never become an empty or partial successful claims list. The fixed read-only reader runs in a separate Node process, outside the router's filesystem worker pool. It accepts no shell commands and never starts agents or writes state. @@ -60,8 +86,10 @@ run concurrently. Late thread responses cannot advance to turn creation after the caller records an unknown outcome. These are **stage** deadlines, not an end-to-end service-level guarantee. Workspace -validation, dispatch-lock I/O and durable state writes can still be delayed by an -unhealthy filesystem. State writes are not raced against a timeout, since a late +validation and identity resolution (including claimed paths), dispatch-lock I/O +and durable state writes can still be delayed by an unhealthy filesystem. The +final ledger rescan is bounded by the reader limits above, not by +`--stage-timeout-ms`. State writes are not raced against a timeout, since a late write could corrupt the evidence used by a subsequent attempt. Existing lock contention controls remain. The router never removes an ownership check or claims it cancelled an external lifecycle request just because its wait ended. @@ -73,7 +101,9 @@ node --test conductor/tests/*.test.mjs conductor/providers/*.test.mjs ``` The regression suite covers missing telemetry/usage, active dispatched claims, -workspace/work-ID conflicts, unsafe/corrupt/oversize receipts, native reader +workspace/work-ID conflicts, local-ledger enforcement under a command collector, +workspace alias/case/firmlink/parent/child overlap with a sibling-worktree +control, unsafe/corrupt/oversize receipts, native reader timeout, pre-start diagnostic persistence, lifecycle timeout with late response, duplicate suppression, status lookup and the real command-line status path. Optional native Codex initialization uses `BORG_TEST_CODEX_BIN` with an isolated diff --git a/conductor/router/router.mjs b/conductor/router/router.mjs index 3b6a53d..2669122 100644 --- a/conductor/router/router.mjs +++ b/conductor/router/router.mjs @@ -137,6 +137,10 @@ async function readJsonFiles(directory) { export async function observeClaims(machine, statePath, nowMs = Date.now()) { if (machine.claims.kind === 'command') return commandJson(machine.claims); if (machine.claims.kind !== 'router-state') throw new Error('CLAIMS_KIND_UNSUPPORTED'); + return observeLocalReceiptClaims(statePath, nowMs); +} + +async function observeLocalReceiptClaims(statePath, nowMs = Date.now()) { const receipts = await readJsonFiles(path.join(statePath, 'dispatch-receipts')); const terminal = new Set(['PRE_START_FAILED', 'COMPLETED', 'FAILED', 'CANCELLED', 'RECONCILED_NO_START']); for (const receipt of receipts) { @@ -146,6 +150,9 @@ export async function observeClaims(machine, statePath, nowMs = Date.now()) { if (ACTIVE_RECEIPT_STATES.has(receipt.state) && (typeof receipt.workId !== 'string' || !receipt.workId.trim())) { throw new Error('RECEIPT_WORK_ID_INVALID'); } + if (ACTIVE_RECEIPT_STATES.has(receipt.state) && (typeof receipt.cwd !== 'string' || !path.isAbsolute(receipt.cwd))) { + throw new Error('RECEIPT_CWD_INVALID'); + } } return { schemaVersion: 1, @@ -431,18 +438,68 @@ function intentDigest(workId, cwd) { return sha256(JSON.stringify({ workId, cwd })); } -async function createReceipt(config, receipt) { +// Canonical workspace path: native realpath resolves symlinks and, on +// case-insensitive volumes, letter case. A missing tail (a claimed workspace +// that was since deleted) is kept below its deepest existing ancestor. +async function canonicalWorkspace(value) { + let existing = path.resolve(value); + const missing = []; + while (true) { + try { + const real = await fs.realpath(existing); + return { path: path.join(real, ...missing), existing: real, exists: missing.length === 0 }; + } catch (error) { + const parent = path.dirname(existing); + if (!['ENOENT', 'ENOTDIR'].includes(error.code) || parent === existing) throw error; + missing.unshift(path.basename(existing)); + existing = parent; + } + } +} + +// Device/inode of the deepest existing directory and each ancestor also +// matches spellings realpath keeps distinct, such as firmlinks and bind mounts. +async function workspaceIdentity(value) { + const canonical = await canonicalWorkspace(value); + const ancestry = []; + for (let current = canonical.existing; ; current = path.dirname(current)) { + const stat = await fs.stat(current, { bigint: true }); + ancestry.push(`${stat.dev}:${stat.ino}`); + if (path.dirname(current) === current) break; + } + return { ...canonical, ancestry }; +} + +function isWithin(parent, child) { + const relative = path.relative(parent, child); + return relative === '' || (relative !== '..' && !relative.startsWith(`..${path.sep}`) + && !path.isAbsolute(relative)); +} + +// Same directory, ancestor or descendant. Segment-aware, so sibling +// worktrees such as `repo` and `repo-other` stay independent. +function workspacesOverlap(left, right) { + return isWithin(left.path, right.path) || isWithin(right.path, left.path) + || (left.exists && right.ancestry.includes(left.ancestry[0])) + || (right.exists && left.ancestry.includes(right.ancestry[0])); +} + +async function createReceipt(config, receipt, lexicalCwd) { const directory = path.join(config.statePath, 'dispatch-receipts'); const intents = path.join(config.statePath, 'intents'); await ensurePrivateDirectory(directory); await ensurePrivateDirectory(intents); const receiptPath = path.join(directory, `${receipt.attemptedAt.replace(/[:.]/g, '-')}-${receipt.attemptId}.json`); const intentPath = path.join(intents, `${intentDigest(receipt.workId, receipt.cwd)}.json`); - try { - const existing = JSON.parse(await fs.readFile(intentPath, 'utf8')); - throw new Error(`duplicate intent; do not retry: ${existing.receiptPath}`); - } catch (error) { - if (error.code !== 'ENOENT') throw error; + // Intents recorded before canonical workspace identity are keyed by the + // lexical path and still block replay. + for (const cwd of new Set([receipt.cwd, lexicalCwd])) { + try { + const existing = JSON.parse(await fs.readFile(path.join(intents, `${intentDigest(receipt.workId, cwd)}.json`), 'utf8')); + throw new Error(`duplicate intent; do not retry: ${existing.receiptPath}`); + } catch (error) { + if (error.code !== 'ENOENT') throw error; + } } await atomicPrivateWrite(receiptPath, receipt, { createOnly: true }); try { @@ -482,15 +539,27 @@ export async function rank(configInput, options = {}) { }); } -function assertClaimsAvailable(snapshot, receipt) { - for (const admission of Object.values(snapshot.admissions)) { - for (const claim of admission.claims?.active ?? []) { - if (claim.attemptId === receipt.attemptId) continue; - if (claim.workId === receipt.workId) throw new Error('WORK_ID_CLAIM_CONFLICT'); - if (typeof claim.cwd === 'string' && path.resolve(claim.cwd) === receipt.cwd) { - throw new Error('WORKSPACE_CLAIM_CONFLICT'); +// Work IDs are reserved across every claim source and machine. Workspace +// overlap is checked against every claim naming a workspace, resolved on this +// router's filesystem, where dispatch validates cwd; paths reported for +// different machines are never assumed disjoint. A claim without cwd reserves +// only its work ID; a malformed or unresolvable cwd refuses dispatch. +async function assertClaimsAvailable(claimSets, receipt, workspace) { + const active = claimSets.flatMap((claims) => claims?.active ?? []) + .filter((claim) => claim.attemptId !== receipt.attemptId); + if (active.some((claim) => claim.workId === receipt.workId)) throw new Error('WORK_ID_CLAIM_CONFLICT'); + const identities = new Map(); + for (const claim of active) { + if (claim.cwd === undefined || claim.cwd === null) continue; + if (typeof claim.cwd !== 'string' || !path.isAbsolute(claim.cwd)) throw new Error('WORKSPACE_CLAIM_UNRESOLVED'); + if (!identities.has(claim.cwd)) { + try { + identities.set(claim.cwd, await workspaceIdentity(claim.cwd)); + } catch { + throw new Error('WORKSPACE_CLAIM_UNRESOLVED'); } } + if (workspacesOverlap(workspace, identities.get(claim.cwd))) throw new Error('WORKSPACE_CLAIM_CONFLICT'); } } @@ -499,15 +568,22 @@ export async function inspectDispatch(configInput, options = {}) { if (typeof options.workId !== 'string' || !options.workId.trim() || !path.isAbsolute(options.cwd ?? '')) { throw new Error('workId and absolute cwd are required'); } - const cwd = path.resolve(options.cwd); - const intentPath = path.join(config.statePath, 'intents', `${intentDigest(options.workId.trim(), cwd)}.json`); - const intent = await readPrivateJsonRecord(intentPath); + const workId = options.workId.trim(); + const lexicalCwd = path.resolve(options.cwd); + const { path: cwd } = await canonicalWorkspace(lexicalCwd); + let intent = null; + for (const candidate of new Set([cwd, lexicalCwd])) { + intent = await readPrivateJsonRecord(path.join(config.statePath, 'intents', `${intentDigest(workId, candidate)}.json`)); + if (intent !== null) break; + } if (intent === null) return { found: false, state: 'NOT_FOUND', noStartProven: false, completionVerified: false }; const receiptPath = intent.receiptPath; if (typeof receiptPath !== 'string' || path.dirname(receiptPath) !== path.join(config.statePath, 'dispatch-receipts') || !path.basename(receiptPath).endsWith('.json')) throw new Error('UNSAFE_RECEIPT_REFERENCE'); const receipt = await readPrivateJsonRecord(receiptPath); - if (!receipt || receipt.workId !== options.workId.trim() || receipt.cwd !== cwd) throw new Error('DISPATCH_RECEIPT_MISMATCH'); + if (!receipt || receipt.workId !== workId || (receipt.cwd !== cwd && receipt.cwd !== lexicalCwd)) { + throw new Error('DISPATCH_RECEIPT_MISMATCH'); + } return { found: true, receipt, receiptPath, noStartProven: receipt.state === 'PRE_START_FAILED' && receipt.nativeStartAttempted === false, completionVerified: false }; @@ -518,9 +594,13 @@ export async function dispatch(configInput, options = {}) { if (typeof options.workId !== 'string' || !options.workId.trim()) throw new Error('workId is required for duplicate-safe dispatch'); if (typeof options.prompt !== 'string' || !options.prompt.trim()) throw new Error('prompt is required'); if (!path.isAbsolute(options.cwd ?? '')) throw new Error('absolute cwd is required'); - const cwd = path.resolve(options.cwd); - const stat = await fs.stat(cwd); + const lexicalCwd = path.resolve(options.cwd); + const stat = await fs.stat(lexicalCwd); if (!stat.isDirectory()) throw new Error('cwd must be a directory'); + // Receipt, intent, claim checks and the native thread share one filesystem + // identity, so an alias or nested path is not a new workspace. + const workspace = await workspaceIdentity(lexicalCwd); + const { path: cwd } = workspace; const stageTimeoutMs = options.stageTimeoutMs ?? config.routing.timeoutMs; if (!Number.isInteger(stageTimeoutMs) || stageTimeoutMs < 1 || stageTimeoutMs > 60000) throw new Error('STAGE_TIMEOUT_INVALID'); const release = await acquireLock(config.statePath, config.routing.lockTimeoutMs); @@ -548,14 +628,19 @@ export async function dispatch(configInput, options = {}) { promptSha256: sha256(options.prompt), events: [{ phase: 'PREPARING', at: clock() }], }; // Persist the work intent BEFORE any admission scans or native lifecycle call. - paths = await createReceipt(config, receipt); + paths = await createReceipt(config, receipt, lexicalCwd); await progress('ADMISSION'); const preliminary = await rank(config, common); await progress('RECHECK'); const final = await rank(config, { ...common, nowMs: options.recheckNowMs ?? options.nowMs ?? Date.now() }); const selected = final.candidates.find((candidate) => candidate.eligible); if (!selected) throw new Error(noEligibleMessage(final)); - assertClaimsAvailable(final, receipt); + // This router's own submitted and uncertain receipts are claims whichever + // collector each machine uses; a failed ledger scan refuses dispatch. + const localClaims = await observeLocalReceiptClaims(config.statePath); + await assertClaimsAvailable([ + ...Object.values(final.admissions).map((admission) => admission.claims), localClaims, + ], receipt, workspace); const lane = config.conductors.find((candidate) => candidate.id === selected.laneId); await progress('THREAD_START_PENDING', { laneId: lane.id, machineId: lane.machineId, nativeStartAttempted: true, diff --git a/conductor/tests/router-ownership.test.mjs b/conductor/tests/router-ownership.test.mjs new file mode 100644 index 0000000..a764abe --- /dev/null +++ b/conductor/tests/router-ownership.test.mjs @@ -0,0 +1,158 @@ +import assert from 'node:assert/strict'; +import crypto from 'node:crypto'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import test from 'node:test'; +import * as router from '../router/router.mjs'; +import { buildDefaultConfig } from '../config.mjs'; + +// Workspace and work-ID ownership regressions for ROUTER-REVIEW R1/R2. The +// configured claims collector runs unless a test injects external claims. +const NOW = Date.parse('2026-09-16T12:00:00Z'); +const account = { type: 'chatgpt', email: 'owner@example.test', planType: 'pro' }; +const claims = (active = []) => ({ observedAt: new Date(NOW).toISOString(), active }); +function fixture(t, { collector = 'router-state' } = {}) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), 'borg-ownership-')); + fs.chmodSync(root, 0o700); + t.after(() => fs.rmSync(root, { recursive: true, force: true })); + const workspace = path.join(root, 'workspace'); + fs.mkdirSync(path.join(workspace, 'child'), { recursive: true }); + fs.mkdirSync(path.join(root, 'workspace-other')); + fs.symlinkSync(workspace, path.join(root, 'alias')); + const config = buildDefaultConfig(root, process.execPath, { nodeBin: process.execPath }); + config.conductors[0].accountPin = router.accountIdentityDigest(account); + if (collector === 'command') { + config.machines[0].claims = { kind: 'command', command: process.execPath, maxAgeMs: 15000, + args: ['-e', `console.log(JSON.stringify({observedAt:${JSON.stringify(new Date(NOW).toISOString())},active:[]}))`] }; + } + const calls = []; const bodies = []; + const native = async (lane, op) => { + calls.push(op.kind); + if (op.kind === 'status') return { ok: true, port: lane.port, codexHome: lane.codexHome, supportedRoles: ['leaf'], threads: {} }; + if (op.method === 'account/read') return { account }; + if (op.method === 'account/rateLimits/read') return { rateLimits: { limitId: 'codex', primary: { usedPercent: 5, windowDurationMins: 300, resetsAt: '2026-09-17T12:00:00Z' } } }; + if (op.kind === 'thread-start') { bodies.push(op.body); return { threadId: `thread-${bodies.length}` }; } + if (op.kind === 'turn-start') return { turnId: `turn-${bodies.length}` }; + throw new Error('unexpected operation'); + }; + const options = { cwd: workspace, workId: 'first', prompt: 'Synthetic ownership check only', nowMs: NOW, stageTimeoutMs: 5000, + capacityProvider: async () => ({ observedAt: new Date(NOW).toISOString(), reachable: true, loadPerCore: 0.2, memoryUsePercent: 35 }), + conductorProvider: native }; + const starts = () => calls.filter((kind) => kind === 'thread-start').length; + return { root, workspace, config, options, native, calls, bodies, starts }; +} +const route = (f, changes) => router.dispatch(f.config, { ...f.options, ...changes }); + +test('command claims collector cannot hide this router\'s dispatched workspace claim', async (t) => { + const f = fixture(t, { collector: 'command' }); + assert.equal((await route(f, {})).receipt.state, 'DISPATCHED'); + await assert.rejects(route(f, { workId: 'second' }), /WORKSPACE_CLAIM_CONFLICT/); + assert.equal(f.starts(), 1); + assert.equal(f.calls.filter((kind) => kind === 'turn-start').length, 1); + const refused = await router.inspectDispatch(f.config, { cwd: f.workspace, workId: 'second' }); + assert.equal(refused.receipt.state, 'PRE_START_FAILED'); + assert.equal(refused.noStartProven, true); +}); +test('command claims collector cannot hide an uncertain launch or its work ID', async (t) => { + const f = fixture(t, { collector: 'command' }); + await assert.rejects(route(f, { conductorProvider: async (lane, op) => { + if (op.kind === 'thread-start') { f.calls.push(op.kind); throw Object.assign(new Error('socket hang up'), { code: 'ECONNRESET' }); } + return f.native(lane, op); + } }), /outcome unknown.*do not retry/i); + assert.equal((await router.inspectDispatch(f.config, { cwd: f.workspace, workId: 'first' })).receipt.state, 'UNKNOWN_DO_NOT_RETRY'); + await assert.rejects(route(f, { workId: 'second' }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { cwd: path.join(f.root, 'workspace-other') }), /WORK_ID_CLAIM_CONFLICT/); + await assert.rejects(route(f, {}), /duplicate intent.*do not retry/i); + assert.equal(f.starts(), 1); +}); +test('a failed local receipt scan refuses dispatch instead of trusting empty external claims', async (t) => { + const f = fixture(t, { collector: 'command' }); + const receipts = path.join(f.config.statePath, 'dispatch-receipts'); + fs.mkdirSync(receipts, { recursive: true, mode: 0o700 }); + fs.writeFileSync(path.join(receipts, 'shared.json'), JSON.stringify({ state: 'DISPATCHED', workId: 'owned', cwd: f.workspace }), { mode: 0o644 }); + await assert.rejects(route(f, { cwd: path.join(f.root, 'workspace-other') }), /UNSAFE_RECEIPT/); + const status = await router.inspectDispatch(f.config, { cwd: path.join(f.root, 'workspace-other'), workId: 'first' }); + assert.equal(status.receipt.state, 'PRE_START_FAILED'); + assert.equal(status.receipt.errorClass, 'UNSAFE_RECEIPT'); + assert.equal(f.starts(), 0); +}); +test('an active local receipt without a workspace fails the ledger scan closed', async (t) => { + const f = fixture(t, { collector: 'command' }); + const receipts = path.join(f.config.statePath, 'dispatch-receipts'); + fs.mkdirSync(receipts, { recursive: true, mode: 0o700 }); + fs.writeFileSync(path.join(receipts, 'no-cwd.json'), JSON.stringify({ state: 'UNKNOWN_DO_NOT_RETRY', workId: 'owned', attemptId: 'owned' }), { mode: 0o600 }); + await assert.rejects(route(f, { cwd: path.join(f.root, 'workspace-other') }), /RECEIPT_CWD_INVALID/); + assert.equal(f.starts(), 0); +}); +test('an alias of a claimed workspace is the same workspace for claims, intents and status', async (t) => { + const f = fixture(t); const alias = path.join(f.root, 'alias'); const canonical = fs.realpathSync.native(f.workspace); + const first = await route(f, {}); + assert.equal(first.receipt.cwd, canonical); + assert.equal(f.bodies[0].cwd, canonical, 'the native thread starts in the claimed identity'); + await assert.rejects(route(f, { cwd: alias, workId: 'second' }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { cwd: alias }), /duplicate intent.*do not retry/i); + for (const cwd of [alias, f.workspace, canonical]) { + assert.equal((await router.inspectDispatch(f.config, { cwd, workId: 'first' })).receipt.attemptId, first.receipt.attemptId); + } + assert.equal(f.starts(), 1); +}); +test('parent and child workspaces conflict in both directions, including through an alias', async (t) => { + const f = fixture(t); const child = path.join(f.workspace, 'child'); + await route(f, {}); + await assert.rejects(route(f, { cwd: child, workId: 'nested' }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { cwd: path.join(f.root, 'alias', 'child'), workId: 'nested-alias' }), /WORKSPACE_CLAIM_CONFLICT/); + const g = fixture(t); + await route(g, { cwd: path.join(g.workspace, 'child') }); + await assert.rejects(route(g, { workId: 'parent' }), /WORKSPACE_CLAIM_CONFLICT/); + assert.equal(f.starts() + g.starts(), 2); +}); +test('sibling worktrees sharing a name prefix remain independently admissible', async (t) => { + const f = fixture(t); + assert.equal((await route(f, {})).receipt.state, 'DISPATCHED'); + assert.equal((await route(f, { cwd: path.join(f.root, 'workspace-other'), workId: 'sibling' })).receipt.state, 'DISPATCHED'); + assert.equal(f.starts(), 2); +}); +test('a letter-case spelling of a claimed workspace conflicts on case-insensitive volumes', async (t) => { + const f = fixture(t); const upper = path.join(f.root, 'WORKSPACE'); + if (!fs.existsSync(upper)) { t.skip('temporary volume is case-sensitive'); return; } + await route(f, {}); + await assert.rejects(route(f, { cwd: upper, workId: 'second' }), /WORKSPACE_CLAIM_CONFLICT/); + assert.equal(f.starts(), 1); +}); +test('a firmlinked spelling realpath keeps distinct conflicts by device and inode', async (t) => { + const f = fixture(t); const firmlinked = path.join('/System/Volumes/Data', fs.realpathSync.native(f.workspace)); + let same = false; + try { const [a, b] = [fs.statSync(f.workspace), fs.statSync(firmlinked)]; same = a.dev === b.dev && a.ino === b.ino; } catch { /* not macOS */ } + if (!same || fs.realpathSync.native(firmlinked) === fs.realpathSync.native(f.workspace)) { t.skip('no distinct firmlinked spelling here'); return; } + await route(f, {}); + await assert.rejects(route(f, { cwd: firmlinked, workId: 'second' }), /WORKSPACE_CLAIM_CONFLICT/); + assert.equal(f.starts(), 1); +}); +test('external claims resolve aliases and deleted paths; malformed ones refuse, absent ones reserve only work IDs', async (t) => { + const f = fixture(t); const other = path.join(f.root, 'workspace-other'); const deleted = path.join(f.workspace, 'deleted', 'deeper'); + const external = (claim) => ({ claimsProvider: async () => claims([{ workId: 'external', state: 'DISPATCHED', attemptId: 'external', ...claim }]) }); + await assert.rejects(route(f, { workId: 'w1', ...external({ cwd: path.join(f.root, 'alias') }) }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { workId: 'w2', ...external({ cwd: deleted }) }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { workId: 'w3', ...external({ cwd: 'relative/workspace' }) }), /WORKSPACE_CLAIM_UNRESOLVED/); + await assert.rejects(route(f, { workId: 'w4', ...external({ cwd: 42 }) }), /WORKSPACE_CLAIM_UNRESOLVED/); + assert.equal(f.starts(), 0); + assert.equal((await route(f, { cwd: other, workId: 'w5', ...external({ cwd: deleted }) })).receipt.state, 'DISPATCHED'); + assert.equal((await route(f, { workId: 'w6', ...external({}) })).receipt.state, 'DISPATCHED'); + await assert.rejects(route(f, { cwd: other, workId: 'w7', ...external({}) }), /WORKSPACE_CLAIM_CONFLICT/); + await assert.rejects(route(f, { workId: 'external', cwd: path.join(f.workspace, 'child'), ...external({}) }), /WORK_ID_CLAIM_CONFLICT/); + assert.equal(f.starts(), 2); +}); +test('an intent recorded under the lexical alias path still blocks replay and stays visible', async (t) => { + const f = fixture(t); const alias = path.join(f.root, 'alias'); + const receipts = path.join(f.config.statePath, 'dispatch-receipts'); const intents = path.join(f.config.statePath, 'intents'); + fs.mkdirSync(receipts, { recursive: true, mode: 0o700 }); fs.mkdirSync(intents, { mode: 0o700 }); + const receiptPath = path.join(receipts, 'legacy.json'); + fs.writeFileSync(receiptPath, JSON.stringify({ schemaVersion: 1, attemptId: 'legacy', workId: 'first', cwd: alias, state: 'PRE_START_FAILED', nativeStartAttempted: false }), { mode: 0o600 }); + const digest = crypto.createHash('sha256').update(JSON.stringify({ workId: 'first', cwd: alias })).digest('hex'); + fs.writeFileSync(path.join(intents, `${digest}.json`), JSON.stringify({ schemaVersion: 1, receiptPath, state: 'PRE_START_FAILED' }), { mode: 0o600 }); + await assert.rejects(route(f, { cwd: alias }), /duplicate intent.*do not retry/i); + const status = await router.inspectDispatch(f.config, { cwd: alias, workId: 'first' }); + assert.equal(status.receipt.attemptId, 'legacy'); + assert.equal(f.starts(), 0); +}); From 5732b722391da56c826a6bb879f6e3716a8ef398 Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 04:16:28 -0600 Subject: [PATCH 3/6] Verify public inventory and native conductor in CI --- .github/workflows/verify.yml | 47 ++++++++++++++++++++++++++++++++++++ 1 file changed, 47 insertions(+) create mode 100644 .github/workflows/verify.yml diff --git a/.github/workflows/verify.yml b/.github/workflows/verify.yml new file mode 100644 index 0000000..9b9b241 --- /dev/null +++ b/.github/workflows/verify.yml @@ -0,0 +1,47 @@ +name: Verify BORG source + +on: + pull_request: + push: + branches: [master] + workflow_dispatch: + +permissions: + contents: read + +jobs: + verify: + name: Source and conductor (${{ matrix.os }}) + strategy: + fail-fast: false + matrix: + os: [ubuntu-24.04, macos-14] + runs-on: ${{ matrix.os }} + timeout-minutes: 15 + steps: + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 # v6 + with: + persist-credentials: false + - uses: actions/setup-node@249970729cb0ef3589644e2896645e5dc5ba9c38 # v6 + with: + node-version: '24.21.0' + - name: Verify exact public distribution + shell: bash + run: | + mkdir "$RUNNER_TEMP/borg-source" + git archive HEAD | tar -x -C "$RUNNER_TEMP/borg-source" + python3 -B "$RUNNER_TEMP/borg-source/installer/release_guard.py" \ + --root "$RUNNER_TEMP/borg-source" \ + --allowlist "$RUNNER_TEMP/borg-source/RELEASE-ALLOWLIST.json" \ + --inventory "$RUNNER_TEMP/borg-source/RELEASE-INVENTORY.json" + - name: Verify release guard and adapter pins + run: python3 -B -m unittest installer.test_release_guard installer.test_adapters + - name: Install pinned native Codex runtime + run: npm ci --prefix installer/npm --ignore-scripts --no-audit --no-fund + - name: Check conductor syntax + run: npm run check --prefix conductor + - name: Test conductor and isolated native initialization + working-directory: conductor + env: + BORG_TEST_CODEX_BIN: ${{ github.workspace }}/installer/npm/node_modules/.bin/codex + run: node --test --test-concurrency=1 tests/*.test.mjs providers/*.test.mjs From 3b8a82b5699ee3b9fcbbec8b570eb79f9109efaa Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 04:26:53 -0600 Subject: [PATCH 4/6] Add independent collision regressions and portable test fixtures --- conductor/docs/LAUNCH-RELIABILITY.md | 13 +- conductor/tests/resume-bookkeeping.test.mjs | 2 +- .../tests/router-collision-review.test.mjs | 367 ++++++++++++++++++ installer/test_fleet_installation.py | 5 + 4 files changed, 381 insertions(+), 6 deletions(-) create mode 100644 conductor/tests/router-collision-review.test.mjs diff --git a/conductor/docs/LAUNCH-RELIABILITY.md b/conductor/docs/LAUNCH-RELIABILITY.md index 5e69b93..1fd838c 100644 --- a/conductor/docs/LAUNCH-RELIABILITY.md +++ b/conductor/docs/LAUNCH-RELIABILITY.md @@ -8,8 +8,8 @@ borg-conductor route-status --config /absolute/BORG_HOME/conductors/config.json --cwd /absolute/workspace --work-id stable-work-id ``` -Use the workspace (any spelling of it resolves to the same filesystem-canonical -path) and stable work ID supplied to `route`. This is +Use the workspace and stable work ID supplied to `route`. For intents created by +this version, every spelling resolves to the same filesystem-canonical path. This is read-only reconciliation: it does not create, resume, cancel or complete a worker. The returned receipt contains a bounded phase history, attempt ID, selected lane, native thread/turn IDs when proved, and error classification. Prompt text is not @@ -61,9 +61,12 @@ where `route` validates the workspace; paths reported for different machines are never assumed disjoint. A claimed workspace that no longer exists still blocks its ancestors. A claim without `cwd` reserves only its work ID; a malformed or unresolvable claimed `cwd` refuses dispatch (`WORKSPACE_CLAIM_UNRESOLVED`). -Intents recorded under a lexical path by earlier versions still block replay and -remain visible to `route-status`. Symlinks inside a workspace that point -elsewhere are not traced. +External collectors must report paths that the router host can resolve; even one +inaccessible path blocks dispatch until the claim or filesystem access is repaired. +Intents recorded by earlier versions are keyed by the spelling used then; that +spelling still blocks replay and finds status. A different spelling is protected +only while the legacy receipt is active, through work-ID and workspace claims. +Symlinks inside a workspace that point elsewhere are not traced. A receipt scan is all-or-error. Malformed JSON, unsafe file permissions, symlinks, unknown receipt states, active receipts without a work ID or absolute workspace, diff --git a/conductor/tests/resume-bookkeeping.test.mjs b/conductor/tests/resume-bookkeeping.test.mjs index 316e107..61d8eb5 100644 --- a/conductor/tests/resume-bookkeeping.test.mjs +++ b/conductor/tests/resume-bookkeeping.test.mjs @@ -24,7 +24,7 @@ test('resume rebuilds bookkeeping from effective native settings and complete re const port = await unusedLoopbackPort(); const historicalCwd = path.join(root, 'historical-worktree'); const effectiveCwd = path.join(root, 'relocated-worktree'); - fs.writeFileSync(fakeCodex, `#!/usr/bin/env node + fs.writeFileSync(fakeCodex, `#!${process.execPath} let input = ''; process.stdin.setEncoding('utf8'); process.stdin.on('data', (chunk) => { diff --git a/conductor/tests/router-collision-review.test.mjs b/conductor/tests/router-collision-review.test.mjs new file mode 100644 index 0000000..2d61441 --- /dev/null +++ b/conductor/tests/router-collision-review.test.mjs @@ -0,0 +1,367 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { describe, test } from 'node:test'; +import * as router from '../router/router.mjs'; +import { buildDefaultConfig } from '../config.mjs'; + +// Independent regressions for review findings R1 (command collector hides the +// router's own durable claims) and R2 (workspace identity is lexical only). +// Every case goes through public dispatch()/inspectDispatch() with the DEFAULT +// claims provider, so a fix may live in observeClaims() or in dispatch(). +// Capacity and the native conductor are synthetic; nothing is launched. +// A fix may either canonicalize workspaces or reject noncanonical spellings: +// both are accepted as long as no second native thread/turn request is made. + +const account = { type: 'chatgpt', email: 'owner@example.test', planType: 'pro' }; +const LIFECYCLE = new Set(['thread-start', 'turn-start']); +const WORKSPACE_REFUSAL = /WORKSPACE|CONFLICT|OVERLAP|ALIAS|SYMLINK|CANONICAL/i; +const UNRESOLVED_REFUSAL = /WORKSPACE|CONFLICT|OVERLAP|ALIAS|SYMLINK|CANONICAL|CLAIM|IDENTITY|UNRESOLV/i; +const WORK_ID_REFUSAL = /WORK_ID|CONFLICT|duplicate intent/i; +const UNREADABLE_CLAIMS = /RECEIPT_[A-Z_]+|UNSAFE_RECEIPT|CLAIMS_[A-Z_]+|LEDGER/; +// Fixed external collector: prints the claims it is given, observed now. +const COLLECTOR = 'process.stdout.write(JSON.stringify({ schemaVersion: 1, observedAt: new Date().toISOString(), ' + + 'source: "synthetic-external", active: JSON.parse(process.argv[1]) }))'; +const TMP_REAL = fs.realpathSync.native(os.tmpdir()); +const CASE_INSENSITIVE = (() => { + const probe = fs.mkdtempSync(path.join(TMP_REAL, 'borg-case-')); + try { fs.mkdirSync(path.join(probe, 'Probe')); return fs.existsSync(path.join(probe, 'PROBE')); } + finally { fs.rmSync(probe, { recursive: true, force: true }); } +})(); + +function fixture(t, { collector = 'router-state' } = {}) { + // Canonical root, so that only the aliases each case creates are noncanonical. + const root = fs.realpathSync.native(fs.mkdtempSync(path.join(os.tmpdir(), 'borg-collision-'))); + fs.chmodSync(root, 0o700); + t.after(() => fs.rmSync(root, { recursive: true, force: true })); + const config = buildDefaultConfig(root, process.execPath, { nodeBin: process.execPath }); + config.conductors[0].accountPin = router.accountIdentityDigest(account); + const external = (active) => { + config.machines[0].claims = { kind: 'command', command: process.execPath, + args: ['-e', COLLECTOR, JSON.stringify(active)], maxAgeMs: 15000 }; + }; + if (collector === 'command') external([]); + const calls = []; + let serial = 0; + const native = (mode) => async (lane, op) => { + calls.push(op.kind); + if (op.kind === 'status') return { ok: true, port: lane.port, codexHome: lane.codexHome, supportedRoles: ['leaf', 'lead'], threads: {} }; + if (op.method === 'account/read') return { account }; + if (op.method === 'account/rateLimits/read') return { rateLimits: { limitId: 'codex', + primary: { usedPercent: 5, windowDurationMins: 300, resetsAt: new Date(Date.now() + 86_400_000).toISOString() } } }; + if (op.kind === 'thread-start') { + if (mode === 'thread-error') throw new Error('synthetic lost thread response'); + return mode === 'thread-no-id' ? {} : { threadId: `thread-${++serial}` }; + } + if (op.kind === 'turn-start') { + if (mode === 'turn-error') throw new Error('synthetic lost turn response'); + return mode === 'turn-no-id' ? {} : { turnId: `turn-${serial}` }; + } + throw new Error('unexpected operation'); + }; + const tree = (...parts) => path.join(root, 'tree', ...parts); + const dir = (...parts) => { fs.mkdirSync(tree(...parts), { recursive: true }); return tree(...parts); }; + const link = (target, ...parts) => { + fs.mkdirSync(path.dirname(tree(...parts)), { recursive: true }); + fs.symlinkSync(target, tree(...parts)); + return tree(...parts); + }; + const send = (cwd, workId, mode = 'ok', extra = {}) => router.dispatch(config, { + cwd, workId, prompt: 'Synthetic collision review only', stageTimeoutMs: 10000, + capacityProvider: async () => ({ observedAt: new Date().toISOString(), reachable: true, loadPerCore: 0.2, memoryUsePercent: 35 }), + conductorProvider: native(mode), ...extra, + }); + return { root, config, calls, tree, dir, link, external, send, + lifecycle: () => calls.filter((kind) => LIFECYCLE.has(kind)).length, + threadStarts: () => calls.filter((kind) => kind === 'thread-start').length }; +} + +const DISPATCHED = { name: 'DISPATCHED', mode: 'ok', state: 'DISPATCHED' }; +const PRIORS = [ + DISPATCHED, + { name: 'UNKNOWN_DO_NOT_RETRY (lost thread response)', mode: 'thread-error', state: 'UNKNOWN_DO_NOT_RETRY' }, + { name: 'UNKNOWN_DO_NOT_RETRY (thread response without ID)', mode: 'thread-no-id', state: 'UNKNOWN_DO_NOT_RETRY' }, + { name: 'STARTED_TURN_UNKNOWN (lost turn response)', mode: 'turn-error', state: 'STARTED_TURN_UNKNOWN' }, + { name: 'STARTED_TURN_UNKNOWN (turn response without ID)', mode: 'turn-no-id', state: 'STARTED_TURN_UNKNOWN' }, + // Crash fixtures: the last state the router persisted before the process died. + { name: 'THREAD_STARTED (router died before turn start)', mode: 'turn-error', state: 'THREAD_STARTED', + crash: { state: 'THREAD_STARTED', phase: 'TURN_START_PENDING', errorClass: null } }, + { name: 'ATTEMPTING (router died during thread start)', mode: 'thread-error', state: 'ATTEMPTING', + crash: { state: 'ATTEMPTING', phase: 'THREAD_START_PENDING', errorClass: null } }, +]; + +// Leave one owned claim in `prior.state` through a real public dispatch. +async function claim(f, cwd, workId, prior = DISPATCHED, { aliasInput = false } = {}) { + const before = f.lifecycle(); + let refusal = null; + await f.send(cwd, workId, prior.mode).catch((error) => { refusal = error; }); + if (aliasInput && f.lifecycle() === before) { + // Rejecting a noncanonical spelling before any native request is safe. + assert.match(String(refusal), WORKSPACE_REFUSAL); + return null; + } + assert.equal(f.threadStarts(), before + 1, `fixture claim did not reach native start: ${refusal}`); + let status = await router.inspectDispatch(f.config, { cwd, workId }); + assert.equal(status.found, true, 'status lookup must find the claim through the spelling used to create it'); + if (prior.crash) { + fs.writeFileSync(status.receiptPath, JSON.stringify({ ...status.receipt, ...prior.crash }), { mode: 0o600 }); + status = await router.inspectDispatch(f.config, { cwd, workId }); + } + assert.equal(status.receipt.state, prior.state); + return status; +} + +async function refused(f, cwd, workId, pattern = WORKSPACE_REFUSAL) { + const before = f.lifecycle(); + await assert.rejects(f.send(cwd, workId), pattern, + `work ${workId} was dispatched into ${path.relative(f.root, cwd)} despite an overlapping active claim`); + assert.equal(f.lifecycle(), before, 'refusal must precede every native thread/turn request'); +} + +async function admitted(f, cwd, workId) { + const before = f.threadStarts(); + const result = await f.send(cwd, workId); + assert.equal(result.receipt.state, 'DISPATCHED'); + assert.equal(f.threadStarts(), before + 1); +} + +function ledger(f) { + const directory = path.join(f.config.statePath, 'dispatch-receipts'); + fs.mkdirSync(directory, { recursive: true, mode: 0o700 }); + return directory; +} + +function putReceipt(f, name, body, mode = 0o600) { + const target = path.join(ledger(f), name); + fs.writeFileSync(target, typeof body === 'string' ? body : JSON.stringify(body), { mode }); + fs.chmodSync(target, mode); + return target; +} + +// Each case owns a private temporary tree, so cases can run concurrently. +describe('router collision review (R1/R2)', { concurrency: 4 }, () => { + + // --- R1: every active local receipt state is a claim, whatever the collector. + + for (const collector of ['router-state', 'command']) { + for (const prior of PRIORS) { + test(`[${collector}] ${prior.name} claim blocks a different work ID in the same workspace`, async (t) => { + const f = fixture(t, { collector }); + const ws = f.dir('ws'); + await claim(f, ws, 'first', prior); + await refused(f, ws, 'second'); + }); + } + for (const prior of [DISPATCHED, PRIORS[1]]) { + test(`[${collector}] ${prior.name} claim reserves its work ID in another workspace`, async (t) => { + const f = fixture(t, { collector }); + await claim(f, f.dir('wt', 'a'), 'shared', prior); + await refused(f, f.dir('wt', 'b'), 'shared', WORK_ID_REFUSAL); + }); + } + } + + test('[command] an empty external view still admits the first dispatch exactly once', async (t) => { + const f = fixture(t, { collector: 'command' }); + await admitted(f, f.dir('ws'), 'only'); + assert.equal(f.lifecycle(), 2); + }); + + test('[command] a failing external collector refuses before native start', async (t) => { + const f = fixture(t, { collector: 'command' }); + f.config.machines[0].claims.args = ['-e', 'process.exit(3)']; + await refused(f, f.dir('ws'), 'first', UNREADABLE_CLAIMS); + }); + + // --- Corrupt local ledger: an unreadable ledger is never an empty claim list. + + const CORRUPT = [ + ['malformed JSON', (f, ws) => putReceipt(f, 'broken.json', `{"state":"DISPATCHED","workId":"owner","cwd":${JSON.stringify(ws)}`)], + ['a non-object JSON record', (f) => putReceipt(f, 'array.json', '[]')], + ['a group/world-readable receipt', (f, ws) => putReceipt(f, 'shared.json', { state: 'DISPATCHED', workId: 'owner', cwd: ws, attemptId: 'a1' }, 0o644)], + ['a symlinked receipt entry', (f, ws) => { + const outside = path.join(f.root, 'outside.json'); + fs.writeFileSync(outside, JSON.stringify({ state: 'DISPATCHED', workId: 'owner', cwd: ws, attemptId: 'a1' }), { mode: 0o600 }); + fs.symlinkSync(outside, path.join(ledger(f), 'redirect.json')); + }], + ['an unknown receipt state', (f, ws) => putReceipt(f, 'future.json', { state: 'RECONCILING', workId: 'owner', cwd: ws, attemptId: 'a1' })], + ['an active receipt without a work ID', (f, ws) => putReceipt(f, 'anonymous.json', { state: 'UNKNOWN_DO_NOT_RETRY', cwd: ws, attemptId: 'a1' })], + ['an oversize receipt', (f, ws) => putReceipt(f, 'huge.json', { state: 'DISPATCHED', workId: 'owner', cwd: ws, attemptId: 'a1', pad: 'x'.repeat(70000) })], + ]; + for (const collector of ['router-state', 'command']) { + for (const [name, corrupt] of CORRUPT) { + test(`[${collector}] local ledger with ${name} refuses before native start`, async (t) => { + const f = fixture(t, { collector }); + const ws = f.dir('ws'); + corrupt(f, ws); + await refused(f, ws, 'second', UNREADABLE_CLAIMS); + }); + } + } + + // --- R2: filesystem identity, alias and ancestor/descendant overlap. + + const OVERLAPS = [ + ['a directory symlink alias of the claimed workspace', (f) => { const ws = f.dir('ws'); return [ws, f.link(ws, 'via-link')]; }], + ['the real path after a claim made through its symlink alias', (f) => { const ws = f.dir('ws'); return [f.link(ws, 'via-link'), ws, true]; }], + ['a child of the claimed workspace', (f) => [f.dir('ws'), f.dir('ws', 'child')]], + ['the parent of a claimed child workspace', (f) => [f.dir('ws', 'child'), f.dir('ws')]], + ['a deep descendant of the claimed workspace', (f) => [f.dir('ws'), f.dir('ws', 'a', 'b', 'c')]], + ['a distant ancestor of a claimed deep workspace', (f) => [f.dir('ws', 'a', 'b', 'c'), f.dir('ws')]], + ['a child reached through a symlinked parent', (f) => { + const ws = f.dir('ws'); f.dir('ws', 'child'); + return [ws, path.join(f.link(ws, 'via-link'), 'child')]; + }], + ['a parent reached through a symlink after its child was claimed', (f) => { + const ws = f.dir('ws'); + return [f.dir('ws', 'child'), f.link(ws, 'via-link')]; + }], + ['an unrelated-looking symlink that resolves inside the claimed tree', (f) => [f.dir('ws'), f.link(f.dir('ws', 'child'), 'elsewhere', 'into')]], + ['the owning parent after a child was claimed through an outside symlink', (f) => [f.link(f.dir('ws', 'child'), 'elsewhere', 'into'), f.dir('ws'), true]], + ]; + const RELATION = Object.fromEntries(OVERLAPS); + const CORE = ['a directory symlink alias of the claimed workspace', 'a child of the claimed workspace', 'the parent of a claimed child workspace']; + + for (const [name, build] of OVERLAPS) { + test(`[router-state] DISPATCHED claim blocks ${name}`, async (t) => { + const f = fixture(t); + const [first, second, aliasInput] = build(f); + if (await claim(f, first, 'first', DISPATCHED, { aliasInput }) === null) return; + await refused(f, second, 'second'); + }); + } + for (const prior of PRIORS.filter((row) => row.mode.endsWith('-error'))) { + for (const name of CORE) { + test(`[router-state] ${prior.name} claim blocks ${name}`, async (t) => { + const f = fixture(t); + const [first, second] = RELATION[name](f); + await claim(f, first, 'first', prior); + await refused(f, second, 'second'); + }); + } + } + for (const name of CORE) { + test(`[command] DISPATCHED local claim blocks ${name}`, async (t) => { + const f = fixture(t, { collector: 'command' }); + const [first, second] = RELATION[name](f); + await claim(f, first, 'first'); + await refused(f, second, 'second'); + }); + } + + test('[router-state] a system-level ancestor symlink (tmpdir spelling) cannot re-enter a claimed workspace', + { skip: os.tmpdir() === TMP_REAL && 'tmpdir is already canonical here' }, async (t) => { + const f = fixture(t); + const ws = f.dir('ws'); + await claim(f, ws, 'first'); + const spelled = path.join(os.tmpdir(), path.relative(TMP_REAL, ws)); + assert.notEqual(spelled, ws); + await refused(f, spelled, 'second'); + }); + + test('[router-state] a case-variant spelling on a case-insensitive volume cannot re-enter a claimed workspace', + { skip: !CASE_INSENSITIVE && 'volume is case-sensitive' }, async (t) => { + const f = fixture(t); + const ws = f.dir('ws'); + await claim(f, ws, 'first'); + await refused(f, f.tree('WS'), 'second'); + }); + + test('[router-state] a claimed child stays owned after its directory disappears', async (t) => { + const f = fixture(t); + const child = f.dir('ws', 'child'); + await claim(f, child, 'first'); + fs.rmSync(child, { recursive: true }); + await refused(f, f.tree('ws'), 'second', UNRESOLVED_REFUSAL); + }); + + test('[router-state] a claim made through a symlink stays owned after the symlink is removed', async (t) => { + const f = fixture(t); + const ws = f.dir('ws'); + const alias = f.link(ws, 'via-link'); + if (await claim(f, alias, 'first', DISPATCHED, { aliasInput: true }) === null) return; + fs.rmSync(alias); + await refused(f, ws, 'second', UNRESOLVED_REFUSAL); + }); + + for (const [name, build] of [ + ['parent and child', (f) => [f.dir('ws'), f.dir('ws', 'child')]], + ['workspace and its symlink alias', (f) => { const ws = f.dir('ws'); return [ws, f.link(ws, 'via-link')]; }], + ]) { + test(`[router-state] concurrent dispatches into ${name} start exactly one native thread`, async (t) => { + const f = fixture(t); + const [one, two] = build(f); + const results = await Promise.allSettled([f.send(one, 'first'), f.send(two, 'second')]); + assert.equal(results.filter((result) => result.status === 'fulfilled').length, 1); + assert.match(String(results.find((result) => result.status === 'rejected').reason), WORKSPACE_REFUSAL); + assert.equal(f.threadStarts(), 1); + }); + } + + test('[router-state] status lookup through a symlink alias finds the same attempt or refuses the alias', async (t) => { + const f = fixture(t); + const ws = f.dir('ws'); + const sent = await claim(f, ws, 'first'); + const alias = f.link(ws, 'via-link'); + let status; + try { status = await router.inspectDispatch(f.config, { cwd: alias, workId: 'first' }); } + catch (error) { assert.match(String(error), WORKSPACE_REFUSAL); return; } + assert.equal(status.found, true, 'an alias of an owned workspace must not read as NOT_FOUND'); + assert.equal(status.receipt.attemptId, sent.receipt.attemptId); + }); + + // --- External collector claims use the same identity rules. + + const externalClaim = (cwd) => [{ workId: 'external-owner', attemptId: 'external-attempt', cwd, state: 'DISPATCHED' }]; + for (const [name, build, pattern] of [ + ['the exact workspace', (f) => { const ws = f.dir('ws'); return [ws, ws]; }], + ['a symlink alias of the workspace', (f) => { const ws = f.dir('ws'); return [f.link(ws, 'via-link'), ws]; }], + ['a parent of the workspace', (f) => [f.dir('ws'), f.dir('ws', 'child')]], + ['a child of the workspace', (f) => [f.dir('ws', 'child'), f.dir('ws')]], + ['a missing child of the workspace', (f) => { const ws = f.dir('ws'); return [path.join(ws, 'gone'), ws]; }, UNRESOLVED_REFUSAL], + ]) { + test(`[command] same-machine external claim on ${name} refuses before native start`, async (t) => { + const f = fixture(t, { collector: 'command' }); + const [claimed, target] = build(f); + f.external(externalClaim(claimed)); + await refused(f, target, 'mine', pattern); + }); + } + + // --- Controls: a fix must not conflate siblings, prefixes or terminal receipts. + + const SIBLINGS = [ + ['sibling worktrees', (f) => [[f.dir('wt', 'a'), 'first'], [f.dir('wt', 'b'), 'second']]], + ['string-prefix siblings (repo then repo-other)', (f) => [[f.dir('repo'), 'first'], [f.dir('repo-other'), 'second']]], + ['string-prefix siblings (repo-other then repo)', (f) => [[f.dir('repo-other'), 'first'], [f.dir('repo'), 'second']]], + ['sibling children of an unclaimed parent', (f) => [[f.dir('ws', 'child-a'), 'first'], [f.dir('ws', 'child-b'), 'second']]], + ]; + for (const collector of ['router-state', 'command']) { + for (const [name, build] of SIBLINGS) { + test(`[${collector}] control: ${name} are independently admissible`, async (t) => { + const f = fixture(t, { collector }); + const [[one, idOne], [two, idTwo]] = build(f); + await admitted(f, one, idOne); + await admitted(f, two, idTwo); + }); + } + test(`[${collector}] control: a PRE_START_FAILED receipt is not an active workspace claim`, async (t) => { + const f = fixture(t, { collector }); + const ws = f.dir('ws'); + await assert.rejects(f.send(ws, 'failed', 'ok', { capacityProvider: async () => null }), /CAPACITY_UNKNOWN/); + assert.equal((await router.inspectDispatch(f.config, { cwd: ws, workId: 'failed' })).receipt.state, 'PRE_START_FAILED'); + await admitted(f, f.dir('ws', 'child'), 'child'); + }); + } + test('[command] control: an external claim on a string-prefix sibling does not block', async (t) => { + const f = fixture(t, { collector: 'command' }); + f.external(externalClaim(f.dir('repo'))); + await admitted(f, f.dir('repo-other'), 'mine'); + }); + test('[command] control: an external claim on a sibling worktree does not block', async (t) => { + const f = fixture(t, { collector: 'command' }); + f.external(externalClaim(f.dir('wt', 'a'))); + await admitted(f, f.dir('wt', 'b'), 'mine'); + }); +}); diff --git a/installer/test_fleet_installation.py b/installer/test_fleet_installation.py index 5bebbc1..6728851 100644 --- a/installer/test_fleet_installation.py +++ b/installer/test_fleet_installation.py @@ -4,6 +4,7 @@ import json import os from pathlib import Path +import sys import tempfile import threading from types import SimpleNamespace @@ -14,6 +15,10 @@ from installer import config from installer import clients from installer.fleet import manage + +# Match the connector import used by installer.fleet without requiring a caller's +# PYTHONPATH to include the source tree's connector directory. +sys.path.insert(0, str(Path(__file__).resolve().parents[1] / "connector")) from fleet_tools import Fleet From c0a7242038d6f0787343c8fb78a30c91802b0646 Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 04:34:35 -0600 Subject: [PATCH 5/6] Check selected ports before provisioning and report unverified startup --- .github/workflows/verify.yml | 4 +- docs/INSTALL.md | 8 ++ install.sh | 12 ++- installer/cli.py | 20 +++- installer/installation.py | 7 +- installer/services.py | 32 +++++- installer/test_startup_preflight.py | 153 ++++++++++++++++++++++++++++ 7 files changed, 221 insertions(+), 15 deletions(-) create mode 100644 installer/test_startup_preflight.py diff --git a/.github/workflows/verify.yml b/.github/workflows/verify.yml index 9b9b241..6501628 100644 --- a/.github/workflows/verify.yml +++ b/.github/workflows/verify.yml @@ -34,8 +34,8 @@ jobs: --root "$RUNNER_TEMP/borg-source" \ --allowlist "$RUNNER_TEMP/borg-source/RELEASE-ALLOWLIST.json" \ --inventory "$RUNNER_TEMP/borg-source/RELEASE-INVENTORY.json" - - name: Verify release guard and adapter pins - run: python3 -B -m unittest installer.test_release_guard installer.test_adapters + - name: Verify release, adapter and installer guards + run: python3 -B -m unittest installer.test_release_guard installer.test_adapters installer.test_startup_preflight - name: Install pinned native Codex runtime run: npm ci --prefix installer/npm --ignore-scripts --no-audit --no-fund - name: Check conductor syntax diff --git a/docs/INSTALL.md b/docs/INSTALL.md index 48f303b..b3c4cae 100644 --- a/docs/INSTALL.md +++ b/docs/INSTALL.md @@ -45,6 +45,11 @@ The home must be an absolute canonical path without symlinks, owned by you and m consecutive numbers starting at `--port-base`. Another installation needs a different home and a non-overlapping port range. +Setup checks the selected service ports before provisioning dependencies. If a +host Python is available, `install.sh` checks before bootstrap downloads too. +An occupied port produces an error naming the service and port; setup does not +stop its current listener. + `--no-start` installs source, dependencies and native configuration without launching services. Run the same installer again without that flag to pull the models, initialize the empty stores and complete setup. @@ -96,6 +101,9 @@ finish with `provider_sign_in_required`. Optional public OAuth still needs a rea client sign-in test. A healthy service process alone does not prove an extraction or remote action succeeded. +`start` reports `started_not_verified` with each service's process state. Run +`doctor` to verify readiness after the services and first brain cycle finish. + ## Reruns, recovery and removal Rerunning the same source release preserves owner credentials, identities and data. diff --git a/install.sh b/install.sh index b2a02fe..e9fdcc3 100755 --- a/install.sh +++ b/install.sh @@ -66,13 +66,15 @@ while IFS= read -r borg_relative; do fi done < "$borg_script_dir/installer/managed-directories.txt" if [ -L "$borg_home/config.json" ]; then echo "BORG config cannot be a symlink" >&2; exit 2; fi -# Blueprint imports must fail before bootstrap creates a home or downloads files. -# A host Python is used only for stdlib validation, never as the installed runtime. -if "$borg_has_blueprint"; then - if ! command -v python3 >/dev/null 2>&1; then +# Use an available host Python for read-only validation before bootstrap writes or +# downloads. Hosts without Python still validate after the minimal bootstrap; +# blueprint imports require validation before any bootstrap work. +if "$borg_has_blueprint" || [ ! -f "$borg_home/config.json" ]; then + if command -v python3 >/dev/null 2>&1; then + PYTHONPATH= PYTHONHOME= python3 -B "$borg_script_dir/borg.py" install --validate-only --home "$borg_home" "$@" + elif "$borg_has_blueprint"; then echo "Blueprint validation requires Python 3 on PATH before bootstrap" >&2; exit 2 fi - PYTHONPATH= PYTHONHOME= python3 -B "$borg_script_dir/borg.py" install --validate-only --home "$borg_home" "$@" fi # An established installation reaches read-only source preflight before downloads. if [ -f "$borg_home/config.json" ] && [ -x "$borg_home/mem0/venv/bin/python" ]; then diff --git a/installer/cli.py b/installer/cli.py index 20c7479..33c1efb 100644 --- a/installer/cli.py +++ b/installer/cli.py @@ -94,6 +94,11 @@ def main(argv: list[str] | None = None) -> int: if existing: from installer.installation import preflight preflight(existing) + # Occupied ports refuse before credentials, downloads or services are created. + from installer import services + services.check_ports(existing or { + "ports": {name: args.port_base + offset for name, offset in config.PORT_OFFSETS.items()}, + "external_access": {"enabled": False}, **({"blueprint": selection} if selection else {})}) if args.validate_only: return 0 if args.command in {"init", "install"}: @@ -119,10 +124,19 @@ def main(argv: list[str] | None = None) -> int: elif args.command == "fleet": from installer.fleet import manage print(json.dumps(manage(doc, args), indent=2)) - elif args.command in {"start", "stop"}: + elif args.command == "start": + from installer import services + result = services.start(doc, args.components or None) + # A live PID is not readiness; only doctor verifies the services. + print(json.dumps({"state": "started_not_verified", + "processes": {name: "process_running" if value == "running" else "process_starting" + for name, value in result.items()}, + "readiness": "not_verified", + "verify": str(Path(doc["home"]) / "bin/borg") + " doctor"}, indent=2)) + elif args.command == "stop": from installer import services - result = getattr(services, args.command)(doc, args.components or None) - print(json.dumps(result or {"state": "stopped"}, indent=2)) + services.stop(doc, args.components or None) + print(json.dumps({"state": "stopped"}, indent=2)) elif args.command == "install": from installer.installation import install return install(doc, system_dependencies=args.system_dependencies, start=not args.no_start) diff --git a/installer/installation.py b/installer/installation.py index 5d8cfe7..af4d0c3 100644 --- a/installer/installation.py +++ b/installer/installation.py @@ -182,6 +182,7 @@ def login(doc: dict, provider: str) -> int: def install(doc: dict, *, system_dependencies: bool = False, start: bool = True) -> int: root = Path(doc["home"]) preflight(doc) + services.check_ports(doc) config.write_connector_config(doc) dependencies.prepare(doc, system_dependencies=system_dependencies) # Bootstrap Python intentionally has no application packages. All native @@ -261,4 +262,8 @@ def complete_install(doc: dict, *, start: bool = True) -> int: parser.add_argument("--no-start", action="store_true") args = parser.parse_args() os.umask(0o077) - raise SystemExit(complete_install(config.load(args.home), start=not args.no_start)) + try: + raise SystemExit(complete_install(config.load(args.home), start=not args.no_start)) + except (ValueError, OSError, RuntimeError) as exc: + print("BORG: " + str(exc), file=sys.stderr) + raise SystemExit(1) diff --git a/installer/services.py b/installer/services.py index ed9a0cd..b065921 100644 --- a/installer/services.py +++ b/installer/services.py @@ -103,6 +103,30 @@ def running(doc: dict, component: str) -> bool: return result.returncode == 0 +def service_ports(doc: dict) -> dict[str, list[int]]: + """Loopback ports each selected service binds, from configuration alone.""" + from installer.blueprint import service_names + keys = {"qdrant": ["qdrant", "qdrant_grpc"], "tunnel": ["tunnel_metrics"]} + return {name: [doc["ports"][key] for key in keys.get(name, [name])] + for name in service_names(doc) if name not in {"brain", "watchdog"}} + + +def occupied(port: int) -> bool: + with socket.socket() as sock: + sock.settimeout(1) + return sock.connect_ex(("127.0.0.1", port)) == 0 + + +def check_ports(doc: dict) -> None: + """Read-only preflight: refuse selected ports held by anything but this installation's own service.""" + conflicts = [f"{port} ({name})" for name, ports in service_ports(doc).items() + if not (doc.get("instance_id") and running(doc, name)) + for port in ports if occupied(port)] + if conflicts: + raise RuntimeError("Loopback ports needed by BORG are already in use: " + ", ".join(conflicts) + + ". BORG will not stop another listener; free those ports or choose another --port-base.") + + def install_definition(doc: dict, name: str, spec: dict) -> Path: root = Path(doc["home"]) logs = root / "logs" @@ -151,10 +175,10 @@ def start(doc: dict, components: list[str] | None = None) -> dict: spec = rows[name] if running(doc, name): continue - if "port" in spec: - with socket.socket() as sock: - if sock.connect_ex(("127.0.0.1", spec["port"])) == 0: - raise RuntimeError(f"Port {spec['port']} is occupied outside this BORG service: {name}") + # Recheck at start: a listener may have appeared since the install preflight. + for port in service_ports(doc).get(name, []): + if occupied(port): + raise RuntimeError(f"Port {port} is occupied outside this BORG service: {name}") target = install_definition(doc, name, spec) if platform.system() == "Darwin": domain = f"gui/{os.getuid()}" diff --git a/installer/test_startup_preflight.py b/installer/test_startup_preflight.py new file mode 100644 index 0000000..fc56def --- /dev/null +++ b/installer/test_startup_preflight.py @@ -0,0 +1,153 @@ +"""Selected-port preflight and start output, using temporary homes and this test's own loopback listeners.""" +import contextlib +import io +import json +import os +from pathlib import Path +import socket +import subprocess +import sys +import tempfile +import unittest +from unittest.mock import patch + +from installer import blueprint, cli, config, installation, services + + +def free_base() -> int: + """Reserve and release a range without depending on the code under test.""" + span = max(config.PORT_OFFSETS.values()) + 1 + for base in range(38000, 60000, span * 3): + with contextlib.ExitStack() as stack: + try: + for offset in range(span): + sock = stack.enter_context(socket.socket()) + sock.bind(("127.0.0.1", base + offset)) + except OSError: + continue + return base + raise unittest.SkipTest("No free loopback port range for this test") + + +@contextlib.contextmanager +def listener(port: int): + with socket.socket() as sock: + sock.bind(("127.0.0.1", port)) + sock.listen() + yield + + +class StartupPreflightTests(unittest.TestCase): + def setUp(self): + temp = tempfile.TemporaryDirectory() + self.addCleanup(temp.cleanup) + self.base = Path(temp.name).resolve() + self.root = self.base / "home" + self.port_base = free_base() + + def doc(self, **extra): + ports = {name: self.port_base + offset for name, offset in config.PORT_OFFSETS.items()} + return {"home": str(self.root), "ports": ports, "external_access": {"enabled": False}, **extra} + + def test_cli_refuses_occupied_selected_port_before_owner_state_or_downloads(self): + args = ["install", "--home", str(self.root), "--owner", "owner", "--port-base", str(self.port_base)] + memory = self.port_base + config.PORT_OFFSETS["memory"] + stderr = io.StringIO() + with listener(memory), \ + patch.object(config, "initialize", side_effect=AssertionError("created owner state")), \ + patch.object(installation.dependencies, "prepare", side_effect=AssertionError("downloaded")), \ + contextlib.redirect_stderr(stderr): + self.assertEqual(cli.main(args), 1) + self.assertEqual(cli.main([*args, "--validate-only"]), 1) + self.assertIn(f"{memory} (memory)", stderr.getvalue()) + self.assertNotIn("Traceback", stderr.getvalue()) + self.assertFalse(self.root.exists()) + + def test_install_refuses_before_connector_config_and_dependencies(self): + doc = self.doc(instance_id="00000000-0000-4000-8000-000000000000") + with listener(doc["ports"]["connector"]), \ + patch.object(installation, "preflight"), \ + patch.object(services, "running", return_value=False), \ + patch.object(config, "write_connector_config", side_effect=AssertionError("wrote config")), \ + patch.object(installation.dependencies, "prepare", side_effect=AssertionError("downloaded")), \ + self.assertRaisesRegex(RuntimeError, "connector"): + installation.install(doc) + + def test_shell_entry_refuses_before_bootstrap_writes_or_downloads(self): + bin_dir = self.base / "bin" + bin_dir.mkdir() + (bin_dir / "python3").symlink_to(sys.executable) + curl = bin_dir / "curl" + curl.write_text('#!/bin/sh\n: > "$BORG_TEST_DOWNLOADED"\nexit 81\n') + curl.chmod(0o700) + downloaded = self.base / "download-attempted" + script = Path(__file__).resolve().parents[1] / "install.sh" + with listener(self.port_base + config.PORT_OFFSETS["qdrant_grpc"]): + result = subprocess.run(["/bin/sh", str(script), "--home", str(self.root), + "--owner", "owner", "--port-base", str(self.port_base)], + capture_output=True, text=True, timeout=10, + env={**os.environ, "PATH": str(bin_dir) + ":/usr/bin:/bin", + "BORG_TEST_DOWNLOADED": str(downloaded)}) + self.assertEqual(result.returncode, 1, result.stderr) + self.assertIn("already in use", result.stderr) + self.assertFalse(self.root.exists()) + self.assertFalse(downloaded.exists()) + + def test_qdrant_grpc_port_conflict_is_refused(self): + doc = self.doc() + with listener(doc["ports"]["qdrant_grpc"]), self.assertRaisesRegex(RuntimeError, r"\(qdrant\)"): + services.check_ports(doc) + + def test_unselected_services_do_not_block(self): + doc = self.doc() + with listener(doc["ports"]["gateway"]), listener(doc["ports"]["tunnel_metrics"]): + services.check_ports(doc) # External access is not enabled. + tools = {"profile": "tools", "components": ["beads"]} + with patch.object(blueprint, "selected_machine", return_value=tools), \ + listener(doc["ports"]["memory"]), listener(doc["ports"]["inbox"]), listener(doc["ports"]["conductor"]): + services.check_ports(doc) + with listener(doc["ports"]["connector"]), self.assertRaisesRegex(RuntimeError, "connector"): + services.check_ports(doc) + with listener(doc["ports"]["inbox"]), self.assertRaisesRegex(RuntimeError, "inbox"): + services.check_ports(doc) + + def test_this_instance_may_rerun_but_an_unowned_listener_is_refused(self): + doc = self.doc(instance_id="00000000-0000-4000-8000-000000000000") + with listener(doc["ports"]["memory"]): + with patch.object(services, "running", side_effect=lambda d, name: name == "memory"): + services.check_ports(doc) + with patch.object(services, "running", return_value=False), self.assertRaisesRegex(RuntimeError, "memory"): + services.check_ports(doc) + # Without an instance the running check is never consulted. + with listener(doc["ports"]["memory"]), \ + patch.object(services, "running", side_effect=AssertionError("consulted")), \ + self.assertRaises(RuntimeError): + services.check_ports(self.doc()) + + def test_start_rechecks_ports_for_a_race_and_installs_nothing(self): + doc = self.doc(instance_id="00000000-0000-4000-8000-000000000000") + rows = {"qdrant": {"port": doc["ports"]["qdrant"], "label": "fixture"}} + with listener(doc["ports"]["qdrant_grpc"]), \ + patch.object(services, "specifications", return_value=rows), \ + patch.object(services, "running", return_value=False), \ + patch.object(services, "install_definition", side_effect=AssertionError("installed")), \ + self.assertRaisesRegex(RuntimeError, "occupied outside this BORG service: qdrant"): + services.start(doc, ["qdrant"]) + + def test_start_output_reports_process_state_without_claiming_readiness(self): + doc = config.initialize(self.root, "owner", port_base=self.port_base) + out = io.StringIO() + with patch.object(services, "start", return_value={"qdrant": "running", "memory": "starting"}), \ + contextlib.redirect_stdout(out): + self.assertEqual(cli.main(["start", "--home", str(self.root)]), 0) + result = json.loads(out.getvalue()) + self.assertEqual(result["state"], "started_not_verified") + self.assertEqual(result["readiness"], "not_verified") + self.assertEqual(result["processes"], {"qdrant": "process_running", "memory": "process_starting"}) + self.assertEqual(result["verify"], str(Path(doc["home"]) / "bin/borg") + " doctor") + self.assertNotIn('"ready"', out.getvalue()) + self.assertNotIn('"running"', out.getvalue()) + + +if __name__ == "__main__": + unittest.main() From 335d34b54d8e80d42790012ebb74f510e73c74f5 Mon Sep 17 00:00:00 2001 From: CryptoJym Date: Thu, 24 Sep 2026 04:36:14 -0600 Subject: [PATCH 6/6] Seal reviewed hardening release inventory --- RELEASE-INVENTORY.json | 95 +++++++++++++++++++++++++++++++----------- 1 file changed, 70 insertions(+), 25 deletions(-) diff --git a/RELEASE-INVENTORY.json b/RELEASE-INVENTORY.json index 53ec96b..55e228a 100644 --- a/RELEASE-INVENTORY.json +++ b/RELEASE-INVENTORY.json @@ -1,12 +1,17 @@ { "excluded_self": "RELEASE-INVENTORY.json", - "file_count": 363, + "file_count": 372, "files": [ { "bytes": 1914, "path": ".github/workflows/site.yml", "sha256": "6bfdccbe0a3994dba7c402c2f7a2e0713de367e92476c68f680bc87439e7d1cf" }, + { + "bytes": 1756, + "path": ".github/workflows/verify.yml", + "sha256": "d97bb88a2b4944ebdaa16b70322c11f1ac568c35d9321cdb9ed133fa4573e38b" + }, { "bytes": 29, "path": ".gitignore", @@ -38,9 +43,9 @@ "sha256": "c9b42bdae7f17d1e4f8bb89728cd7c999a0c8cfe0b7b7b7d91501e49bb25b1d8" }, { - "bytes": 6840, + "bytes": 7171, "path": "adapters/README.md", - "sha256": "7f0062d67950f84fb35bfa3ac58d12f882ad172cc65ff7e76867a3a1b01699e1" + "sha256": "c9ecfef571dff7eef55a1dc30a58cf4e4e032d7dda66c86092140370af5291b9" }, { "bytes": 2767, @@ -128,9 +133,9 @@ "sha256": "e80f558ee37bbaa1fa62f1cf7a3a3831e0c00d1265c09734e800e6e74a0bce81" }, { - "bytes": 4832, + "bytes": 5446, "path": "conductor/README.md", - "sha256": "9f10c8c64e6eda8e3a585eb21146f0e8f496a7d536d08e8ebcb130636d63db3f" + "sha256": "851a3a2f15704bd767d5652d5290f4a14067ee54fcc6f4acac5fa84c6132a37e" }, { "bytes": 1007, @@ -138,9 +143,9 @@ "sha256": "dcac2c44247aea3237cc366a1c97fad7a13dbcdcfe3410957836fe664982fa90" }, { - "bytes": 11381, + "bytes": 11832, "path": "conductor/borg-conductor.mjs", - "sha256": "69595a2e02d1ff6b75ae217cf6c976b850313ec6f2213c5a1eb95b07e101d2e9" + "sha256": "16bfcbc021204e2086955c420d65b7e384cabaa74f9a032438bcf62467771f9d" }, { "bytes": 26418, @@ -162,6 +167,11 @@ "path": "conductor/docs/EVENTS-PAGINATION.md", "sha256": "b90a21d7341dbfd3f8279eaedbb0367c2eed019e9eca3906ecb27f339d3d7b04" }, + { + "bytes": 7276, + "path": "conductor/docs/LAUNCH-RELIABILITY.md", + "sha256": "88e58e20a4817726a851a902193f3700c0c93bcffb471e664c0abb99c156dac8" + }, { "bytes": 2673, "path": "conductor/docs/conductor-0144.md", @@ -198,9 +208,19 @@ "sha256": "dc2f397239ace8ad8f31375c1c9c064148ec7956396be9408b32df15d0a35a4f" }, { - "bytes": 22924, + "bytes": 3139, + "path": "conductor/router/receipt-reader-worker.mjs", + "sha256": "5810e5fbcf95a09c089f74befde10bf691dfad3ebee79c476a2c193cc9dc23a9" + }, + { + "bytes": 2740, + "path": "conductor/router/receipt-reader.mjs", + "sha256": "529a8ddfc636174029dbc45b92f62edc6ce21362ec0c7081ac337ce4f412066d" + }, + { + "bytes": 30846, "path": "conductor/router/router.mjs", - "sha256": "9d6689e9253e89e24351a2c84373b2075593efbed58382f97da18d154702a115" + "sha256": "21dae5b306a32828eaa467dbd906d04229bb93d13997b6709c1298464ac12296" }, { "bytes": 7600, @@ -223,9 +243,29 @@ "sha256": "8e0491bf9e3665e971707d995ed52e0f59d4883d5b03450c57d535178b5b6340" }, { - "bytes": 5925, + "bytes": 3201, + "path": "conductor/tests/receipt-reader.test.mjs", + "sha256": "9ee874660a8b68c2138ffe79a00dc004ae5a1046d1eb5faf5ca3faf37dc008be" + }, + { + "bytes": 5927, "path": "conductor/tests/resume-bookkeeping.test.mjs", - "sha256": "6e098f3d598174a764c07abfa105ec816ccbf66f2b098497ea10bb4c4e43ac62" + "sha256": "da0d0fdb77326135e7a59805d566788872bec30f54562f793894835ec29228ec" + }, + { + "bytes": 19262, + "path": "conductor/tests/router-collision-review.test.mjs", + "sha256": "887b542efcd87b1ee7fab0eff8f0f2c453d4dd90f531bbe1fff9d2618bd27996" + }, + { + "bytes": 11193, + "path": "conductor/tests/router-ownership.test.mjs", + "sha256": "300c2d5bf9b26ffbe1690c0c3af35325007a1622baf15530053f212b345a39d3" + }, + { + "bytes": 12400, + "path": "conductor/tests/router-reliability.test.mjs", + "sha256": "0d6d76b2af5b57a83a82d8a7584a5a4c1015c8cacc4d2370d0bda327872ebea4" }, { "bytes": 10470, @@ -773,9 +813,9 @@ "sha256": "c43600f06704928a08381c34d31471dff166b548f0b242247d06fae377b4e00e" }, { - "bytes": 6630, + "bytes": 7051, "path": "docs/INSTALL.md", - "sha256": "e862aa605e8f83c677254edd7c1c2559ebc78b1848c5d059f584f12ffb57e0b7" + "sha256": "67254a624350f612570ceb0daddac93a1ffbe09533400e5144dbe18cbca24b1e" }, { "bytes": 16398, @@ -908,9 +948,9 @@ "sha256": "6d600ae548f2a53d9f20104e9780ac9d1a1df5d88b8920f34b7fc8a75576c20e" }, { - "bytes": 8772, + "bytes": 8905, "path": "install.sh", - "sha256": "d6e8ac05411b82e78a3125a6919d3441d7e630f8279dbc8da4165db0482eae4e" + "sha256": "2b433c18631eecaa352e08fca17beded2433baf0818c6bbb1179e35816218acf" }, { "bytes": 65, @@ -953,9 +993,9 @@ "sha256": "4a9edbc2d7704f8d88eb4d797c7d50f411533152cfb44da6b6eae45acc43fe9e" }, { - "bytes": 8726, + "bytes": 9714, "path": "installer/cli.py", - "sha256": "848873640e9ca365a4b595cb07a873a4918d596ea00387fc982a0275ec8a797b" + "sha256": "cebd7eacfcb53efe76b666bef26beda8a789afb80dc1f2e6f3b53d170056743f" }, { "bytes": 10302, @@ -998,9 +1038,9 @@ "sha256": "561d1c1c2864a6709189b975e116760fa50342e9b266a6ab7522eb512dba4fb0" }, { - "bytes": 14199, + "bytes": 14377, "path": "installer/installation.py", - "sha256": "ee9b3ee64dfd870b9775912f2c80693b387855bb3d80194f3fe94ccd25597479" + "sha256": "9817bc906501500154933b7f8bdde593aa5ec29aa4830abf3a120093aeff3684" }, { "bytes": 479, @@ -1058,9 +1098,9 @@ "sha256": "fc04cd7fd80c989480c498a11da93f199db8864e9706d6437014789552ab5f9b" }, { - "bytes": 10652, + "bytes": 11836, "path": "installer/services.py", - "sha256": "57221bb53ec2a368bb17491d3a209613539d618914dc27f00e3e4789f0b495c9" + "sha256": "54cfbfaa6dfa7a6e72bd7fbfc1c99501e0c66e2ae972db2c13522a1a9e44de49" }, { "bytes": 1354, @@ -1088,9 +1128,9 @@ "sha256": "3e3c6f2e760d213b9e4405b159b69c60d2ac5bf3a7d51ab8d80fbcb392b0d591" }, { - "bytes": 14901, + "bytes": 15133, "path": "installer/test_fleet_installation.py", - "sha256": "afd694d9d4f8240575f6e5bc633b044a05be0d0782f63c62746def167b926c96" + "sha256": "4148a149631941d2faa59db44672c06638b931221381e9f9f70229bdcd52a942" }, { "bytes": 7977, @@ -1112,6 +1152,11 @@ "path": "installer/test_service_limits.py", "sha256": "387ae7d5523e977c9af1e3ae56936c194887661b0df4ad9e377dcf5337bcfeef" }, + { + "bytes": 7769, + "path": "installer/test_startup_preflight.py", + "sha256": "98a00d31d2dc217a515dca492c82aea9e746acdf7aabe2f1dfd034624ce73709" + }, { "bytes": 3605, "path": "installer/test_web.py", @@ -1818,7 +1863,7 @@ "sha256": "cd4bdb4529012e0cfcd38e059215f9ea433b8ee1fa636276b36c6515e7949e28" } ], - "inventory_sha256": "c7c9d2f6152074f32aa6ab42f474eae1dac852c0c86f652dec91bc281dda232b", + "inventory_sha256": "787241ae9bf10d37224ad69423f5b37cbc3d8110df4c4e7d9461320adf1afbe6", "schema": "borg-public-inventory/v1", - "total_bytes": 87653251 + "total_bytes": 87734443 }