diff --git a/.github/workflows/verify.yml b/.github/workflows/verify.yml index 6501628..8411368 100644 --- a/.github/workflows/verify.yml +++ b/.github/workflows/verify.yml @@ -45,3 +45,102 @@ jobs: 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 + + python-tests: + name: Python tests (${{ matrix.suite }}) + runs-on: ${{ matrix.os }} + timeout-minutes: 30 + env: + PYTHONDONTWRITEBYTECODE: '1' + strategy: + fail-fast: false + matrix: + include: + - suite: connector + os: macos-14 + requirements: installer/requirements.lock + - suite: coordination + os: ubuntu-24.04 + requirements: coordination/requirements-core.lock + - suite: installer + os: macos-14 + requirements: installer/requirements.lock + - suite: memory + os: ubuntu-24.04 + requirements: memory/requirements.txt + - suite: graph + os: ubuntu-24.04 + requirements: graph/requirements.txt + - suite: training + os: ubuntu-24.04 + requirements: training/requirements.txt + steps: + - uses: actions/checkout@d23441a48e516b6c34aea4fa41551a30e30af803 # v6 + with: + persist-credentials: false + - uses: actions/setup-python@a26af69be951a213d495a4c3e4e4022e16d87065 # v5 + with: + python-version: '3.12' + cache: pip + cache-dependency-path: ${{ matrix.requirements }} + - name: Install pinned suite dependencies + shell: bash + run: | + if [[ "${{ matrix.requirements }}" == "installer/requirements.lock" ]]; then + python -m pip install --require-hashes -r "${{ matrix.requirements }}" + else + python -m pip install -r "${{ matrix.requirements }}" + fi + - name: Install pinned coordination test runtime + if: matrix.suite == 'coordination' + env: + BORG_CI_RUNTIME_HOME: ${{ runner.temp }}/borg-coordination-runtime + run: | + python -B - <<'PY' + import os + from pathlib import Path + from installer.downloads import executable, unpack + + root = Path(os.environ["BORG_CI_RUNTIME_HOME"]) + unpack(root, "beads") + print(executable(root, "beads", "bd")) + PY + - name: Prepare portable graph test fixture + if: matrix.suite == 'graph' + env: + BORG_CI_GRAPH_HOME: ${{ runner.temp }}/borg-graph-fixture + run: | + PYTHONPATH=graph/tests python -B - <<'PY' + import os + from pathlib import Path + from test_portable_graph import install_fixture + + install_fixture(Path(os.environ["BORG_CI_GRAPH_HOME"])) + PY + - name: Run ${{ matrix.suite }} suite + shell: bash + run: | + case "${{ matrix.suite }}" in + connector) + PYTHONPATH=connector python -B -m unittest discover -s connector -p 'test_*.py' + ;; + coordination) + PATH="$RUNNER_TEMP/borg-coordination-runtime/runtime/beads:$PATH" \ + PYTHONPATH=coordination python -B -m unittest discover -s coordination/tests -p 'test_*.py' + PATH="$RUNNER_TEMP/borg-coordination-runtime/runtime/beads:$PATH" \ + PYTHONPATH=coordination python -B -m unittest discover -s coordination/comms/hub/tests -p 'test_*.py' + ;; + installer) + python -B -m unittest discover -s installer -p 'test_*.py' + ;; + memory) + python -B -m unittest discover -s memory/tests -p 'test_*.py' + ;; + graph) + BORG_HOME="$RUNNER_TEMP/borg-graph-fixture" GRAPH_FEED_LIVE=1 \ + python -B -m unittest discover -s graph/tests -p 'test_*.py' + ;; + training) + python -B -m unittest discover -s training/tests -p 'test_*.py' + ;; + esac diff --git a/RELEASE-INVENTORY.json b/RELEASE-INVENTORY.json index 39f9c74..1bd1565 100644 --- a/RELEASE-INVENTORY.json +++ b/RELEASE-INVENTORY.json @@ -1,6 +1,6 @@ { "excluded_self": "RELEASE-INVENTORY.json", - "file_count": 396, + "file_count": 400, "files": [ { "bytes": 1914, @@ -8,9 +8,9 @@ "sha256": "6bfdccbe0a3994dba7c402c2f7a2e0713de367e92476c68f680bc87439e7d1cf" }, { - "bytes": 1756, + "bytes": 5498, "path": ".github/workflows/verify.yml", - "sha256": "d97bb88a2b4944ebdaa16b70322c11f1ac568c35d9321cdb9ed133fa4573e38b" + "sha256": "cb97b2990c3af1b3831c9a940df64e7f6def1287ae73e0e718c5a2d2e380a902" }, { "bytes": 29, @@ -128,14 +128,14 @@ "sha256": "da78048f938980354dc1572d7c7b74693a5f672979f2453ab71ddd70eb41ae0f" }, { - "bytes": 5952, + "bytes": 6343, "path": "conductor/INTEGRATION.md", - "sha256": "e80f558ee37bbaa1fa62f1cf7a3a3831e0c00d1265c09734e800e6e74a0bce81" + "sha256": "696bd18dcc8961ebcad6bc32ba4d539465683004761a7766c1de9eb11fc0eaf1" }, { - "bytes": 5446, + "bytes": 6147, "path": "conductor/README.md", - "sha256": "851a3a2f15704bd767d5652d5290f4a14067ee54fcc6f4acac5fa84c6132a37e" + "sha256": "40137a28aba61cf7e3e07fe436a9705b19826888227175669ed321368dcbfb28" }, { "bytes": 1007, @@ -148,19 +148,19 @@ "sha256": "16bfcbc021204e2086955c420d65b7e384cabaa74f9a032438bcf62467771f9d" }, { - "bytes": 26418, + "bytes": 30170, "path": "conductor/conductor.mjs", - "sha256": "7e01bd436fe0c1242ff0128290ed43f05a4123c688b06b6b4d721e56ebbc06c0" + "sha256": "188db1b61294f411ef87bdbf54596585469776c281b7f3ddde98aacdefe29870" }, { - "bytes": 1794, + "bytes": 1827, "path": "conductor/config.example.json", - "sha256": "7c448cc4badca7720cedc54cd5f7fc4e39b28fcdc70d835a46a075e13f3bc3d9" + "sha256": "86e3459fd8ecf7f619323e549d4a910638698c47d0ef4e15eae7a27d6462802a" }, { - "bytes": 12451, + "bytes": 12901, "path": "conductor/config.mjs", - "sha256": "70dd4729fc315ac92566cc9c07b99221e402bad73d43e06a9f02b24e11461db6" + "sha256": "f65fdc54b8ef4d8e153299ab478e87e36a8b7e2813bd5ac77eceb7086baaeacb" }, { "bytes": 730, @@ -168,9 +168,9 @@ "sha256": "b90a21d7341dbfd3f8279eaedbb0367c2eed019e9eca3906ecb27f339d3d7b04" }, { - "bytes": 7276, + "bytes": 8448, "path": "conductor/docs/LAUNCH-RELIABILITY.md", - "sha256": "88e58e20a4817726a851a902193f3700c0c93bcffb471e664c0abb99c156dac8" + "sha256": "c6944b0ff09d2951a502d2bf151c1f58b2a73525f97e8b3f0fec29329789425f" }, { "bytes": 2673, @@ -182,6 +182,11 @@ "path": "conductor/docs/protocol-notes.md", "sha256": "b9cb66dc771323af67a9080c30385ccae63faf209b03a2a3832c719527f857bc" }, + { + "bytes": 3398, + "path": "conductor/http-auth.mjs", + "sha256": "fd0b2eec1968b5d9aa9dbbc39d38f6f88c310d7dc7bfdf372e81cbb8ea5068f6" + }, { "bytes": 506, "path": "conductor/package.json", @@ -208,19 +213,19 @@ "sha256": "dc2f397239ace8ad8f31375c1c9c064148ec7956396be9408b32df15d0a35a4f" }, { - "bytes": 3139, + "bytes": 4169, "path": "conductor/router/receipt-reader-worker.mjs", - "sha256": "5810e5fbcf95a09c089f74befde10bf691dfad3ebee79c476a2c193cc9dc23a9" + "sha256": "87692446392c7454a68d1b9d6d9bef123a1a948d434fd1439305b6ae72d0810f" }, { - "bytes": 2740, + "bytes": 3163, "path": "conductor/router/receipt-reader.mjs", - "sha256": "529a8ddfc636174029dbc45b92f62edc6ce21362ec0c7081ac337ce4f412066d" + "sha256": "6a6e00355c1635f9057d5744553ad9adf562134f43ddc15e97d694d4888053e8" }, { - "bytes": 30846, + "bytes": 47501, "path": "conductor/router/router.mjs", - "sha256": "21dae5b306a32828eaa467dbd906d04229bb93d13997b6709c1298464ac12296" + "sha256": "a162d59f963a63a6d16b87466eeb6e477851ea8042fa321d4cc4484e9862109b" }, { "bytes": 7600, @@ -233,30 +238,40 @@ "sha256": "e88a460e968cd42453b691fade79e97f5ef62af081362bd506fde0b1134742a2" }, { - "bytes": 4707, + "bytes": 5403, "path": "conductor/tests/config.test.mjs", - "sha256": "de4869238830787b4a97e439d6a692b44e3a04c56ab553f8f8cd3c7e090380cf" + "sha256": "031e06a8af00143db161eb4cf9fee785dd274d6e7b18f9c121999bdad1bba8b7" + }, + { + "bytes": 8281, + "path": "conductor/tests/http-auth.test.mjs", + "sha256": "d1363c04e8adcc92b9ce85eb874602aa0ded863bedd7606a08e81ff32178ed70" }, { - "bytes": 3761, + "bytes": 4013, "path": "conductor/tests/native-app-server.test.mjs", - "sha256": "8e0491bf9e3665e971707d995ed52e0f59d4883d5b03450c57d535178b5b6340" + "sha256": "8ba0f91c04a709178958e651568ff07e480a3ad66e8e7bae0c8c28e97840174e" }, { - "bytes": 3201, + "bytes": 4129, "path": "conductor/tests/receipt-reader.test.mjs", - "sha256": "9ee874660a8b68c2138ffe79a00dc004ae5a1046d1eb5faf5ca3faf37dc008be" + "sha256": "aaa5d2f4e69da5b517e917b5b9ee10d284f1cacf8f15acef9592c7ea725d9a76" }, { - "bytes": 5927, + "bytes": 6190, "path": "conductor/tests/resume-bookkeeping.test.mjs", - "sha256": "da0d0fdb77326135e7a59805d566788872bec30f54562f793894835ec29228ec" + "sha256": "6314d14eb14f2e39365a3a1cc124787deccb42da834c7c736da4de634f43fde9" }, { "bytes": 19262, "path": "conductor/tests/router-collision-review.test.mjs", "sha256": "887b542efcd87b1ee7fab0eff8f0f2c453d4dd90f531bbe1fff9d2618bd27996" }, + { + "bytes": 17062, + "path": "conductor/tests/router-lifecycle.test.mjs", + "sha256": "269042e90c330435720b4c6e49aa78d95bce78631a287b3ca793bb6323a45b90" + }, { "bytes": 11193, "path": "conductor/tests/router-ownership.test.mjs", @@ -923,9 +938,9 @@ "sha256": "364603a5328293467b00765f5aa989d294a5ebd8414d1d9b307997b531b3bc00" }, { - "bytes": 2957, + "bytes": 3280, "path": "graph/tests/test_legacy_graphiti_scope_guard.py", - "sha256": "de7da0bf2baae31ab74e1f3afb0f978353fd436834b4be68d0c8c398599007e2" + "sha256": "a107021fef8b8c8d64a465dd5612aa756a037b8bb3f0191747413269aac058a9" }, { "bytes": 7624, @@ -933,9 +948,9 @@ "sha256": "fd62acdc688193d2b0d42c6582882132e735bca9f7384521c47e107809b2dd3f" }, { - "bytes": 18441, + "bytes": 19636, "path": "graph/tests/test_portable_graph.py", - "sha256": "6773dfde9b2187a5e67b261382910e586ad9c8282d8937c13797ec81413a823d" + "sha256": "7233df41ebc3d84167e3de1444fbcb6403ec4d24ca4cc31be45c1ac6b36f15d3" }, { "bytes": 1308, @@ -1033,9 +1048,9 @@ "sha256": "b8573095ab48abf8bede9582452e598afd7118078f6a571d5ee86832f58a7826" }, { - "bytes": 17826, + "bytes": 19270, "path": "installer/health.py", - "sha256": "561d1c1c2864a6709189b975e116760fa50342e9b266a6ab7522eb512dba4fb0" + "sha256": "88b70d35ff8e029a90c0ad70a80871a8aa51aff2098e56acdea7213ca9670c2e" }, { "bytes": 14377, @@ -1122,6 +1137,11 @@ "path": "installer/test_browser.py", "sha256": "5153cd2ec6e58ee83bf1bf43c4f0e1f22a4ef2cdef6ac7285eed1ea6b9188a57" }, + { + "bytes": 1463, + "path": "installer/test_conductor_auth.py", + "sha256": "44d3aaa890404ef1c6d2a4d4796b84cca10b755d71ba1c9e21d5a85a7860f655" + }, { "bytes": 2884, "path": "installer/test_config.py", @@ -1983,7 +2003,7 @@ "sha256": "cd4bdb4529012e0cfcd38e059215f9ea433b8ee1fa636276b36c6515e7949e28" } ], - "inventory_sha256": "23769a5f267e127828ca7376832656052cd92035ab9af3b2537bd13d00cce849", + "inventory_sha256": "927bc4b7fc93fcc354eea91f24337b98d289d616292f38fab72c87a95eb85cac", "schema": "borg-public-inventory/v1", - "total_bytes": 87983389 + "total_bytes": 88047043 } diff --git a/conductor/INTEGRATION.md b/conductor/INTEGRATION.md index 2c809e5..12d8616 100644 --- a/conductor/INTEGRATION.md +++ b/conductor/INTEGRATION.md @@ -18,6 +18,7 @@ $BORG_HOME/ runtime/npm/node_modules/.bin/codex conductors/config.json # conductor-owned schema conductors/primary/profile/ + .conductor/http-token # generated bearer credential, mode 0600 conductors/primary/logs/ policies/SEAT-RULES.md policies/LEAD-RULES.md @@ -92,7 +93,10 @@ CONDUCTOR_LOGS=/absolute/path/to/borg/conductors/primary/logs \ The root installer owns service management. Readiness is the interaction-backed `GET /status` response with `ok: true`, the configured port and the exact -configured `codexHome`; a listening process alone is not readiness proof. +configured `codexHome`; a listening process alone is not readiness proof. The +shipped CLI reads the profile-local bearer token automatically. External local +supervisors may use token-free `GET /healthz` only for process liveness; it is +not conductor initialization or provider readiness proof. ## 3. Status and native authentication @@ -136,10 +140,11 @@ BORG_HOME=/absolute/path/to/borg \ --effort high ``` -`rank` is read-only. `route` is the only dispatch entrypoint; it writes private -intent/receipt evidence and then uses the existing conductor thread and turn -protocol. Root should never dispatch directly to a port when fleet admission -or duplicate protection is required. +`rank` never dispatches work, but it reconciles and may archive private receipt +evidence before scanning claims. `route` is the only dispatch entrypoint; it +writes private intent/receipt evidence and then uses the existing conductor +thread and turn protocol. Root should never dispatch directly to a port when +fleet admission or duplicate protection is required. ## Multi-machine or multi-account extension diff --git a/conductor/README.md b/conductor/README.md index 66d5bee..00abcff 100644 --- a/conductor/README.md +++ b/conductor/README.md @@ -33,6 +33,15 @@ or external queue. A one-machine, one-lane configuration is valid. - Node at the root-provided executable inside `$BORG_HOME/runtime/` - Codex at `$BORG_HOME/runtime/npm/node_modules/.bin/codex` - HTTP bound only to `127.0.0.1:$PORT` +- Per-profile bearer token at `$CODEX_HOME/.conductor/http-token`, mode `0600` + +The conductor creates a 32-byte random token on first start and defaults to +`CONDUCTOR_AUTH=enforce`. All shipped router and installer health calls read +that profile-local token and send it as a bearer credential. `GET /healthz` is +the only token-free endpoint; it returns only `{ "ok": true }`. Host and browser +origin checks apply to every endpoint in both `report` and `enforce` modes. +`POST /rpc` accepts only `account/read`, `account/rateLimits/read`, `model/list`, +`thread/read`, `hooks/list`, and `config/read`. Fresh installs contain no provider authentication. The owner authenticates the dedicated profile with native `codex login`, then pins the observed account @@ -57,8 +66,9 @@ 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, persists an intent before admission scans, -rechecks the configured fleet before selection, and refuses duplicate +Dispatch holds a private lock, reconciles active receipts from exact native +`thread/read` evidence, archives aged terminal receipts, persists an intent +before admission scans, rechecks the configured fleet before selection, and refuses duplicate `{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 diff --git a/conductor/conductor.mjs b/conductor/conductor.mjs index aa80ff7..963f3f5 100644 --- a/conductor/conductor.mjs +++ b/conductor/conductor.mjs @@ -11,7 +11,7 @@ // GET /status conductor + child health, known threads, supported roles // GET /threads?limit=20 thread/list passthrough // GET /events?threadId=&afterSeq= buffered notification stream (per thread) -// POST /rpc {method, params, timeoutMs?} raw JSON-RPC passthrough +// POST /rpc {method, params, timeoutMs?} allowlisted JSON-RPC passthrough // POST /thread/start {cwd, model?, instructions?, role?, sandbox?, personality?} // POST /lead/thread/start {cwd, model?, instructions?, sandbox?, personality?} // POST /thread/resume {threadId, cwd?, role?, model?, predecessor?} @@ -25,10 +25,12 @@ // conductor DENIES it and logs loudly — it never silently grants. import { spawn } from 'node:child_process'; +import crypto from 'node:crypto'; import http from 'node:http'; import fs from 'node:fs'; import path from 'node:path'; import { fileURLToPath } from 'node:url'; +import { ensureHttpToken } from './http-auth.mjs'; export function ownerPolicyPaths(borgHome) { if (typeof borgHome !== 'string' || !path.isAbsolute(borgHome) || path.normalize(borgHome) !== borgHome || path.resolve(borgHome) !== borgHome) { @@ -101,6 +103,45 @@ export const LOG_DIRECTORY_MODE = 0o700; export const LOG_FILE_MODE = 0o600; export const MAX_BODY_BYTES = 64 * 1024; export const MAX_EVENT_PAGE = 200; +export const RPC_METHOD_ALLOWLIST = new Set([ + 'account/read', + 'account/rateLimits/read', + 'model/list', + 'thread/read', + 'hooks/list', + 'config/read', +]); + +function allowedHost(value) { + if (typeof value !== 'string') return false; + return /^(?:(?:127\.0\.0\.1|localhost)(?::[0-9]+)?|\[::1\](?::[0-9]+)?)$/i.test(value); +} + +function browserRequestReason(headers = {}) { + if (Object.hasOwn(headers, 'origin')) return 'origin_forbidden'; + if (Object.hasOwn(headers, 'sec-fetch-site') + && String(headers['sec-fetch-site']).toLowerCase() !== 'none') { + return 'sec_fetch_site_forbidden'; + } + return null; +} + +function jsonContentType(headers = {}) { + const value = headers['content-type']; + return typeof value === 'string' && /^application\/json(?:\s*;|$)/i.test(value.trim()); +} + +function bearerReason(header, expected) { + if (header === undefined) return 'missing_token'; + if (typeof header !== 'string') return 'bad_token'; + const match = /^Bearer ([a-f0-9]{64})$/i.exec(header); + if (!match) return 'bad_token'; + const supplied = Buffer.from(match[1], 'utf8'); + const wanted = Buffer.from(expected, 'utf8'); + return supplied.length === wanted.length && crypto.timingSafeEqual(supplied, wanted) + ? null + : 'bad_token'; +} export async function readBody(req) { const declared = Number(req.headers?.['content-length'] || 0); @@ -390,6 +431,7 @@ const BORG_HOME = options.borgHome || process.env.BORG_HOME; const LOG_DIR = options.logsPath || process.env.CONDUCTOR_LOGS || path.join(process.cwd(), 'logs'); const POLICIES = options.policies || loadOwnerPolicies(BORG_HOME); const EVENT_RING_MAX = 5000; +const AUTH_MODE = options.authMode ?? process.env.CONDUCTOR_AUTH ?? 'enforce'; if (!Number.isInteger(PORT) || PORT < 1024 || PORT > 65535) { throw new Error('CONDUCTOR_PORT must be an integer from 1024 to 65535'); @@ -397,11 +439,23 @@ if (!Number.isInteger(PORT) || PORT < 1024 || PORT > 65535) { if (!CODEX_HOME || !path.isAbsolute(CODEX_HOME)) { throw new Error('CODEX_HOME must be an absolute dedicated profile path'); } +if (!['report', 'enforce'].includes(AUTH_MODE)) { + throw new Error('CONDUCTOR_AUTH must be report or enforce'); +} +const HTTP_TOKEN = ensureHttpToken(CODEX_HOME); ensurePrivateLogDirectory(LOG_DIR); ensurePrivateLogDirectory(path.join(LOG_DIR, 'threads')); const bootTs = new Date().toISOString().replace(/[:.]/g, '-'); const mainLog = openPrivateLogStream(path.join(LOG_DIR, `conductor-${bootTs}.jsonl`)); +const authLog = openPrivateLogStream(path.join(LOG_DIR, `http-auth-${bootTs}.jsonl`)); +const authCounters = { + mode: AUTH_MODE, + unauthenticatedCount: 0, + lastUnauthenticatedAt: null, + disallowedRpcCount: 0, +}; +const authLogTimes = new Map(); function log(kind, data) { const entry = { ts: new Date().toISOString(), kind, data }; @@ -411,6 +465,29 @@ function log(kind, data) { } } +function recordHttpViolation(req, pathname, reason) { + const nowMs = Date.now(); + const ts = new Date(nowMs).toISOString(); + if (reason === 'missing_token' || reason === 'bad_token') { + authCounters.unauthenticatedCount += 1; + authCounters.lastUnauthenticatedAt = ts; + } + if (reason === 'disallowed_rpc_method') authCounters.disallowedRpcCount += 1; + const key = `${pathname}\u0000${reason}`; + const previous = authLogTimes.get(key); + if (previous !== undefined && nowMs - previous < 60_000) return; + authLogTimes.set(key, nowMs); + authLog.write(`${JSON.stringify({ + ts, + method: req.method || null, + path: pathname, + reason, + userAgent: typeof req.headers['user-agent'] === 'string' + ? req.headers['user-agent'].slice(0, 512) + : null, + })}\n`); +} + // ---------- app-server child ---------- const child = spawn(CODEX_BIN, ['app-server'], { stdio: ['pipe', 'pipe', 'pipe'], @@ -551,6 +628,19 @@ function json(res, code, obj) { const server = http.createServer(async (req, res) => { const url = new URL(req.url, `http://${HOST}:${PORT}`); try { + if (!allowedHost(req.headers.host)) return json(res, 403, { error: 'host not allowed' }); + const browserReason = browserRequestReason(req.headers); + if (browserReason) return json(res, 403, { error: 'browser request not allowed' }); + if (req.method === 'GET' && url.pathname === '/healthz') return json(res, 200, { ok: true }); + const authReason = bearerReason(req.headers.authorization, HTTP_TOKEN); + if (authReason) { + recordHttpViolation(req, url.pathname, authReason); + if (AUTH_MODE === 'enforce') return json(res, 401, { error: 'bearer token required' }); + } + if (req.method === 'POST' && !jsonContentType(req.headers)) { + recordHttpViolation(req, url.pathname, 'invalid_content_type'); + if (AUTH_MODE === 'enforce') return json(res, 415, { error: 'application/json required' }); + } await initialized; if (req.method === 'GET' && url.pathname === '/status') { return json(res, 200, { @@ -558,6 +648,7 @@ const server = http.createServer(async (req, res) => { codexHome: CODEX_HOME, supportedRoles: POLICIES.lead ? ['leaf', 'lead'] : ['leaf'], threads: Object.fromEntries(threads), eventSeq: seq, + auth: { ...authCounters }, }); } if (req.method === 'GET' && url.pathname === '/threads') { @@ -576,6 +667,10 @@ const server = http.createServer(async (req, res) => { } if (req.method === 'POST' && url.pathname === '/rpc') { const b = await readBody(req); + if (!RPC_METHOD_ALLOWLIST.has(b.method)) { + recordHttpViolation(req, url.pathname, 'disallowed_rpc_method'); + if (AUTH_MODE === 'enforce') return json(res, 403, { error: 'RPC method not allowed' }); + } return json(res, 200, await rpc(b.method, b.params || {}, b.timeoutMs || 120000)); } if (req.method === 'POST' @@ -659,6 +754,7 @@ server.listen(PORT, HOST, () => { }); const close = () => { + if (closing) return; closing = true; child.kill(); server.close(); diff --git a/conductor/config.example.json b/conductor/config.example.json index 2979742..d3e967e 100644 --- a/conductor/config.example.json +++ b/conductor/config.example.json @@ -19,6 +19,7 @@ "routing": { "timeoutMs": 5000, "lockTimeoutMs": 30000, + "archiveAfterMs": 604800000, "defaultLimitId": "codex", "modelLimitIds": {} }, diff --git a/conductor/config.mjs b/conductor/config.mjs index e197f9e..6d83e2a 100644 --- a/conductor/config.mjs +++ b/conductor/config.mjs @@ -34,6 +34,13 @@ function positiveInteger(value, label, low, high) { return number; } +function strictInteger(value, label, low, high) { + if (typeof value !== 'number' || !Number.isInteger(value) || value < low || value > high) { + throw new Error(`${label} must be an integer from ${low} to ${high}`); + } + return value; +} + function assertNoCredentialKeys(value, label = 'config') { if (!value || typeof value !== 'object') return; for (const [key, child] of Object.entries(value)) { @@ -199,6 +206,8 @@ export function validateInstallConfig(raw, options = {}) { routing: { timeoutMs: positiveInteger(raw.routing?.timeoutMs ?? 5_000, 'routing.timeoutMs', 500, 60_000), lockTimeoutMs: positiveInteger(raw.routing?.lockTimeoutMs ?? 30_000, 'routing.lockTimeoutMs', 1_000, 60_000), + archiveAfterMs: strictInteger(raw.routing?.archiveAfterMs ?? 7 * 24 * 60 * 60 * 1000, + 'routing.archiveAfterMs', 1_000, 365 * 24 * 60 * 60 * 1000), defaultLimitId: requireString(raw.routing?.defaultLimitId ?? 'codex', 'routing.defaultLimitId'), modelLimitIds: { ...(raw.routing?.modelLimitIds ?? {}) }, }, @@ -232,6 +241,7 @@ export function buildDefaultConfig(borgHomeInput, codexBinInput, options = {}) { routing: { timeoutMs: 5_000, lockTimeoutMs: 30_000, + archiveAfterMs: 7 * 24 * 60 * 60 * 1000, defaultLimitId: 'codex', modelLimitIds: {}, }, diff --git a/conductor/docs/LAUNCH-RELIABILITY.md b/conductor/docs/LAUNCH-RELIABILITY.md index 1fd838c..df956ec 100644 --- a/conductor/docs/LAUNCH-RELIABILITY.md +++ b/conductor/docs/LAUNCH-RELIABILITY.md @@ -13,8 +13,9 @@ this version, every spelling resolves to the same filesystem-canonical path. Thi 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. +persisted, only its SHA-256. `completionVerified` is true only for a `COMPLETED` +receipt carrying exact native `thread/read` evidence for its thread and turn. It +does not assert that tests, review or deployment passed. ## States and decisions @@ -26,6 +27,9 @@ is not evidence that a task finished, passed review or was deployed. | `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. | +| `COMPLETED` | Exact native thread/turn readback reported completion. | Inspect the worker artifact and its independent acceptance evidence. | +| `FAILED` | Exact native thread/turn readback reported failure. | Diagnose the native failure before a new dispatch. | +| `CANCELLED` | Exact native thread/turn readback reported cancellation. | Confirm why it was cancelled before a new dispatch. | | `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 @@ -38,8 +42,11 @@ uncertain lifecycle calls are not retried on a different account or machine. 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. +native reconciliation proves a terminal turn state. Before every rank or +dispatch claim scan, the router reads each active receipt's exact native thread +and turn through its conductor. Ambiguous, mismatched or unavailable readback +leaves the receipt active. 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 @@ -81,6 +88,14 @@ 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. +The reader emits early warnings at 8,000 entries and 3.2 MiB. Terminal receipts +older than `routing.archiveAfterMs` (seven days by default) move to the private +`dispatch-receipts/archive/` directory and remain available to `route-status`. + +The private dispatch lock records its owning process. A live holder is never +disturbed. A dead holder, or a lock older than ten minutes with no live holder, +is preserved as `dispatch.lock.stale-` before one acquisition retry. + ## Deadlines and limits `route --stage-timeout-ms N` bounds each admission/provider operation (1–60,000 diff --git a/conductor/http-auth.mjs b/conductor/http-auth.mjs new file mode 100644 index 0000000..90290c5 --- /dev/null +++ b/conductor/http-auth.mjs @@ -0,0 +1,99 @@ +import crypto from 'node:crypto'; +import fs from 'node:fs'; +import path from 'node:path'; + +export const HTTP_TOKEN_DIRECTORY = '.conductor'; +export const HTTP_TOKEN_FILE = 'http-token'; + +function tokenPaths(codexHome) { + if (typeof codexHome !== 'string' || !path.isAbsolute(codexHome)) { + throw new Error('CODEX_HOME must be absolute before reading its conductor token'); + } + const directory = path.join(codexHome, HTTP_TOKEN_DIRECTORY); + return { directory, tokenPath: path.join(directory, HTTP_TOKEN_FILE) }; +} + +function owned(stat, label) { + if (typeof process.getuid === 'function' && stat.uid !== process.getuid()) { + throw new Error(`${label} is owned by another user`); + } +} + +function privateTokenDirectory(directory, allowMissing) { + let stat; + try { + stat = fs.lstatSync(directory); + } catch (error) { + if (allowMissing && error.code === 'ENOENT') return null; + throw error; + } + if (!stat.isDirectory() || stat.isSymbolicLink()) { + throw new Error('conductor token directory must be a regular directory'); + } + owned(stat, 'conductor token directory'); + if ((stat.mode & 0o077) !== 0) throw new Error('conductor token directory must be mode 0700'); + return stat; +} + +export function httpTokenPath(codexHome) { + return tokenPaths(codexHome).tokenPath; +} + +export function readHttpToken(codexHome, options = {}) { + const { directory, tokenPath } = tokenPaths(codexHome); + const allowMissing = options.allowMissing === true; + if (!privateTokenDirectory(directory, allowMissing)) return null; + let fd; + try { + fd = fs.openSync(tokenPath, fs.constants.O_RDONLY | (fs.constants.O_NOFOLLOW ?? 0)); + } catch (error) { + if (allowMissing && error.code === 'ENOENT') return null; + throw error; + } + try { + const stat = fs.fstatSync(fd); + if (!stat.isFile() || stat.nlink !== 1 || (stat.mode & 0o077) !== 0 || stat.size > 128) { + throw new Error('conductor token must be one owner-only regular file'); + } + owned(stat, 'conductor token'); + const buffer = Buffer.alloc(129); + const bytes = fs.readSync(fd, buffer, 0, buffer.length, 0); + if (bytes > 128) throw new Error('conductor token is oversized'); + const token = buffer.subarray(0, bytes).toString('utf8').trim(); + if (!/^[a-f0-9]{64}$/.test(token)) throw new Error('conductor token is malformed'); + return token; + } finally { + fs.closeSync(fd); + } +} + +export function ensureHttpToken(codexHome) { + const { directory, tokenPath } = tokenPaths(codexHome); + fs.mkdirSync(directory, { recursive: true, mode: 0o700 }); + const stat = fs.lstatSync(directory); + if (!stat.isDirectory() || stat.isSymbolicLink()) { + throw new Error('conductor token directory must be a regular directory'); + } + owned(stat, 'conductor token directory'); + fs.chmodSync(directory, 0o700); + try { + return readHttpToken(codexHome); + } catch (error) { + if (error.code !== 'ENOENT') throw error; + } + + const token = crypto.randomBytes(32).toString('hex'); + let fd; + try { + fd = fs.openSync(tokenPath, fs.constants.O_WRONLY | fs.constants.O_CREAT + | fs.constants.O_EXCL | (fs.constants.O_NOFOLLOW ?? 0), 0o600); + fs.writeSync(fd, `${token}\n`); + fs.fsyncSync(fd); + fs.fchmodSync(fd, 0o600); + } catch (error) { + if (error.code !== 'EEXIST') throw error; + } finally { + if (fd !== undefined) fs.closeSync(fd); + } + return readHttpToken(codexHome); +} diff --git a/conductor/router/receipt-reader-worker.mjs b/conductor/router/receipt-reader-worker.mjs index d4adcc3..5702406 100644 --- a/conductor/router/receipt-reader-worker.mjs +++ b/conductor/router/receipt-reader-worker.mjs @@ -5,6 +5,7 @@ import path from 'node:path'; const MAX_ENTRIES = 10000; const MAX_FILE_BYTES = 65536; const MAX_TOTAL_BYTES = 4 * 1024 * 1024; +const WARNING_FRACTION = 0.8; function reject(code) { const error = new Error(code); error.code = code; throw error; } function privateStat(target, directory = false) { const stat = fs.lstatSync(target); @@ -33,10 +34,27 @@ function record(target) { return { value, bytes }; } finally { fs.closeSync(fd); } } -function directory(target) { +function scanWarnings(entries, totalBytes) { + const warnings = []; + if (entries >= Math.ceil(MAX_ENTRIES * WARNING_FRACTION)) warnings.push({ + code: 'RECEIPT_ENTRY_LIMIT_WARNING', + message: `receipt scan has ${entries} entries; hard cap is ${MAX_ENTRIES}`, + entries, + limit: MAX_ENTRIES, + }); + if (totalBytes >= Math.ceil(MAX_TOTAL_BYTES * WARNING_FRACTION)) warnings.push({ + code: 'RECEIPT_TOTAL_LIMIT_WARNING', + message: `receipt scan has ${totalBytes} JSON bytes; hard cap is ${MAX_TOTAL_BYTES}`, + totalBytes, + limit: MAX_TOTAL_BYTES, + }); + return warnings; +} +function directory(target, includeEntries = false) { privateStat(target, true); const handle = fs.opendirSync(target); const values = []; + const entries = []; let seen = 0; let total = 0; try { let entry; @@ -48,24 +66,30 @@ function directory(target) { total += read.bytes; if (total > MAX_TOTAL_BYTES) reject('RECEIPT_TOTAL_LIMIT'); values.push(read.value); + if (includeEntries) entries.push({ name: entry.name, value: read.value, bytes: read.bytes }); } } finally { handle.closeSync(); } - return values; + return { value: includeEntries ? entries : values, warnings: scanWarnings(seen, total) }; } const [mode, target] = process.argv.slice(2); try { - if (!['record','directory'].includes(mode) || !path.isAbsolute(target ?? '')) reject('RECEIPT_REQUEST_INVALID'); + if (!['record','directory','entries'].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.stdout.write(JSON.stringify({ ok: true, value: mode === 'record' ? null : [], warnings: [] })); process.exit(0); } - value = mode === 'directory' ? directory(target) : record(target).value; - process.stdout.write(JSON.stringify({ ok: true, value })); + if (mode === 'record') { + value = record(target).value; + process.stdout.write(JSON.stringify({ ok: true, value })); + } else { + const result = directory(target, mode === 'entries'); + process.stdout.write(JSON.stringify({ ok: true, value: result.value, warnings: result.warnings })); + } } 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 index d0be981..e833c45 100644 --- a/conductor/router/receipt-reader.mjs +++ b/conductor/router/receipt-reader.mjs @@ -56,6 +56,11 @@ function readInChild(mode, target, options = {}) { const errorCode = /^[A-Z_]+$/.test(envelope.code ?? '') ? envelope.code : 'RECEIPT_READ_FAILED'; fail(errorCode); return; } + const warnings = Array.isArray(envelope.warnings) ? envelope.warnings : []; + for (const warning of warnings) { + if (typeof options.onWarning === 'function') options.onWarning(warning); + else process.emitWarning(warning.message, { code: warning.code }); + } finish(null, envelope.value); }); }); @@ -65,6 +70,10 @@ export async function readPrivateJsonDirectory(directory, options = {}) { return readInChild('directory', directory, options); } +export async function readPrivateJsonDirectoryEntries(directory, options = {}) { + return readInChild('entries', 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 2669122..782d9d9 100644 --- a/conductor/router/router.mjs +++ b/conductor/router/router.mjs @@ -1,12 +1,17 @@ import { execFile as execFileCallback } from 'node:child_process'; import crypto from 'node:crypto'; -import { promises as fs } from 'node:fs'; +import { constants as fsConstants, promises as fs } from 'node:fs'; import os from 'node:os'; import path from 'node:path'; import { promisify } from 'node:util'; import { validateInstallConfig } from '../config.mjs'; -import { readPrivateJsonDirectory, readPrivateJsonRecord } from './receipt-reader.mjs'; +import { readHttpToken } from '../http-auth.mjs'; +import { + readPrivateJsonDirectory, + readPrivateJsonDirectoryEntries, + readPrivateJsonRecord, +} from './receipt-reader.mjs'; // Portable extraction of the supplied conductor-usage-router and // fleet-placement invariants. Estate-specific roster, SSH shipping, and @@ -21,6 +26,14 @@ const ACTIVE_RECEIPT_STATES = new Set([ 'UNKNOWN_DO_NOT_RETRY', 'STARTED_TURN_UNKNOWN', ]); +const TERMINAL_RECEIPT_STATES = new Set([ + 'PRE_START_FAILED', + 'COMPLETED', + 'FAILED', + 'CANCELLED', + 'RECONCILED_NO_START', +]); +const STALE_LOCK_AGE_MS = 10 * 60 * 1000; function sha256(value) { return crypto.createHash('sha256').update(value).digest('hex'); @@ -142,9 +155,8 @@ export async function observeClaims(machine, statePath, nowMs = Date.now()) { 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) { - if (!ACTIVE_RECEIPT_STATES.has(receipt.state) && !terminal.has(receipt.state)) { + if (!ACTIVE_RECEIPT_STATES.has(receipt.state) && !TERMINAL_RECEIPT_STATES.has(receipt.state)) { throw new Error('RECEIPT_STATE_UNKNOWN'); } if (ACTIVE_RECEIPT_STATES.has(receipt.state) && (typeof receipt.workId !== 'string' || !receipt.workId.trim())) { @@ -169,8 +181,8 @@ async function observeLocalReceiptClaims(statePath, nowMs = Date.now()) { }; } -async function fetchJson(url, init, timeoutMs) { - const response = await fetch(url, { ...init, signal: AbortSignal.timeout(timeoutMs) }); +async function fetchJson(url, init, timeoutMs, fetchImpl = globalThis.fetch) { + const response = await fetchImpl(url, { ...init, signal: AbortSignal.timeout(timeoutMs) }); if (!response.ok) throw new Error(`HTTP_${response.status}`); try { return await response.json(); @@ -179,34 +191,36 @@ async function fetchJson(url, init, timeoutMs) { } } -export async function nativeConductorProvider(lane, operation, timeoutMs = 5_000) { +export async function nativeConductorProvider(lane, operation, timeoutMs = 5_000, fetchImpl = globalThis.fetch) { const base = `http://${lane.host}:${lane.port}`; + const token = readHttpToken(lane.codexHome, { allowMissing: true }); + const headers = token ? { authorization: `Bearer ${token}` } : {}; if (operation.kind === 'status') { - const status = await fetchJson(`${base}/status`, { method: 'GET' }, timeoutMs); - const nativeThreads = await fetchJson(`${base}/threads?limit=100`, { method: 'GET' }, timeoutMs); + const status = await fetchJson(`${base}/status`, { method: 'GET', headers }, timeoutMs, fetchImpl); + const nativeThreads = await fetchJson(`${base}/threads?limit=100`, { method: 'GET', headers }, timeoutMs, fetchImpl); return { ...status, nativeThreads }; } if (operation.kind === 'rpc') { return fetchJson(`${base}/rpc`, { method: 'POST', - headers: { 'content-type': 'application/json' }, + headers: { ...headers, 'content-type': 'application/json' }, body: JSON.stringify({ method: operation.method, params: operation.params ?? {}, timeoutMs }), - }, timeoutMs); + }, timeoutMs, fetchImpl); } if (operation.kind === 'thread-start') { const endpoint = operation.role === 'lead' ? '/lead/thread/start' : '/thread/start'; return fetchJson(`${base}${endpoint}`, { method: 'POST', - headers: { 'content-type': 'application/json' }, + headers: { ...headers, 'content-type': 'application/json' }, body: JSON.stringify(operation.body), - }, 60_000); + }, 60_000, fetchImpl); } if (operation.kind === 'turn-start') { return fetchJson(`${base}/turn/start`, { method: 'POST', - headers: { 'content-type': 'application/json' }, + headers: { ...headers, 'content-type': 'application/json' }, body: JSON.stringify(operation.body), - }, 60_000); + }, 60_000, fetchImpl); } throw new Error(`unknown conductor operation: ${operation.kind}`); } @@ -416,18 +430,229 @@ async function atomicPrivateWrite(target, value, options = {}) { await fs.rename(temporary, target); } -async function acquireLock(statePath, timeoutMs) { +function processIsAlive(pid) { + if (!Number.isInteger(pid) || pid < 1) return false; + try { + process.kill(pid, 0); + return true; + } catch (error) { + if (error.code === 'ESRCH') return false; + // EPERM proves that a process exists even when this user cannot signal it. + return true; + } +} + +async function readProcessStartTime(pid) { + if (!Number.isInteger(pid) || pid < 1) return null; + try { + if (process.platform === 'linux') { + const [stat, bootIdBody] = await Promise.all([ + fs.readFile(`/proc/${pid}/stat`, 'utf8'), + fs.readFile('/proc/sys/kernel/random/boot_id', 'utf8'), + ]); + const commandEnd = stat.lastIndexOf(')'); + if (commandEnd < 0) return null; + // Fields following the command start at field 3; process start time is field 22. + const fields = stat.slice(commandEnd + 1).trim().split(/\s+/); + const startTime = fields[19]; + const bootId = bootIdBody.trim().toLowerCase(); + return /^\d+$/.test(startTime ?? '') && /^[0-9a-f-]+$/.test(bootId) + ? `linux-proc:${bootId}:${startTime}` + : null; + } + if (process.platform === 'darwin') { + const { stdout } = await execFile('/bin/ps', ['-o', 'lstart=', '-p', String(pid)], { + encoding: 'utf8', + env: { LC_ALL: 'C', TZ: 'UTC' }, + maxBuffer: 16 * 1024, + }); + const startTime = stdout.trim().replace(/\s+/g, ' '); + return startTime ? `darwin-ps:${startTime}` : null; + } + } catch { + // A vanished process or unavailable native reading is not identity evidence. + } + return null; +} + +function lockOwner(body) { + const trimmed = body.trim(); + if (!trimmed) return null; + if (/^[1-9][0-9]*$/.test(trimmed)) { + return { pid: Number(trimmed), processStartTime: null }; + } + try { + const value = JSON.parse(trimmed); + if (!Number.isInteger(value?.pid) || value.pid < 1) return null; + const processStartTime = typeof value.processStartTime === 'string' + && value.processStartTime.trim() + ? value.processStartTime.trim() + : null; + return { pid: value.pid, processStartTime }; + } catch { + return null; + } +} + +const LOCK_READ_FLAGS = fsConstants.O_RDONLY | (fsConstants.O_NOFOLLOW ?? 0); + +async function readLockSnapshot(lockPath) { + let handle; + try { + handle = await fs.open(lockPath, LOCK_READ_FLAGS); + const before = await handle.stat({ bigint: true }); + if (!before.isFile() || (before.mode & 0o077n) !== 0n || before.size > 4096n) { + return { valid: false, stat: before }; + } + const bytes = await handle.readFile(); + const after = await handle.stat({ bigint: true }); + const unchanged = before.dev === after.dev + && before.ino === after.ino + && before.size === after.size + && before.mtimeNs === after.mtimeNs + && before.ctimeNs === after.ctimeNs + && BigInt(bytes.length) === after.size; + if (!unchanged) return { valid: false, stat: after }; + const body = bytes.toString('utf8'); + return { + valid: true, + stat: after, + body: bytes, + owner: lockOwner(body), + mtimeMs: Number(after.mtimeNs / 1_000_000n), + }; + } catch (error) { + if (error.code === 'ENOENT') return null; + if (error.code === 'ELOOP') return { valid: false, stat: null }; + throw error; + } finally { + await handle?.close(); + } +} + +function sameLockInode(left, right) { + return left?.stat && right?.stat + && left.stat.dev === right.stat.dev + && left.stat.ino === right.stat.ino; +} + +function sameLockOwner(left, right) { + return left?.valid === true + && right?.valid === true + && sameLockInode(left, right) + && left.stat.mtimeNs === right.stat.mtimeNs + && left.body.equals(right.body); +} + +async function lockIsStale(snapshot, options, nowMs) { + if (snapshot?.valid !== true) return false; + if (snapshot.owner === null) return nowMs - snapshot.mtimeMs > STALE_LOCK_AGE_MS; + const alive = await (options.processAlive ?? processIsAlive)(snapshot.owner.pid); + if (!alive) return true; + if (snapshot.owner.processStartTime === null) return false; + const processStartTime = await (options.processStartTime ?? readProcessStartTime)(snapshot.owner.pid); + const normalizedStartTime = typeof processStartTime === 'string' ? processStartTime.trim() : ''; + return normalizedStartTime !== '' && normalizedStartTime !== snapshot.owner.processStartTime; +} + +async function recoverStaleLock(lockPath, options = {}) { + const observed = await readLockSnapshot(lockPath); + if (observed === null) return true; + const nowMs = options.nowMs ?? Date.now(); + if (!await lockIsStale(observed, options, nowMs)) return false; + + // The inode modification time makes the recovery name stable across contenders, + // so only one hard link can elect a reaper for this exact lock. + const base = `${lockPath}.stale-${new Date(observed.mtimeMs).toISOString()}`; + let stalePath; + for (let suffix = 0; ; suffix += 1) { + const candidate = suffix === 0 ? base : `${base}-${suffix}`; + try { + await fs.link(lockPath, candidate); + stalePath = candidate; + break; + } catch (error) { + if (error.code === 'ENOENT') return true; + if (error.code === 'EEXIST') { + const existing = await readLockSnapshot(candidate); + if (sameLockInode(existing, observed)) return false; + continue; + } + throw error; + } + } + + const removeRecoveryLink = async () => { + await fs.unlink(stalePath).catch((error) => { + if (error.code !== 'ENOENT') throw error; + }); + }; + const [linked, current] = await Promise.all([ + readLockSnapshot(stalePath), + readLockSnapshot(lockPath), + ]); + const safeToRetire = sameLockOwner(linked, observed) + && sameLockOwner(current, observed) + && await lockIsStale(current, options, nowMs); + if (!safeToRetire) { + await removeRecoveryLink(); + return false; + } + + // Re-read immediately before retirement so a replacement inode or owner is + // never removed based on an earlier pathname observation. + const finalCurrent = await readLockSnapshot(lockPath); + if (!sameLockOwner(finalCurrent, observed) + || !await lockIsStale(finalCurrent, options, nowMs)) { + await removeRecoveryLink(); + return false; + } + try { + await fs.unlink(lockPath); + return true; + } catch (error) { + if (error.code === 'ENOENT') return true; + await removeRecoveryLink(); + throw error; + } +} + +export async function acquireDispatchLock(statePath, timeoutMs, options = {}) { await ensurePrivateDirectory(statePath); const lockPath = path.join(statePath, 'dispatch.lock'); + const lockId = crypto.randomUUID(); + const processStartTime = await (options.processStartTime ?? readProcessStartTime)(process.pid); + if (typeof processStartTime !== 'string' || !processStartTime.trim()) { + throw new Error('cannot read this process start time for dispatch lock ownership'); + } const deadline = Date.now() + timeoutMs; + let recoveryAttempted = false; while (true) { try { const handle = await fs.open(lockPath, 'wx', 0o600); - await handle.writeFile(`${process.pid}\n`); - await handle.close(); - return () => fs.unlink(lockPath); + try { + await handle.writeFile(`${JSON.stringify({ + schemaVersion: 1, + lockId, + pid: process.pid, + processStartTime: processStartTime.trim(), + acquiredAt: new Date(options.nowMs ?? Date.now()).toISOString(), + }, null, 2)}\n`); + await handle.sync(); + } finally { + await handle.close(); + } + return async () => { + const current = JSON.parse(await fs.readFile(lockPath, 'utf8')); + if (current?.lockId !== lockId) throw new Error('dispatch lock ownership changed'); + await fs.unlink(lockPath); + }; } catch (error) { if (error.code !== 'EEXIST') throw error; + if (!recoveryAttempted && await recoverStaleLock(lockPath, options)) { + recoveryAttempted = true; + continue; + } if (Date.now() >= deadline) throw new Error(`dispatch lock busy; inspect ${lockPath}`); await new Promise((resolve) => setTimeout(resolve, 50)); } @@ -517,6 +742,179 @@ async function updateReceipt(receiptPath, intentPath, receipt) { await atomicPrivateWrite(intentPath, { schemaVersion: 1, receiptPath, state: receipt.state }); } +function statusWords(value) { + const values = []; + if (typeof value === 'string') values.push(value); + if (value && typeof value === 'object' && !Array.isArray(value)) { + for (const key of ['type', 'status', 'state']) { + if (typeof value[key] === 'string') values.push(value[key]); + } + } + return values.map((item) => item.toLowerCase().replace(/[^a-z]+/g, '_')); +} + +function terminalStateFromTurn(turn) { + const words = new Set(statusWords(turn?.status)); + if (['interrupted', 'cancelled', 'canceled', 'aborted'].some((word) => words.has(word))) { + return 'CANCELLED'; + } + if (['failed', 'failure', 'error', 'errored'].some((word) => words.has(word))) return 'FAILED'; + if (['completed', 'complete', 'succeeded', 'success', 'done'].some((word) => words.has(word))) { + return 'COMPLETED'; + } + return null; +} + +function normalizedNativeTime(value, fallback) { + if (typeof value === 'number' && Number.isFinite(value)) { + const milliseconds = value < 10_000_000_000 ? value * 1000 : value; + if (Number.isFinite(milliseconds)) return new Date(milliseconds).toISOString(); + } + if (typeof value === 'string' && Number.isFinite(Date.parse(value))) { + return new Date(value).toISOString(); + } + return fallback; +} + +async function updateIntentState(config, receiptPath, receipt) { + if (typeof receipt.workId !== 'string' || typeof receipt.cwd !== 'string') return; + const intentPath = path.join(config.statePath, 'intents', `${intentDigest(receipt.workId, receipt.cwd)}.json`); + const intent = await readPrivateJsonRecord(intentPath); + if (intent === null || intent.receiptPath !== receiptPath) return; + await atomicPrivateWrite(intentPath, { schemaVersion: 1, receiptPath, state: receipt.state }); +} + +async function reconcileReceipt(config, entry, options) { + const receipt = entry.value; + if (!ACTIVE_RECEIPT_STATES.has(receipt.state) + || typeof receipt.threadId !== 'string' || !receipt.threadId + || typeof receipt.laneId !== 'string') return { receipt, changed: false }; + const lane = config.conductors.find((candidate) => candidate.id === receipt.laneId); + if (!lane) return { receipt, changed: false, error: 'LANE_NOT_FOUND' }; + let readback; + try { + readback = await stageDeadline(() => options.conductorProvider(lane, { + kind: 'rpc', + method: 'thread/read', + params: { threadId: receipt.threadId, includeTurns: true }, + }), 'RECONCILE_READ', options.stageTimeoutMs); + } catch (error) { + return { receipt, changed: false, error: error.code || error.name || 'READ_FAILED' }; + } + const thread = readback?.thread; + if (!thread || thread.id !== receipt.threadId || !Array.isArray(thread.turns)) { + return { receipt, changed: false, error: 'THREAD_READ_INVALID' }; + } + const turn = receipt.turnId + ? thread.turns.find((candidate) => candidate?.id === receipt.turnId) + : thread.turns.at(-1); + if (!turn || typeof turn.id !== 'string' || !turn.id) { + return { receipt, changed: false, error: 'TURN_NOT_FOUND' }; + } + const state = terminalStateFromTurn(turn); + if (!state) return { receipt, changed: false }; + const reconciledAt = new Date(options.nowMs).toISOString(); + const updated = { + ...receipt, + state, + phase: 'RECONCILED', + turnId: turn.id, + reconciledAt, + terminalAt: normalizedNativeTime(turn.completedAt, reconciledAt), + nativeEvidence: { + method: 'thread/read', + threadId: thread.id, + turnId: turn.id, + status: statusWords(turn.status)[0] ?? null, + observedAt: reconciledAt, + }, + events: [...(Array.isArray(receipt.events) ? receipt.events : []), { + phase: 'RECONCILED', at: reconciledAt, state, + }].slice(-16), + }; + const receiptPath = path.join(config.statePath, 'dispatch-receipts', entry.name); + await atomicPrivateWrite(receiptPath, updated); + await updateIntentState(config, receiptPath, updated); + return { receipt: updated, changed: true }; +} + +function terminalTimestamp(receipt) { + for (const value of [receipt.terminalAt, receipt.reconciledAt, receipt.dispatchedAt, receipt.attemptedAt]) { + const timestamp = timestampMillis(value); + if (timestamp !== null) return timestamp; + } + return null; +} + +async function archiveReceipt(config, entry, receipt) { + const receiptsDirectory = path.join(config.statePath, 'dispatch-receipts'); + const archiveDirectory = path.join(receiptsDirectory, 'archive'); + await ensurePrivateDirectory(archiveDirectory); + const receiptPath = path.join(receiptsDirectory, entry.name); + const archivePath = path.join(archiveDirectory, entry.name); + try { + await fs.access(archivePath); + throw new Error('ARCHIVE_RECEIPT_EXISTS'); + } catch (error) { + if (error.code !== 'ENOENT') throw error; + } + await fs.rename(receiptPath, archivePath); + if (typeof receipt.workId === 'string' && typeof receipt.cwd === 'string') { + const intentPath = path.join(config.statePath, 'intents', `${intentDigest(receipt.workId, receipt.cwd)}.json`); + const intent = await readPrivateJsonRecord(intentPath); + if (intent?.receiptPath === receiptPath) { + await atomicPrivateWrite(intentPath, { schemaVersion: 1, receiptPath: archivePath, state: receipt.state }); + } + } +} + +async function reconcileReceiptsLocked(config, options) { + const directory = path.join(config.statePath, 'dispatch-receipts'); + const entries = await readPrivateJsonDirectoryEntries(directory, { onWarning: options.onWarning }); + const errors = []; + let reconciled = 0; + let archived = 0; + const observed = new Map(); + for (const entry of entries) { + const result = await reconcileReceipt(config, entry, options); + observed.set(entry.name, result.receipt); + if (result.changed) reconciled += 1; + if (result.error) errors.push({ attemptId: entry.value.attemptId ?? null, code: result.error }); + } + for (const entry of entries) { + const receipt = observed.get(entry.name) ?? entry.value; + const timestamp = terminalTimestamp(receipt); + if (!TERMINAL_RECEIPT_STATES.has(receipt.state) || timestamp === null + || options.nowMs - timestamp < options.archiveAfterMs) continue; + await archiveReceipt(config, entry, receipt); + archived += 1; + } + return { scanned: entries.length, reconciled, archived, errors }; +} + +export async function reconcileReceipts(configInput, options = {}) { + const config = validateInstallConfig(configInput, { expectedBorgHome: configInput.borgHome }); + const nowMs = options.nowMs ?? Date.now(); + if (!Number.isFinite(nowMs)) throw new Error('RECONCILE_TIME_INVALID'); + const archiveAfterMs = options.archiveAfterMs ?? config.routing.archiveAfterMs; + if (!Number.isInteger(archiveAfterMs) || archiveAfterMs < 1_000) throw new Error('ARCHIVE_AGE_INVALID'); + const settings = { + nowMs, + archiveAfterMs, + stageTimeoutMs: options.stageTimeoutMs ?? config.routing.timeoutMs, + conductorProvider: options.conductorProvider + ?? ((lane, operation) => nativeConductorProvider(lane, operation, config.routing.timeoutMs)), + onWarning: options.onWarning, + }; + if (options.lockHeld === true) return reconcileReceiptsLocked(config, settings); + const release = await acquireDispatchLock(config.statePath, config.routing.lockTimeoutMs, { nowMs }); + try { + return await reconcileReceiptsLocked(config, settings); + } finally { + await release(); + } +} + function noEligibleMessage(result) { const evidence = result.candidates.map((candidate) => `${candidate.laneId}:${candidate.issues.join('+')}`).join(','); return `no eligible conductor${evidence ? `: ${evidence}` : ''}`; @@ -525,7 +923,7 @@ function noEligibleMessage(result) { export async function rank(configInput, options = {}) { const config = validateInstallConfig(configInput, { expectedBorgHome: configInput.borgHome }); const nowMs = options.nowMs ?? Date.now(); - return snapshot(config, { + const settings = { nowMs, usageMaxAgeMs: options.usageMaxAgeMs ?? 15_000, stageTimeoutMs: options.stageTimeoutMs ?? config.routing.timeoutMs, @@ -536,7 +934,15 @@ export async function rank(configInput, options = {}) { claimsProvider: options.claimsProvider ?? observeClaims, conductorProvider: options.conductorProvider ?? ((lane, operation) => nativeConductorProvider(lane, operation, config.routing.timeoutMs)), + }; + await reconcileReceipts(config, { + nowMs, + archiveAfterMs: options.archiveAfterMs ?? config.routing.archiveAfterMs, + stageTimeoutMs: settings.stageTimeoutMs, + conductorProvider: settings.conductorProvider, + onWarning: options.onWarning, }); + return snapshot(config, settings); } // Work IDs are reserved across every claim source and machine. Workspace @@ -577,16 +983,23 @@ export async function inspectDispatch(configInput, options = {}) { 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') + let receiptPath = intent.receiptPath; + const receiptsDirectory = path.join(config.statePath, 'dispatch-receipts'); + const allowedDirectories = new Set([receiptsDirectory, path.join(receiptsDirectory, 'archive')]); + if (typeof receiptPath !== 'string' || !allowedDirectories.has(path.dirname(receiptPath)) || !path.basename(receiptPath).endsWith('.json')) throw new Error('UNSAFE_RECEIPT_REFERENCE'); - const receipt = await readPrivateJsonRecord(receiptPath); + let receipt = await readPrivateJsonRecord(receiptPath); + if (receipt === null && path.dirname(receiptPath) === receiptsDirectory) { + const archivedPath = path.join(receiptsDirectory, 'archive', path.basename(receiptPath)); + receipt = await readPrivateJsonRecord(archivedPath); + if (receipt !== null) receiptPath = archivedPath; + } 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 }; + completionVerified: receipt.state === 'COMPLETED' && receipt.nativeEvidence?.method === 'thread/read' }; } export async function dispatch(configInput, options = {}) { @@ -603,7 +1016,9 @@ export async function dispatch(configInput, options = {}) { 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); + const release = await acquireDispatchLock(config.statePath, config.routing.lockTimeoutMs, { + nowMs: options.nowMs ?? Date.now(), + }); let receipt; let paths; let nativeAttempted = false; const clock = () => new Date(options.nowMs ?? Date.now()).toISOString(); const progress = async (phase, changes = {}) => { @@ -627,12 +1042,23 @@ export async function dispatch(configInput, options = {}) { 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. + // Persist the work intent BEFORE reconciliation and every admission scan. + // If either refuses, the catch path records PRE_START_FAILED evidence. paths = await createReceipt(config, receipt, lexicalCwd); + // Retire native terminal work before any admission or claim collector sees + // the active ledger. Failed readback remains an active claim. + await reconcileReceipts(config, { + lockHeld: true, + nowMs: options.nowMs ?? Date.now(), + archiveAfterMs: options.archiveAfterMs ?? config.routing.archiveAfterMs, + stageTimeoutMs, + conductorProvider: common.conductorProvider, + onWarning: options.onWarning, + }); await progress('ADMISSION'); - const preliminary = await rank(config, common); + const preliminary = await snapshot(config, { ...common, nowMs: options.nowMs ?? Date.now() }); await progress('RECHECK'); - const final = await rank(config, { ...common, nowMs: options.recheckNowMs ?? options.nowMs ?? Date.now() }); + const final = await snapshot(config, { ...common, nowMs: options.recheckNowMs ?? options.nowMs ?? Date.now() }); const selected = final.candidates.find((candidate) => candidate.eligible); if (!selected) throw new Error(noEligibleMessage(final)); // This router's own submitted and uncertain receipts are claims whichever diff --git a/conductor/tests/config.test.mjs b/conductor/tests/config.test.mjs index d2eb6fe..8490597 100644 --- a/conductor/tests/config.test.mjs +++ b/conductor/tests/config.test.mjs @@ -41,6 +41,22 @@ test('default config contains one owner-local machine and one dedicated unauthen assert.equal(config.conductors[0].accountPin, null); assert.equal(config.requiredVersions.node, '24.21.0'); assert.equal(config.requiredVersions.codex, '0.146.0'); + assert.equal(config.routing.archiveAfterMs, 7 * 24 * 60 * 60 * 1000); +}); + +test('terminal receipt archive age is configurable and bounded', () => { + const config = buildDefaultConfig(BORG_HOME, '/opt/codex/bin/codex', { nodeBin: '/opt/node/bin/node' }); + const custom = validateInstallConfig({ + ...config, + routing: { ...config.routing, archiveAfterMs: 24 * 60 * 60 * 1000 }, + }); + assert.equal(custom.routing.archiveAfterMs, 24 * 60 * 60 * 1000); + for (const archiveAfterMs of [0, -1, 999, 366 * 24 * 60 * 60 * 1000, '86400000']) { + assert.throws(() => validateInstallConfig({ + ...config, + routing: { ...config.routing, archiveAfterMs }, + }), /archiveAfterMs/); + } }); test('install config rejects noncanonical homes, relative runtime paths, remote listeners, and credentials', () => { diff --git a/conductor/tests/http-auth.test.mjs b/conductor/tests/http-auth.test.mjs new file mode 100644 index 0000000..6391dc4 --- /dev/null +++ b/conductor/tests/http-auth.test.mjs @@ -0,0 +1,192 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import http from 'node:http'; +import net from 'node:net'; +import os from 'node:os'; +import path from 'node:path'; +import test from 'node:test'; + +import * as conductorApi from '../conductor.mjs'; + +async function unusedLoopbackPort() { + const listener = net.createServer(); + await new Promise((resolve, reject) => { + listener.once('error', reject); + listener.listen(0, '127.0.0.1', resolve); + }); + const { port } = listener.address(); + await new Promise((resolve, reject) => listener.close((error) => (error ? reject(error) : resolve()))); + return port; +} + +function request(port, pathname, options = {}) { + return new Promise((resolve, reject) => { + const body = options.body === undefined + ? null + : (typeof options.body === 'string' ? options.body : JSON.stringify(options.body)); + const headers = { + host: `127.0.0.1:${port}`, + ...(options.headers || {}), + ...(body === null ? {} : { 'content-length': Buffer.byteLength(body) }), + }; + const req = http.request({ + host: '127.0.0.1', port, path: pathname, method: options.method || 'GET', headers, + }, (res) => { + const chunks = []; + res.on('data', (chunk) => chunks.push(chunk)); + res.on('end', () => { + const text = Buffer.concat(chunks).toString('utf8'); + resolve({ status: res.statusCode, body: text ? JSON.parse(text) : null }); + }); + }); + req.once('error', reject); + if (body !== null) req.end(body); + else req.end(); + }); +} + +async function conductorFixture(t, authMode) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), `borg-http-${authMode}-`)); + const fakeCodex = path.join(root, 'fake-codex.mjs'); + fs.writeFileSync(fakeCodex, `#!${process.execPath} +let input = ''; +process.stdin.setEncoding('utf8'); +process.stdin.on('data', (chunk) => { + input += chunk; + let newline; + while ((newline = input.indexOf('\\n')) >= 0) { + const line = input.slice(0, newline); + input = input.slice(newline + 1); + if (!line.trim()) continue; + const message = JSON.parse(line); + if (message.id === undefined) continue; + const result = message.method === 'initialize' + ? { userAgent: 'fake-codex', platformFamily: 'unix', platformOs: 'test' } + : { method: message.method, account: { type: 'chatgpt', email: 'synthetic-owner', planType: 'pro' } }; + process.stdout.write(JSON.stringify({ jsonrpc: '2.0', id: message.id, result }) + '\\n'); + } +}); +`, { mode: 0o700 }); + const port = await unusedLoopbackPort(); + const codexHome = path.join(root, 'profile'); + let listeningResolve; + const listening = new Promise((resolve) => { listeningResolve = resolve; }); + const conductor = conductorApi.startConductor({ + port, + authMode, + codexBin: fakeCodex, + codexHome, + borgHome: path.join(root, 'borg'), + logsPath: path.join(root, 'logs'), + policies: { seat: 'owner seat rules', lead: 'owner lead rules' }, + installSignalHandlers: false, + exitOnChildExit: false, + onListening: listeningResolve, + }); + t.after(async () => { + const closed = conductor.server.listening + ? new Promise((resolve) => conductor.server.once('close', resolve)) + : Promise.resolve(); + conductor.close(); + await closed; + fs.rmSync(root, { recursive: true, force: true }); + }); + await listening; + const tokenPath = path.join(codexHome, '.conductor', 'http-token'); + assert.equal(fs.existsSync(tokenPath), true); + const token = fs.readFileSync(tokenPath, 'utf8').trim(); + assert.match(token, /^[a-f0-9]{64}$/); + assert.equal(fs.statSync(tokenPath).mode & 0o777, 0o600); + assert.equal(fs.statSync(path.dirname(tokenPath)).mode & 0o777, 0o700); + return { root, port, token, tokenPath }; +} + +function bearer(token) { + return { authorization: `Bearer ${token}` }; +} + +test('fresh conductors enforce bearer, JSON, RPC, host, and browser boundaries by default', async (t) => { + const f = await conductorFixture(t, undefined); + const health = await request(f.port, '/healthz'); + assert.deepEqual(health, { status: 200, body: { ok: true } }); + + assert.equal((await request(f.port, '/status')).status, 401); + assert.equal((await request(f.port, '/status', { headers: bearer('0'.repeat(64)) })).status, 401); + const status = await request(f.port, '/status', { headers: bearer(f.token) }); + assert.equal(status.status, 200); + assert.equal(status.body.auth.mode, 'enforce'); + assert.equal(status.body.auth.unauthenticatedCount, 2); + assert.ok(status.body.auth.lastUnauthenticatedAt); + + assert.equal((await request(f.port, '/status', { + headers: { ...bearer(f.token), host: 'attacker.example' }, + })).status, 403); + assert.equal((await request(f.port, '/status', { + headers: { ...bearer(f.token), origin: 'https://attacker.example' }, + })).status, 403); + assert.equal((await request(f.port, '/status', { + headers: { ...bearer(f.token), 'sec-fetch-site': 'cross-site' }, + })).status, 403); + assert.equal((await request(f.port, '/status', { + headers: { ...bearer(f.token), 'sec-fetch-site': 'none' }, + })).status, 200); + + const noJson = await request(f.port, '/rpc', { + method: 'POST', headers: bearer(f.token), body: { method: 'account/read', params: {} }, + }); + assert.equal(noJson.status, 415); + const disallowed = await request(f.port, '/rpc', { + method: 'POST', headers: { ...bearer(f.token), 'content-type': 'application/json' }, + body: { method: 'dangerous/write', params: {} }, + }); + assert.equal(disallowed.status, 403); + const allowed = await request(f.port, '/rpc', { + method: 'POST', headers: { ...bearer(f.token), 'content-type': 'application/json; charset=utf-8' }, + body: { method: 'account/read', params: {} }, + }); + assert.equal(allowed.status, 200); + assert.equal(allowed.body.method, 'account/read'); + + const counters = await request(f.port, '/status', { headers: bearer(f.token) }); + assert.equal(counters.body.auth.disallowedRpcCount, 1); + assert.deepEqual([...conductorApi.RPC_METHOD_ALLOWLIST].sort(), [ + 'account/rateLimits/read', 'account/read', 'config/read', 'hooks/list', 'model/list', 'thread/read', + ]); +}); + +test('report mode serves legacy callers, counts violations, and rate-limits audit lines', async (t) => { + const f = await conductorFixture(t, 'report'); + assert.equal((await request(f.port, '/healthz')).status, 200); + assert.equal((await request(f.port, '/status')).status, 200); + assert.equal((await request(f.port, '/status', { headers: bearer(f.token) })).status, 200); + assert.equal((await request(f.port, '/status', { + headers: { host: 'attacker.example' }, + })).status, 403, 'host checks stay enforced in report mode without a token'); + assert.equal((await request(f.port, '/status', { + headers: { ...bearer(f.token), origin: 'https://attacker.example' }, + })).status, 403, 'browser checks stay enforced in report mode with a token'); + + for (let index = 0; index < 2; index += 1) { + const response = await request(f.port, '/rpc', { + method: 'POST', + body: { method: 'dangerous/write', params: {} }, + }); + assert.equal(response.status, 200); + } + const status = await request(f.port, '/status', { headers: bearer(f.token) }); + assert.equal(status.body.auth.mode, 'report'); + assert.equal(status.body.auth.disallowedRpcCount, 2); + assert.equal(status.body.auth.unauthenticatedCount, 3); + + await new Promise((resolve) => setImmediate(resolve)); + const auditFile = fs.readdirSync(path.join(f.root, 'logs')) + .map((name) => path.join(f.root, 'logs', name)) + .find((file) => path.basename(file).startsWith('http-auth-')); + assert.ok(auditFile); + const lines = fs.readFileSync(auditFile, 'utf8').trim().split('\n').map(JSON.parse); + assert.equal(lines.filter((entry) => entry.path === '/rpc' && entry.reason === 'missing_token').length, 1); + assert.equal(lines.filter((entry) => entry.path === '/rpc' && entry.reason === 'invalid_content_type').length, 1); + assert.equal(lines.filter((entry) => entry.path === '/rpc' && entry.reason === 'disallowed_rpc_method').length, 1); + assert.ok(lines.every((entry) => Object.keys(entry).sort().join(',') === 'method,path,reason,ts,userAgent')); + assert.equal(fs.readFileSync(auditFile, 'utf8').includes(f.token), false); +}); diff --git a/conductor/tests/native-app-server.test.mjs b/conductor/tests/native-app-server.test.mjs index b55a5a9..2e799f7 100644 --- a/conductor/tests/native-app-server.test.mjs +++ b/conductor/tests/native-app-server.test.mjs @@ -24,13 +24,15 @@ async function unusedLoopbackPort() { return port; } -async function awaitNativeStatus(port, child, diagnostics, deadlineMs = Date.now() + 20_000) { +async function awaitNativeStatus(port, child, diagnostics, tokenPath, deadlineMs = Date.now() + 20_000) { while (Date.now() < deadlineMs) { if (child.exitCode !== null) { throw new Error(`conductor exited ${child.exitCode}: ${diagnostics.join('')}`); } try { + const token = fs.existsSync(tokenPath) ? fs.readFileSync(tokenPath, 'utf8').trim() : null; const response = await fetch(`http://127.0.0.1:${port}/status`, { + headers: token ? { authorization: `Bearer ${token}` } : {}, signal: AbortSignal.timeout(1_000), }); if (response.ok) return response.json(); @@ -59,6 +61,7 @@ test('real Codex app-server initializes with a separate unauthenticated profile' codexBin: path.resolve(CODEX_BIN), }); const profile = path.join(borgHome, 'conductors/primary/profile'); + const tokenPath = path.join(profile, '.conductor/http-token'); assert.equal(fs.existsSync(path.join(profile, 'auth.json')), false); const child = spawn(process.execPath, [cliPath, 'start', '--config', 'conductors/config.json', '--lane', 'primary'], { @@ -77,7 +80,7 @@ test('real Codex app-server initializes with a separate unauthenticated profile' fs.rmSync(parent, { recursive: true, force: true }); }); - const status = await awaitNativeStatus(port, child, diagnostics); + const status = await awaitNativeStatus(port, child, diagnostics, tokenPath); assert.equal(status.ok, true); assert.equal(status.port, port); assert.equal(status.codexHome, profile); diff --git a/conductor/tests/receipt-reader.test.mjs b/conductor/tests/receipt-reader.test.mjs index 65e8255..b0033d3 100644 --- a/conductor/tests/receipt-reader.test.mjs +++ b/conductor/tests/receipt-reader.test.mjs @@ -54,3 +54,21 @@ test('invalid timeout configuration is rejected without starting a reader', asyn await assert.rejects(readPrivateJsonDirectory(root, { timeoutMs }), /RECEIPT_READ_TIMEOUT_INVALID/); } }); + +test('directory reader warns at eighty percent of both scan caps', async (t) => { + const entriesRoot = fixture(t); + for (let index = 0; index < 8000; index += 1) { + fs.writeFileSync(path.join(entriesRoot, `entry-${index}.txt`), '', { mode: 0o600 }); + } + const entryWarnings = []; + await readPrivateJsonDirectory(entriesRoot, { onWarning: (warning) => entryWarnings.push(warning) }); + assert.ok(entryWarnings.some((warning) => warning.code === 'RECEIPT_ENTRY_LIMIT_WARNING')); + + const bytesRoot = fixture(t); + for (let index = 0; index < 54; index += 1) { + fs.writeFileSync(path.join(bytesRoot, `receipt-${index}.json`), JSON.stringify({ text: 'x'.repeat(62_900) }), { mode: 0o600 }); + } + const byteWarnings = []; + await readPrivateJsonDirectory(bytesRoot, { onWarning: (warning) => byteWarnings.push(warning) }); + assert.ok(byteWarnings.some((warning) => warning.code === 'RECEIPT_TOTAL_LIMIT_WARNING')); +}); diff --git a/conductor/tests/resume-bookkeeping.test.mjs b/conductor/tests/resume-bookkeeping.test.mjs index 61d8eb5..875a595 100644 --- a/conductor/tests/resume-bookkeeping.test.mjs +++ b/conductor/tests/resume-bookkeeping.test.mjs @@ -105,6 +105,9 @@ process.stdin.on('data', (chunk) => { }); await listening; + const token = fs.readFileSync(path.join(root, 'profile', '.conductor', 'http-token'), 'utf8').trim(); + const authHeaders = { 'content-type': 'application/json', authorization: `Bearer ${token}` }; + const predecessor = { laneId: 'prior-lane', threadId: 'thread-relocated', @@ -112,7 +115,7 @@ process.stdin.on('data', (chunk) => { }; const resumed = await fetch(`http://127.0.0.1:${port}/thread/resume`, { method: 'POST', - headers: { 'content-type': 'application/json' }, + headers: authHeaders, body: JSON.stringify({ threadId: 'thread-relocated', cwd: effectiveCwd, @@ -125,7 +128,9 @@ process.stdin.on('data', (chunk) => { assert.equal(resumed.status, 200); assert.equal((await resumed.json()).cwd, effectiveCwd); - const status = await (await fetch(`http://127.0.0.1:${port}/status`)).json(); + const status = await (await fetch(`http://127.0.0.1:${port}/status`, { + headers: { authorization: `Bearer ${token}` }, + })).json(); const record = status.threads['thread-relocated']; assert.equal(record.cwd, effectiveCwd, 'top-level native response cwd wins over historical thread cwd'); assert.equal(record.role, 'lead'); @@ -141,11 +146,13 @@ process.stdin.on('data', (chunk) => { const resumedAgain = await fetch(`http://127.0.0.1:${port}/thread/resume`, { method: 'POST', - headers: { 'content-type': 'application/json' }, + headers: authHeaders, body: JSON.stringify({ threadId: 'thread-relocated', cwd: effectiveCwd }), }); assert.equal(resumedAgain.status, 200); - const nextStatus = await (await fetch(`http://127.0.0.1:${port}/status`)).json(); + const nextStatus = await (await fetch(`http://127.0.0.1:${port}/status`, { + headers: { authorization: `Bearer ${token}` }, + })).json(); const nextRecord = nextStatus.threads['thread-relocated']; assert.equal(nextRecord.startedAt, record.startedAt); assert.equal(nextRecord.role, 'lead', 'an omitted role preserves the attached thread role'); diff --git a/conductor/tests/router-lifecycle.test.mjs b/conductor/tests/router-lifecycle.test.mjs new file mode 100644 index 0000000..2c865db --- /dev/null +++ b/conductor/tests/router-lifecycle.test.mjs @@ -0,0 +1,403 @@ +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'; + +const NOW = Date.parse('2026-09-26T12:00:00.000Z'); +const account = { type: 'chatgpt', email: 'synthetic-owner', planType: 'pro' }; + +function fixture(t) { + const root = fs.mkdtempSync(path.join(os.tmpdir(), 'borg-router-lifecycle-')); + 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); + fs.mkdirSync(path.join(config.statePath, 'dispatch-receipts'), { recursive: true, mode: 0o700 }); + return { root, config }; +} + +function putReceipt(config, name, changes = {}) { + const receipt = { + schemaVersion: 1, + attemptId: name, + workId: `work-${name}`, + cwd: path.dirname(config.statePath), + laneId: 'primary', + machineId: 'local', + state: 'DISPATCHED', + phase: 'TURN_STARTED', + attemptedAt: '2026-09-26T10:00:00.000Z', + dispatchedAt: '2026-09-26T10:00:01.000Z', + threadId: `thread-${name}`, + turnId: `turn-${name}`, + nativeStartAttempted: true, + events: [], + ...changes, + }; + const receiptPath = path.join(config.statePath, 'dispatch-receipts', `${name}.json`); + fs.writeFileSync(receiptPath, `${JSON.stringify(receipt, null, 2)}\n`, { mode: 0o600 }); + return { receipt, receiptPath }; +} + +function readJson(file) { + return JSON.parse(fs.readFileSync(file, 'utf8')); +} + +function intentDigest(workId, cwd) { + return crypto.createHash('sha256').update(JSON.stringify({ workId, cwd })).digest('hex'); +} + +test('reconcile writes terminal receipt states only from exact native turn readback', async (t) => { + const { config } = fixture(t); + const rows = [ + ['complete', 'completed', 'COMPLETED'], + ['failure', { type: 'failed' }, 'FAILED'], + ['cancel', 'interrupted', 'CANCELLED'], + ['running', 'inProgress', 'DISPATCHED'], + ]; + const files = new Map(rows.map(([name]) => [name, putReceipt(config, name).receiptPath])); + const operations = []; + const result = await router.reconcileReceipts(config, { + nowMs: NOW, + archiveAfterMs: 30 * 24 * 60 * 60 * 1000, + conductorProvider: async (_lane, operation) => { + operations.push(operation); + assert.equal(operation.kind, 'rpc'); + assert.equal(operation.method, 'thread/read'); + assert.equal(operation.params.includeTurns, true); + const name = operation.params.threadId.slice('thread-'.length); + const status = rows.find(([candidate]) => candidate === name)[1]; + return { thread: { + id: operation.params.threadId, + turns: [{ id: `turn-${name}`, status, completedAt: '2026-09-26T11:00:00.000Z' }], + } }; + }, + }); + + assert.equal(result.reconciled, 3); + assert.equal(result.archived, 0); + assert.equal(operations.length, 4); + for (const [name, _native, expected] of rows) { + const receipt = readJson(files.get(name)); + assert.equal(receipt.state, expected); + if (expected !== 'DISPATCHED') { + assert.equal(receipt.reconciledAt, new Date(NOW).toISOString()); + assert.equal(receipt.terminalAt, '2026-09-26T11:00:00.000Z'); + assert.equal(receipt.nativeEvidence.method, 'thread/read'); + assert.equal(receipt.nativeEvidence.threadId, `thread-${name}`); + assert.equal(receipt.nativeEvidence.turnId, `turn-${name}`); + } + } +}); + +test('reconcile keeps claims active when native evidence is missing, mismatched, or unreadable', async (t) => { + const { config } = fixture(t); + const missing = putReceipt(config, 'missing'); + const mismatch = putReceipt(config, 'mismatch'); + const unreadable = putReceipt(config, 'unreadable'); + const result = await router.reconcileReceipts(config, { + nowMs: NOW, + conductorProvider: async (_lane, operation) => { + if (operation.params.threadId === 'thread-missing') return { thread: { id: 'thread-missing', turns: [] } }; + if (operation.params.threadId === 'thread-mismatch') { + return { thread: { id: 'different-thread', turns: [{ id: 'turn-mismatch', status: 'completed' }] } }; + } + throw new Error('native unavailable'); + }, + }); + assert.equal(result.reconciled, 0); + assert.equal(result.errors.length, 3); + for (const item of [missing, mismatch, unreadable]) { + assert.equal(readJson(item.receiptPath).state, 'DISPATCHED'); + } +}); + +test('dispatch reconciles native terminal work before the first claims scan', async (t) => { + const { root, config } = fixture(t); + const old = putReceipt(config, 'old', { cwd: root }); + const calls = []; + const conductorProvider = async (lane, operation) => { + if (operation.method === 'thread/read') { + calls.push('thread/read'); + return { thread: { id: 'thread-old', turns: [{ id: 'turn-old', status: 'completed' }] } }; + } + if (operation.kind === 'status') { + calls.push('status'); + return { ok: true, port: lane.port, codexHome: lane.codexHome, supportedRoles: ['leaf'], threads: {} }; + } + if (operation.method === 'account/read') return { account }; + if (operation.method === 'account/rateLimits/read') { + return { rateLimits: { limitId: 'codex', primary: { + usedPercent: 5, windowDurationMins: 300, resetsAt: '2026-09-27T12:00:00.000Z', + } } }; + } + if (operation.kind === 'thread-start') return { threadId: 'thread-new' }; + if (operation.kind === 'turn-start') return { turnId: 'turn-new' }; + throw new Error(`unexpected operation: ${operation.kind}`); + }; + const result = await router.dispatch(config, { + cwd: root, + workId: 'new-work', + prompt: 'synthetic dispatch', + nowMs: NOW, + capacityProvider: async () => ({ + observedAt: new Date(NOW).toISOString(), reachable: true, loadPerCore: 0.1, memoryUsePercent: 20, + }), + claimsProvider: async () => { + calls.push('claims'); + assert.equal(readJson(old.receiptPath).state, 'COMPLETED'); + return { observedAt: new Date(NOW).toISOString(), active: [] }; + }, + conductorProvider, + }); + assert.equal(result.receipt.state, 'DISPATCHED'); + assert.equal(calls[0], 'thread/read'); + assert.ok(calls.indexOf('thread/read') < calls.indexOf('claims')); +}); + +test('standalone rank also reconciles before its claims scan', async (t) => { + const { config } = fixture(t); + const old = putReceipt(config, 'rank-old'); + const calls = []; + await router.rank(config, { + nowMs: NOW, + capacityProvider: async () => ({ + observedAt: new Date(NOW).toISOString(), reachable: true, loadPerCore: 0.1, memoryUsePercent: 20, + }), + claimsProvider: async () => { + calls.push('claims'); + assert.equal(readJson(old.receiptPath).state, 'COMPLETED'); + return { observedAt: new Date(NOW).toISOString(), active: [] }; + }, + conductorProvider: async (lane, operation) => { + if (operation.method === 'thread/read') { + calls.push('thread/read'); + return { thread: { id: 'thread-rank-old', turns: [{ id: 'turn-rank-old', status: 'completed' }] } }; + } + if (operation.kind === 'status') { + return { ok: true, port: lane.port, codexHome: lane.codexHome, supportedRoles: ['leaf'], threads: {} }; + } + if (operation.method === 'account/read') return { account }; + if (operation.method === 'account/rateLimits/read') { + return { rateLimits: { limitId: 'codex', primary: { + usedPercent: 5, windowDurationMins: 300, resetsAt: '2026-09-27T12:00:00.000Z', + } } }; + } + throw new Error(`unexpected operation: ${operation.kind}`); + }, + }); + assert.equal(calls[0], 'thread/read'); + assert.ok(calls.indexOf('thread/read') < calls.indexOf('claims')); +}); + +test('aged terminal receipts move to archive and their durable intent follows them', async (t) => { + const { config } = fixture(t); + const { receipt, receiptPath } = putReceipt(config, 'aged', { + state: 'COMPLETED', + terminalAt: '2026-09-01T00:00:00.000Z', + nativeEvidence: { method: 'thread/read', threadId: 'thread-aged', turnId: 'turn-aged' }, + }); + const intents = path.join(config.statePath, 'intents'); + fs.mkdirSync(intents, { recursive: true, mode: 0o700 }); + const intentPath = path.join(intents, `${intentDigest(receipt.workId, receipt.cwd)}.json`); + fs.writeFileSync(intentPath, `${JSON.stringify({ schemaVersion: 1, receiptPath, state: receipt.state })}\n`, { mode: 0o600 }); + + const result = await router.reconcileReceipts(config, { + nowMs: NOW, + archiveAfterMs: 24 * 60 * 60 * 1000, + conductorProvider: async () => { throw new Error('terminal receipts need no native reread'); }, + }); + const archived = path.join(config.statePath, 'dispatch-receipts', 'archive', 'aged.json'); + assert.equal(result.archived, 1); + assert.equal(fs.existsSync(receiptPath), false); + assert.equal(fs.existsSync(archived), true); + assert.equal(readJson(intentPath).receiptPath, archived); + const inspected = await router.inspectDispatch(config, { cwd: receipt.cwd, workId: receipt.workId }); + assert.equal(inspected.receipt.state, 'COMPLETED'); + assert.equal(inspected.completionVerified, true); +}); + +test('dispatch lock preserves live and young empty locks and recovers only stale owners', async (t) => { + const { config } = fixture(t); + const lockPath = path.join(config.statePath, 'dispatch.lock'); + fs.mkdirSync(config.statePath, { recursive: true, mode: 0o700 }); + fs.writeFileSync(lockPath, `${process.pid}\n`, { mode: 0o600 }); + const liveOld = new Date(NOW - 11 * 60 * 1000); + fs.utimesSync(lockPath, liveOld, liveOld); + await assert.rejects(router.acquireDispatchLock(config.statePath, 20, { + nowMs: NOW, processAlive: () => true, + }), /dispatch lock busy/); + assert.equal(fs.readFileSync(lockPath, 'utf8'), `${process.pid}\n`); + assert.deepEqual(fs.readdirSync(config.statePath).filter((name) => name.startsWith('dispatch.lock.stale-')), []); + + fs.writeFileSync(lockPath, '', { mode: 0o600 }); + fs.utimesSync(lockPath, new Date(NOW), new Date(NOW)); + await assert.rejects(router.acquireDispatchLock(config.statePath, 20, { + nowMs: NOW, processAlive: () => false, + }), /dispatch lock busy/); + assert.equal(fs.readFileSync(lockPath, 'utf8'), ''); + assert.deepEqual(fs.readdirSync(config.statePath).filter((name) => name.startsWith('dispatch.lock.stale-')), []); + + for (const [name, body, aged, processAlive] of [ + ['dead', '999999999\n', false, () => false], + ['aged', 'not-a-pid\n', true, () => false], + ['zero', '', true, () => false], + ]) { + fs.rmSync(lockPath, { force: true }); + fs.writeFileSync(lockPath, body, { mode: 0o600 }); + if (aged) { + const old = new Date(NOW - 11 * 60 * 1000); + fs.utimesSync(lockPath, old, old); + } + const release = await router.acquireDispatchLock(config.statePath, 50, { + nowMs: NOW, processAlive, + }); + const stale = fs.readdirSync(config.statePath) + .filter((entry) => entry.startsWith('dispatch.lock.stale-')); + assert.ok(stale.length >= 1, `${name} lock must be retained under a stale name`); + assert.equal(fs.existsSync(lockPath), true); + await release(); + assert.equal(fs.existsSync(lockPath), false); + } +}); + +test('overlapping lock acquirers cannot both win while the first lock body is empty', async (t) => { + const { config } = fixture(t); + const lockPath = path.join(config.statePath, 'dispatch.lock'); + fs.mkdirSync(config.statePath, { recursive: true, mode: 0o700 }); + const originalOpen = fs.promises.open; + let intercepted = false; + let announceOpen; + let resumeWrite; + const opened = new Promise((resolve) => { announceOpen = resolve; }); + const resume = new Promise((resolve) => { resumeWrite = resolve; }); + fs.promises.open = async (...args) => { + const handle = await originalOpen(...args); + if (!intercepted && args[0] === lockPath && args[1] === 'wx') { + intercepted = true; + const originalWrite = handle.writeFile.bind(handle); + handle.writeFile = async (...writeArgs) => { + announceOpen(); + await resume; + return originalWrite(...writeArgs); + }; + } + return handle; + }; + t.after(() => { + fs.promises.open = originalOpen; + resumeWrite(); + }); + + const settle = (promise) => promise.then( + (release) => ({ acquired: true, release }), + (error) => ({ acquired: false, error }), + ); + const first = settle(router.acquireDispatchLock(config.statePath, 250, { nowMs: NOW })); + await opened; + fs.utimesSync(lockPath, new Date(NOW), new Date(NOW)); + const second = settle(router.acquireDispatchLock(config.statePath, 100, { nowMs: NOW })); + await new Promise((resolve) => setTimeout(resolve, 25)); + resumeWrite(); + const outcomes = await Promise.all([first, second]); + try { + assert.equal(outcomes.filter((outcome) => outcome.acquired).length, 1); + assert.match(outcomes.find((outcome) => !outcome.acquired).error.message, /dispatch lock busy/); + } finally { + await Promise.allSettled(outcomes.filter((outcome) => outcome.acquired) + .map((outcome) => outcome.release())); + } +}); + +test('dispatch lock recovers a recycled live pid with a different process start time', async (t) => { + const { config } = fixture(t); + const lockPath = path.join(config.statePath, 'dispatch.lock'); + fs.mkdirSync(config.statePath, { recursive: true, mode: 0o700 }); + fs.writeFileSync(lockPath, `${JSON.stringify({ + schemaVersion: 1, + lockId: 'previous-process', + pid: process.pid, + processStartTime: 'old-process-start', + acquiredAt: new Date(NOW).toISOString(), + })}\n`, { mode: 0o600 }); + let release; + await assert.doesNotReject(async () => { + release = await router.acquireDispatchLock(config.statePath, 100, { + nowMs: NOW, + processAlive: () => true, + processStartTime: async () => 'current-process-start', + }); + }); + const stale = fs.readdirSync(config.statePath) + .find((entry) => entry.startsWith('dispatch.lock.stale-')); + assert.ok(stale); + assert.equal(readJson(path.join(config.statePath, stale)).lockId, 'previous-process'); + assert.equal(readJson(lockPath).processStartTime, 'current-process-start'); + await release(); +}); + +test('stale recovery never retires a replacement lock inode', async (t) => { + const { config } = fixture(t); + const lockPath = path.join(config.statePath, 'dispatch.lock'); + fs.mkdirSync(config.statePath, { recursive: true, mode: 0o700 }); + fs.writeFileSync(lockPath, `${JSON.stringify({ pid: 999999999, processStartTime: 'dead-start' })}\n`, { + mode: 0o600, + }); + const originalLink = fs.promises.link; + let swapped = false; + fs.promises.link = async (source, destination) => { + if (!swapped && source === lockPath) { + swapped = true; + await fs.promises.rename(lockPath, `${lockPath}.before-swap`); + await fs.promises.writeFile(lockPath, `${JSON.stringify({ + schemaVersion: 1, + lockId: 'replacement', + pid: process.pid, + processStartTime: 'current-process-start', + })}\n`, { mode: 0o600, flag: 'wx' }); + } + return originalLink(source, destination); + }; + t.after(() => { fs.promises.link = originalLink; }); + + await assert.rejects(router.acquireDispatchLock(config.statePath, 70, { + nowMs: NOW, + processAlive: (pid) => pid === process.pid, + processStartTime: async () => 'current-process-start', + }), /dispatch lock busy/); + assert.equal(readJson(lockPath).lockId, 'replacement'); + assert.deepEqual(fs.readdirSync(config.statePath).filter((entry) => ( + entry.startsWith('dispatch.lock.stale-') + )), []); +}); + +test('native conductor calls send each lane profile token and tolerate a missing token', async (t) => { + const { config } = fixture(t); + const lane = config.conductors[0]; + const token = 'b'.repeat(64); + const tokenDirectory = path.join(lane.codexHome, '.conductor'); + fs.mkdirSync(tokenDirectory, { recursive: true, mode: 0o700 }); + fs.writeFileSync(path.join(tokenDirectory, 'http-token'), `${token}\n`, { mode: 0o600 }); + const calls = []; + const fetchImpl = async (url, init) => { + calls.push({ url, init }); + return { ok: true, async json() { return url.includes('/threads') ? { data: [] } : { ok: true }; } }; + }; + await router.nativeConductorProvider(lane, { kind: 'status' }, 5000, fetchImpl); + await router.nativeConductorProvider(lane, { + kind: 'rpc', method: 'thread/read', params: { threadId: 'thread-1', includeTurns: true }, + }, 5000, fetchImpl); + assert.equal(calls.length, 3); + assert.ok(calls.every((call) => call.init.headers.authorization === `Bearer ${token}`)); + + fs.unlinkSync(path.join(tokenDirectory, 'http-token')); + calls.length = 0; + await router.nativeConductorProvider(lane, { kind: 'status' }, 5000, fetchImpl); + assert.ok(calls.every((call) => !call.init.headers?.authorization)); +}); diff --git a/graph/tests/test_legacy_graphiti_scope_guard.py b/graph/tests/test_legacy_graphiti_scope_guard.py index a894b90..739c62b 100644 --- a/graph/tests/test_legacy_graphiti_scope_guard.py +++ b/graph/tests/test_legacy_graphiti_scope_guard.py @@ -27,6 +27,9 @@ def access_token(name, scopes): class LegacyGraphScopeGuardTests(unittest.TestCase): + def legacy_config(self): + return SimpleNamespace(portable=False, values=SERVER.CONFIG.values) + def test_legacy_8767_denies_restricted_principal_before_query(self): query_calls = [] @@ -39,9 +42,11 @@ def falkor_called(): raise AssertionError("restricted principal reached FalkorDB") restricted = access_token("synthetic-restricted", ["team:project"]) - with mock.patch.object(SERVER, "get_access_token", return_value=restricted), mock.patch.object( - SERVER, "graphiti", side_effect=graphiti_called - ), mock.patch.object(SERVER, "falkor_graph", side_effect=falkor_called): + with mock.patch.object(SERVER, "CONFIG", self.legacy_config()), mock.patch.object( + SERVER, "get_access_token", return_value=restricted + ), mock.patch.object(SERVER, "graphiti", side_effect=graphiti_called), mock.patch.object( + SERVER, "falkor_graph", side_effect=falkor_called + ): with self.assertRaises(SERVER.ToolError): asyncio.run(SERVER.graph_search("legacy canary")) with self.assertRaises(SERVER.ToolError): @@ -67,7 +72,9 @@ async def search(self, query, group_ids, num_results): ] full_access = access_token("synthetic-full", ["*"]) - with mock.patch.object(SERVER, "get_access_token", return_value=full_access), mock.patch.object( + with mock.patch.object(SERVER, "CONFIG", self.legacy_config()), mock.patch.object( + SERVER, "GRAPH_KEY", "backfill-v1" + ), mock.patch.object(SERVER, "get_access_token", return_value=full_access), mock.patch.object( SERVER, "graphiti", return_value=FakeGraphiti() ), mock.patch.object(SERVER, "_log"): payload = json.loads(asyncio.run(SERVER.graph_search("legacy canary"))) diff --git a/graph/tests/test_portable_graph.py b/graph/tests/test_portable_graph.py index 6187921..4cd033c 100644 --- a/graph/tests/test_portable_graph.py +++ b/graph/tests/test_portable_graph.py @@ -14,7 +14,6 @@ def install_fixture(home, owner='fixture-owner', port=26383): shutil.copytree(ROOT, home / 'graphiti', ignore=shutil.ignore_patterns('__pycache__')) shutil.copytree(ROOT.parent / 'memory' / 'bin', home / 'mem0' / 'bin', ignore=shutil.ignore_patterns('__pycache__')) - (home / 'config.json').write_text(json.dumps({'schema': 'borg-install/v1'})) settings = dict(BORG_HOME=str(home), BORG_OWNER_ID=owner, BORG_MEMORY_SCOPE='personal:'+owner, BORG_QDRANT_URL='http://127.0.0.1:26333', BORG_QDRANT_COLLECTION='fixture', BORG_HISTORY_DB=str(home/'mem0/data/history.db'), BORG_OLLAMA_URL='http://127.0.0.1:21434', @@ -23,6 +22,7 @@ def install_fixture(home, owner='fixture-owner', port=26383): BORG_EMBED_DIMS='16', BORG_FALKORDB_HOST='127.0.0.1', BORG_FALKORDB_PORT=str(port), BORG_FALKORDB_GRAPH='fixture_legacy', BORG_GRAPH_LLM_URL='http://127.0.0.1:21460/v1', BORG_GRAPH_MODEL='fixture-graph') + (home / 'config.json').write_text(json.dumps({'schema': 'borg-install/v1', **settings})) env = {k:v for k,v in os.environ.items() if not k.startswith(('BORG_', 'MEM0_', 'GRAPH_', 'OLLAMA_', 'SHIM_'))} env.update(settings, PYTHONDONTWRITEBYTECODE='1') return env @@ -34,6 +34,36 @@ def run_code(home, env, code): class PortableGraphTests(unittest.TestCase): + def test_fixture_config_is_self_sufficient_with_only_borg_home(self): + with tempfile.TemporaryDirectory() as temp: + root = Path(temp) + home = root / 'one' + clean_home = root / 'clean-home' + clean_home.mkdir() + install_fixture(home) + env = { + 'PATH': os.environ.get('PATH', ''), + 'HOME': str(clean_home), + 'BORG_HOME': str(home), + 'PYTHONDONTWRITEBYTECODE': '1', + } + result = run_code(home, env, ''' +import sys +sys.path.insert(0, 'graphiti/bin') +import borg_config as config +required = { + 'BORG_OWNER_ID', 'BORG_MEMORY_SCOPE', 'BORG_QDRANT_URL', + 'BORG_QDRANT_COLLECTION', 'BORG_HISTORY_DB', 'BORG_OLLAMA_URL', + 'BORG_EXTRACTION_MODEL', 'BORG_EXTRACTION_MODEL_ID', 'BORG_EMBED_MODEL', + 'BORG_EMBED_MODEL_ID', 'BORG_EMBED_DIMS', 'BORG_FALKORDB_HOST', + 'BORG_FALKORDB_PORT', 'BORG_FALKORDB_GRAPH', 'BORG_GRAPH_LLM_URL', + 'BORG_GRAPH_MODEL', +} +assert required <= set(config.CONFIG.values) +assert config.CONFIG.owner_id == 'fixture-owner' +''') + self.assertEqual(result.returncode, 0, result.stderr) + def test_adapter_uses_explicit_installer_config(self): with tempfile.TemporaryDirectory() as temp: home = Path(temp)/'one' diff --git a/installer/health.py b/installer/health.py index 0f71f5b..2ef5b9e 100644 --- a/installer/health.py +++ b/installer/health.py @@ -3,20 +3,49 @@ import json from datetime import datetime, timezone +import os from pathlib import Path +import re import shlex import socket +import stat import subprocess import time import urllib.error import urllib.request -def get_json(url: str, *, timeout: int = 8) -> dict: - with urllib.request.build_opener(urllib.request.ProxyHandler({})).open(url, timeout=timeout) as response: +def get_json(url: str, *, timeout: int = 8, headers: dict | None = None) -> dict: + request = urllib.request.Request(url, headers=headers or {}) + with urllib.request.build_opener(urllib.request.ProxyHandler({})).open(request, timeout=timeout) as response: return json.load(response) +def conductor_headers(root: Path) -> dict[str, str]: + """Read the conductor token without following a substituted private path.""" + token_path = root / "conductors/primary/profile/.conductor/http-token" + try: + directory = token_path.parent.lstat() + except FileNotFoundError: + return {} + if (not stat.S_ISDIR(directory.st_mode) or token_path.parent.is_symlink() + or directory.st_uid != os.getuid() or directory.st_mode & 0o077): + raise ValueError("Unsafe conductor token directory") + try: + fd = os.open(token_path, os.O_RDONLY | os.O_NOFOLLOW) + except FileNotFoundError: + return {} + with os.fdopen(fd, "r") as handle: + info = os.fstat(handle.fileno()) + if (not stat.S_ISREG(info.st_mode) or info.st_uid != os.getuid() + or info.st_nlink != 1 or info.st_mode & 0o077 or info.st_size > 128): + raise ValueError("Unsafe conductor token file") + token = handle.read(129).strip() + if not re.fullmatch(r"[a-f0-9]{64}", token): + raise ValueError("Malformed conductor token") + return {"Authorization": "Bearer " + token} + + def mcp_call(url: str, name: str, authorization: str, arguments: dict | None = None) -> dict: payload = {"jsonrpc": "2.0", "id": 1, "method": "tools/call", "params": {"name": name, "arguments": arguments or {}}} @@ -137,7 +166,8 @@ def status(doc: dict) -> dict: components["models"] = {"state": "unavailable"} if "conductor" in enabled: try: - native = get_json(f"http://127.0.0.1:{ports['conductor']}/status") + native = get_json(f"http://127.0.0.1:{ports['conductor']}/status", + headers=conductor_headers(root)) expected_profile = str(root / "conductors/primary/profile") matched = native.get("ok") is True and native.get("port") == ports["conductor"] and native.get("codexHome") == expected_profile components["conductor"]["app_server"] = "initialized" if matched else "identity_mismatch" @@ -191,8 +221,9 @@ def status(doc: dict) -> dict: if "conductor" in enabled and blueprint.full(doc): try: payload = {"method": "hooks/list", "params": {"cwd": str(root / "projects")}, "timeoutMs": 8000} + headers = {"Content-Type": "application/json", **conductor_headers(root)} request = urllib.request.Request(f"http://127.0.0.1:{ports['conductor']}/rpc", - data=json.dumps(payload).encode(), headers={"Content-Type": "application/json"}) + data=json.dumps(payload).encode(), headers=headers) with urllib.request.build_opener(urllib.request.ProxyHandler({})).open(request, timeout=10) as response: observed = json.load(response) native = observed.get("result", observed) @@ -271,8 +302,9 @@ def tools_client_health(doc: dict) -> dict: root = Path(doc["home"]) try: payload = {"method": "config/read", "params": {"includeLayers": True}, "timeoutMs": 8000} + headers = {"Content-Type": "application/json", **conductor_headers(root)} request = urllib.request.Request(f"http://127.0.0.1:{doc['ports']['conductor']}/rpc", - data=json.dumps(payload).encode(), headers={"Content-Type": "application/json"}) + data=json.dumps(payload).encode(), headers=headers) with urllib.request.build_opener(urllib.request.ProxyHandler({})).open(request, timeout=10) as response: observed = json.load(response) native = observed.get("result", observed) diff --git a/installer/test_conductor_auth.py b/installer/test_conductor_auth.py new file mode 100644 index 0000000..a6d9804 --- /dev/null +++ b/installer/test_conductor_auth.py @@ -0,0 +1,44 @@ +from pathlib import Path +import tempfile +import unittest + +from installer import health + + +class ConductorAuthTests(unittest.TestCase): + def setUp(self): + self.temp = tempfile.TemporaryDirectory() + self.addCleanup(self.temp.cleanup) + self.root = Path(self.temp.name) / "borg" + self.token = self.root / "conductors/primary/profile/.conductor/http-token" + + def write_token(self, mode=0o600): + self.token.parent.mkdir(parents=True, mode=0o700) + self.token.write_text("a" * 64 + "\n") + self.token.chmod(mode) + + def subject(self): + self.assertTrue(hasattr(health, "conductor_headers"), "installer health must expose conductor_headers") + return health.conductor_headers + + def test_private_conductor_token_becomes_a_bearer_header(self): + self.write_token() + self.assertEqual(self.subject()(self.root), { + "Authorization": "Bearer " + "a" * 64, + }) + + def test_missing_token_preserves_report_mode_compatibility(self): + self.assertEqual(self.subject()(self.root), {}) + + def test_unsafe_or_malformed_token_is_never_sent(self): + self.write_token(mode=0o644) + with self.assertRaises(ValueError): + self.subject()(self.root) + self.token.chmod(0o600) + self.token.write_text("not-a-token\n") + with self.assertRaises(ValueError): + self.subject()(self.root) + + +if __name__ == "__main__": + unittest.main()