diff --git a/.env.example b/.env.example index ffa356f..312ed92 100644 --- a/.env.example +++ b/.env.example @@ -12,6 +12,9 @@ FOREMAN_REQUEST_TIMEOUT_MS=120000 UHP_BASE_URL=http://127.0.0.1:8787 UHP_HARNESS_ID= UHP_MODEL= +# Bearer credential for the bridge or UHP server at UHP_BASE_URL. If you started the bridge +# with LOCAL_CLI_UHP_TOKEN set, use the same value here and in FOREMAN_WORKSPACE_BRIDGE_TOKEN. +# Leave both empty for a bridge started without a token (it then accepts any local process). UHP_TOKEN= HINDSIGHT_BASE_URL=http://127.0.0.1:8888 HINDSIGHT_TOKEN= @@ -19,6 +22,9 @@ HINDSIGHT_TOKEN= # configured as LOCAL_CLI_UHP_SOURCE_REPO in the bridge process. FOREMAN_WORKSPACE_SOURCE_REPO= FOREMAN_WORKSPACE_BRIDGE_URL= +# Optional bearer token for the bridge above (its LOCAL_CLI_UHP_TOKEN). Requires the URL. +# Foreman's own per-project bridges need no configuration: it generates their tokens. +FOREMAN_WORKSPACE_BRIDGE_TOKEN= FOREMAN_WORKSPACE_ALLOWED_SCOPE= # JSON argv list; commands are run by Foreman in a disposable validation # workspace. A workspace workflow requires a source repo, bridge URL, allowed diff --git a/README.md b/README.md index be20d9c..9de8722 100644 --- a/README.md +++ b/README.md @@ -30,7 +30,7 @@ Open `http://127.0.0.1:4399` and choose **Open repository**. Pick a local Git re The selected repository must have a committed HEAD. Uncommitted source changes are visible but runs start from the pinned commit. The folder browser shows folders on the server machine; if Foreman runs remotely, the paths are on that machine. -No `.env` file or separate bridge command is needed for the standard local flow. Foreman starts its own per-project bridge process. +No `.env` file or separate bridge command is needed for the standard local flow. Foreman starts its own per-project bridge process, protects it with a random token, logs it, and restarts it if it crashes; see [Local bridge](#local-bridge). For the manual bridge workflow (custom CLIs, external UHP servers, or disposable-repo testing), see [the host CLI workflow guide](docs/three-harness-workflow.md). @@ -135,7 +135,7 @@ The standard local flow requires no environment variables. The following setting | `FOREMAN_DATA_DIR` | `.foreman-data` | Durable state directory. | | `UHP_BASE_URL` | — | External UHP server base URL. Without this, Foreman uses its own per-project bridge. | | `UHP_HARNESS_ID`, `UHP_MODEL` | — | Optional explicit initial harness/model for the external UHP server. Must be set together. | -| `UHP_TOKEN` | — | Optional bearer credential for the external UHP server. | +| `UHP_TOKEN` | — | Optional bearer credential for the external UHP server, sent on every request to `UHP_BASE_URL` including discovery. For a manual bridge started with `LOCAL_CLI_UHP_TOKEN`, use the same value. | | `HINDSIGHT_BASE_URL` | — | Hindsight API base URL. Outages show degraded memory status and do not stop workflow. | | `HINDSIGHT_TOKEN` | — | Optional credential for Hindsight. | | `FOREMAN_REQUEST_TIMEOUT_MS` | `120000` | Initial HTTP connection timeout (max 120,000 ms). Does not limit turn duration. | @@ -143,6 +143,7 @@ The standard local flow requires no environment variables. The following setting | `FOREMAN_WORKER_TIMEOUT_MS` | `600000` | Turn timeout for Worker roles (max 900,000 ms). | | `FOREMAN_WORKSPACE_SOURCE_REPO` | — | Local Git repository for snapshot verification. Required with the manual bridge workflow. | | `FOREMAN_WORKSPACE_BRIDGE_URL` | — | Loopback URL for an external workspace bridge. | +| `FOREMAN_WORKSPACE_BRIDGE_TOKEN` | — | Optional bearer token for that bridge, for a bridge started with `LOCAL_CLI_UHP_TOKEN` (use the same value). Requires `FOREMAN_WORKSPACE_BRIDGE_URL`; printable ASCII without spaces. It is only ever sent to that loopback URL and is never logged or returned by the API. | | `FOREMAN_WORKSPACE_ALLOWED_SCOPE` | — | Comma-separated exact paths or directory prefixes ending in `/` allowed in Worker results. | | `FOREMAN_VALIDATION_COMMANDS` | — | JSON array of `{"name","command","args","cwd?","network?"}` entries run in the disposable validation workspace. `network` is `true` or `false`; validation runs offline unless it is `true`, and a package-manager install without the field defaults to `true`. See [Validation sandbox](#validation-sandbox). | | `FOREMAN_VALIDATION_TIMEOUT_MS` | `120000` | Per-command time limit (max 600,000 ms). | @@ -152,9 +153,20 @@ The standard local flow requires no environment variables. The following setting Planner and Orchestrator turn timeouts are fixed at 300 seconds and are not configurable via environment variable. The optional configuration shape is recorded in `config.schema.json`. +## Local bridge + +In the standard flow Foreman runs one bridge per project, `investigations/local-cli-uhp/server.mjs`, on a random `127.0.0.1` port. The bridge drives your signed-in Claude Code, Codex and Antigravity CLIs and holds full copies of the repository, so Foreman protects and supervises it: + +- **Authentication.** Each start gets a fresh random 32-byte bearer token, passed to the bridge in its environment and sent by Foreman on every call: UHP requests, workspace seed, overlay and snapshot, and usage. The bridge answers `401` to any request without it, before doing any work (`/v1/uhp` discovery included), and `403` to any request whose `Host` is not `127.0.0.1`, `localhost` or `[::1]` on its port, so neither another local process nor a web page using DNS rebinding can read the repository or use your subscriptions. The token is never logged and never appears in events or `/api/*` responses. +- **Log.** The bridge's stdout and stderr go to `/local-bridges//bridge.log` (mode 0600). At every (re)start, a log over 5 MB is moved to `bridge.log.1`, replacing any older one, so at most one old file is kept. +- **Restart.** If the bridge exits unexpectedly, Foreman restarts it on the same port with the same token, so URLs and credentials already handed out stay valid, after 1 s, then 2 s, 4 s and so on up to 30 s. Anything the bridge was running is marked failed by the bridge, and Foreman's normal reconciliation shows it as failed; nothing is replayed. After 5 consecutive runs that each end within 60 s of starting, including restarts that cannot bind the port, Foreman gives up. A stopped project, such as a deleted one, is never restarted. +- **Status.** `GET /api/projects/:id/workspace-setup` reports `bridgeStatus` (`ready`, `restarting` or `unavailable`) and a `bridgeHealth` object with the last exit code or signal, the restart count and the log path. `GET /api/projects/:id/usage` reports the same `bridgeStatus` while the bridge is not ready. Foreman also prints each change to its own stderr. Reopening the repository starts a fresh bridge for a project that was given up on. + +For a bridge you start yourself (`UHP_BASE_URL`), see [the host CLI workflow guide](docs/three-harness-workflow.md#bridge-authentication): set `LOCAL_CLI_UHP_TOKEN` on the bridge and the same value as `UHP_TOKEN` and `FOREMAN_WORKSPACE_BRIDGE_TOKEN` for Foreman. Without a token that bridge accepts any local process and warns at startup, and Foreman does not supervise it. + ## Security -Foreman has no login; it relies on binding to `127.0.0.1` and on request checks that stop other web pages from driving it through your browser. Every request must carry a `Host` of `localhost`, `127.0.0.1` or `[::1]` (or the `FOREMAN_HOST` value) on `FOREMAN_PORT`, which blocks DNS rebinding. Every request other than `GET` and `HEAD` must also carry an `Origin` that matches `Host`, and a `Sec-Fetch-Site`, if sent, must be `same-origin`; otherwise it gets a 403. Scripts that call the API directly (for example with `curl`) must therefore send `-H 'Origin: http://127.0.0.1:4399'` on writes. `pnpm dev:ui` rewrites `Origin` on proxied requests to `http://127.0.0.1:4399`. GitHub write actions additionally require an explicit confirmation token. Binding to a non-loopback address with `FOREMAN_HOST` exposes an unauthenticated control plane to that network; don't. +Foreman has no login; it relies on binding to `127.0.0.1` and on request checks that stop other web pages from driving it through your browser. Every request must carry a `Host` of `localhost`, `127.0.0.1` or `[::1]` (or the `FOREMAN_HOST` value) on `FOREMAN_PORT`, which blocks DNS rebinding. Every request other than `GET` and `HEAD` must also carry an `Origin` that matches `Host`, and a `Sec-Fetch-Site`, if sent, must be `same-origin`; otherwise it gets a 403. Scripts that call the API directly (for example with `curl`) must therefore send `-H 'Origin: http://127.0.0.1:4399'` on writes. `pnpm dev:ui` rewrites `Origin` on proxied requests to `http://127.0.0.1:4399`. GitHub write actions additionally require an explicit confirmation token. Binding to a non-loopback address with `FOREMAN_HOST` exposes an unauthenticated control plane to that network; don't. The per-project bridge behind it has its own token and `Host` checks; see [Local bridge](#local-bridge). ## Data and cleanup diff --git a/config.schema.json b/config.schema.json index b7f72c3..d18e307 100644 --- a/config.schema.json +++ b/config.schema.json @@ -8,6 +8,7 @@ "FOREMAN_PORT": { "type": "integer", "minimum": 1, "maximum": 65535, "default": 4399 }, "FOREMAN_DATA_DIR": { "type": "string", "default": ".foreman-data" }, "UHP_BASE_URL": { "type": "string", "format": "uri", "description": "Configured UHP server base URL; absent means disconnected" }, + "UHP_TOKEN": { "type": "string", "writeOnly": true, "minLength": 1, "description": "Optional bearer credential for the UHP server, sent on every request including discovery; use the bridge's LOCAL_CLI_UHP_TOKEN for a manually started bridge" }, "UHP_HARNESS_ID": { "type": "string", "description": "Optional explicit initial harness choice" }, "UHP_MODEL": { "type": "string", "description": "Optional explicit initial model choice; pair with UHP_HARNESS_ID" }, "HINDSIGHT_BASE_URL": { "type": "string", "format": "uri", "description": "Configured Hindsight server base URL; absent means degraded memory" }, @@ -16,6 +17,7 @@ "FOREMAN_WORKER_TIMEOUT_MS": { "type": "integer", "minimum": 1000, "maximum": 900000, "default": 600000 }, "FOREMAN_WORKSPACE_SOURCE_REPO": { "type": "string", "description": "Local Git source repository used for pinned Worker snapshot verification" }, "FOREMAN_WORKSPACE_BRIDGE_URL": { "type": "string", "format": "uri", "description": "Loopback URL for the external workspace bridge" }, + "FOREMAN_WORKSPACE_BRIDGE_TOKEN": { "type": "string", "writeOnly": true, "minLength": 1, "maxLength": 512, "pattern": "^[\\x21-\\x7e]+$", "description": "Optional bearer token for the external workspace bridge (its LOCAL_CLI_UHP_TOKEN); printable ASCII without spaces; requires FOREMAN_WORKSPACE_BRIDGE_URL" }, "FOREMAN_WORKSPACE_ALLOWED_SCOPE": { "type": "string", "description": "Comma-separated exact paths or directory prefixes ending in / allowed for Worker changes" }, "FOREMAN_VALIDATION_COMMANDS": { "type": "string", "description": "JSON array of validation command objects {name, command, args, cwd?, network?} with explicit command and args. network is true or false (no other value is accepted): validation commands run without network access unless it is true. When omitted, a pnpm, npm or yarn install (first argument install, i, ci or add, or bare yarn) defaults to true and every other command to false" }, "FOREMAN_VALIDATION_TIMEOUT_MS": { "type": "integer", "minimum": 1, "maximum": 600000, "default": 120000 }, diff --git a/docs/three-harness-workflow.md b/docs/three-harness-workflow.md index d4ab7cf..ee095fd 100644 --- a/docs/three-harness-workflow.md +++ b/docs/three-harness-workflow.md @@ -61,7 +61,14 @@ configuration. Validation commands run in a disposable validation workspace with bounded time and output. Only configured commands run. Do not use a valuable checkout as the task repository. -Start the bridge in one terminal from the repository root: +Choose a bridge token first (see [Bridge authentication](#bridge-authentication)): + +```sh +export LOCAL_CLI_UHP_TOKEN="$(node -e "console.log(require('node:crypto').randomBytes(32).toString('hex'))")" +echo "$LOCAL_CLI_UHP_TOKEN" # you will paste this into Foreman's .env below +``` + +Start the bridge in that terminal from the repository root: ```sh cd investigations/local-cli-uhp @@ -84,7 +91,36 @@ listed on the proof host and completed the live Worker turn. If it is absent on another host, configure an explicitly listed Flash model instead. Any discovered AGY Flash model can be selected per role and per run in Foreman's UI. Do not let the bridge silently substitute a model. The bridge listens on -loopback and has no authentication; keep it local. +loopback only. If it reports `WARNING: LOCAL_CLI_UHP_TOKEN is not set` it is +running unauthenticated; see the next section. + +### Bridge authentication + +The bridge drives your signed-in CLIs and holds full copies of the repository, +so it can be protected with a bearer token. It has two modes: + +- **Token set (recommended).** With `LOCAL_CLI_UHP_TOKEN` set (1-512 printable + ASCII characters, no spaces; 32 random bytes as hex is a good choice), every + request must carry `Authorization: Bearer `, including `/v1/uhp` + discovery and the workspace extension routes. The token is compared in + constant time. A missing or wrong token gets `401` before any other work is + done. Foreman itself uses this mode for the per-project bridges it starts, + with a fresh random token per start. +- **Token unset (legacy manual mode).** The bridge accepts requests from any + local process, and prints one `WARNING: LOCAL_CLI_UHP_TOKEN is not set` line at + startup. The `smoke.mjs`, `workspace-smoke.mjs`, `codex-worker-smoke.mjs` and + `reviewer-smoke.mjs` scripts under `investigations/local-cli-uhp/` send no + token, so they need a bridge in this mode. + +In both modes the bridge checks the `Host` header of every request: only +`127.0.0.1:`, `localhost:` and `[::1]:` (case-insensitive, one +trailing dot tolerated) are accepted, anything else gets `403`. This blocks DNS +rebinding from web pages, and it does not depend on the token. + +When the bridge has a token, give Foreman the same value (see below) in both +`UHP_TOKEN` (UHP calls) and `FOREMAN_WORKSPACE_BRIDGE_TOKEN` (workspace +seed, overlay and snapshot calls). Foreman never logs either value or returns +them from its API. In another terminal, configure the same repo and validation policy for Foreman, then build and start the UI/API service: @@ -98,8 +134,10 @@ Set these entries in `.env` (retain the loopback URL and adjust paths/checks): ```dotenv UHP_BASE_URL=http://127.0.0.1:8787 +UHP_TOKEN= FOREMAN_WORKSPACE_SOURCE_REPO=/absolute/path/to/disposable-repo FOREMAN_WORKSPACE_BRIDGE_URL=http://127.0.0.1:8787 +FOREMAN_WORKSPACE_BRIDGE_TOKEN= FOREMAN_WORKSPACE_ALLOWED_SCOPE=README.md FOREMAN_VALIDATION_COMMANDS=[{"name":"build","command":"pnpm","args":["build"]}] ``` @@ -114,7 +152,10 @@ pnpm start Open `http://127.0.0.1:4399`. The bridge and Foreman have separate processes and configuration; `.env` configures Foreman, while the shell variables above -configure the bridge. Hindsight is optional and advisory. +configure the bridge. If you started the bridge without a token, leave +`UHP_TOKEN` and `FOREMAN_WORKSPACE_BRIDGE_TOKEN` out. Foreman does not restart +a bridge you started yourself; if it stops, start it again. Hindsight is +optional and advisory. If AGY reports that it needs sign-in, use the interactive sign-in flow offered by the installed `agy` CLI in a terminal as the same host user, then restart diff --git a/investigations/local-cli-uhp/README.md b/investigations/local-cli-uhp/README.md index 567060b..a948a5f 100644 --- a/investigations/local-cli-uhp/README.md +++ b/investigations/local-cli-uhp/README.md @@ -72,7 +72,22 @@ provider-selection overrides are excluded. Do not set provider API keys for this experiment. State is stored at `LOCAL_CLI_UHP_STATE` (default `/tmp/local-cli-uhp-state.json`, mode 0600) and work directories under `LOCAL_CLI_UHP_WORK` (default `/tmp/local-cli-uhp-work`). -The HTTP listener binds loopback only and has no authentication; keep it local. +The HTTP listener binds loopback only. + +Authentication: set `LOCAL_CLI_UHP_TOKEN` (1-512 printable ASCII characters, no +spaces) and every request, including `/v1/uhp` discovery and the workspace +extension routes, must carry `Authorization: Bearer `; a missing or wrong +token gets `401` before any other work, compared in constant time. Without it +the bridge accepts requests from any local process and prints a startup +warning. Foreman starts its own bridges with a fresh random token per start. The +token is removed from the bridge's environment, so the CLIs it spawns never see +it. Independently of the token, the `Host` header must be `127.0.0.1:`, +`localhost:` or `[::1]:` (case-insensitive, one trailing dot +tolerated) or the request gets `403`, which stops DNS-rebinding attacks. The +smoke scripts in this directory send no token and need a bridge without one. +If the port cannot be bound the bridge exits with status 1 and a message on +stderr; Foreman's supervisor logs that output to `bridge.log` and treats it as a +failed restart. ### Claude subscription usage diff --git a/investigations/local-cli-uhp/server.mjs b/investigations/local-cli-uhp/server.mjs index a823925..7ce18d4 100644 --- a/investigations/local-cli-uhp/server.mjs +++ b/investigations/local-cli-uhp/server.mjs @@ -4,7 +4,7 @@ import { createServer } from 'node:http'; import { spawn } from 'node:child_process'; import { randomUUID } from 'node:crypto'; -import { createHash } from 'node:crypto'; +import { createHash, timingSafeEqual } from 'node:crypto'; import { mkdir, readFile, rename, writeFile, readdir, lstat, readlink, realpath, rm, access, symlink, chmod, open } from 'node:fs/promises'; import { accessSync, constants as fsConstants, readFileSync } from 'node:fs'; import { tmpdir, homedir } from 'node:os'; @@ -20,6 +20,12 @@ import { readClaudeControlUsage } from './claude-quota.mjs'; const VERSION = '2026-09-12'; const utf8 = new TextDecoder('utf-8', { fatal: true }); const PORT = Number(process.env.LOCAL_CLI_UHP_PORT ?? 8787); +// Bearer-token authentication. When LOCAL_CLI_UHP_TOKEN is set, every request must carry it. The variable is +// removed from the environment so nothing that inherits process.env can ever see it. +const AUTH_TOKEN = (process.env.LOCAL_CLI_UHP_TOKEN ?? '').trim(); +delete process.env.LOCAL_CLI_UHP_TOKEN; +if (AUTH_TOKEN && !/^[\x21-\x7e]{1,512}$/.test(AUTH_TOKEN)) throw new Error('LOCAL_CLI_UHP_TOKEN must be 1-512 printable ASCII characters without spaces'); +const AUTH_TOKEN_DIGEST = AUTH_TOKEN ? createHash('sha256').update(AUTH_TOKEN).digest() : undefined; const STATE = resolve(process.env.LOCAL_CLI_UHP_STATE ?? join(tmpdir(), 'local-cli-uhp-state.json')); const ROOT = resolve(process.env.LOCAL_CLI_UHP_WORK ?? join(tmpdir(), 'local-cli-uhp-work')); if (ROOT !== tmpdir() && !ROOT.startsWith(`${tmpdir()}/`)) throw new Error('LOCAL_CLI_UHP_WORK must be under the system temporary directory'); @@ -1283,8 +1289,25 @@ function safeAgyDiagnostic(stderr) { return 'no safe diagnostic details'; } +// DNS-rebinding guard: only accept the loopback names on the port this server is actually bound to. +// Case-insensitive, and a single trailing dot on the name is tolerated ("localhost.:8787"). +function hostAllowed(header) { + const match = typeof header === 'string' ? /^(?:(?:127\.0\.0\.1|localhost)\.?|\[::1\]):(\d{1,5})$/i.exec(header) : null; + return !!match && Number(match[1]) === (server.address()?.port ?? PORT); +} +// Constant-time bearer check: compare SHA-256 digests so both timingSafeEqual inputs always have the same length. +function tokenAccepted(header) { + if (!AUTH_TOKEN_DIGEST) return true; + const supplied = typeof header === 'string' ? /^Bearer +([\x21-\x7e]+)$/i.exec(header)?.[1] : undefined; + return timingSafeEqual(createHash('sha256').update(supplied ?? '').digest(), AUTH_TOKEN_DIGEST); +} + const server = createServer(async (req, res) => { - const url = new URL(req.url, 'http://localhost'); + // Runs before routing, URL parsing and body handling, so a rejected request costs nothing and touches no state. + if (!hostAllowed(req.headers.host)) return send(res, 403, { error: { code: 'host_forbidden', message: 'Host must be 127.0.0.1, localhost or [::1] on the bridge port' } }, { connection: 'close' }); + if (!tokenAccepted(req.headers.authorization)) return send(res, 401, { error: { code: 'unauthorized', message: 'A valid bearer token is required' } }, { 'www-authenticate': 'Bearer', connection: 'close' }); + let url; + try { url = new URL(req.url, 'http://localhost'); } catch { return send(res, 400, { error: { code: 'invalid_url' } }); } try { if (req.method === 'GET' && url.pathname === '/v1/uhp') return send(res, 200, { protocol: 'uhp', versions: [VERSION], default_version: VERSION, implementation: { name: 'local-cli-uhp', experimental: true }, capabilities: { streaming: true, idempotency: true, sessions: true, cancellation: true, readOnlyReviewer: true, extensions: { foreman_workspace_bridge_v1: { version: 1, seed: !!SOURCE_REPO, complete_snapshot: true, execution_boundary: 'bubblewrap' } } } }); if (req.method === 'POST' && url.pathname === '/extensions/foreman-workspace/v1/workspaces') { @@ -1452,5 +1475,7 @@ if (process.argv[2] === '--preflight-claude-runtime') { console.log(await preflightAgyRuntime()); } else { await load(); - server.listen(PORT, '127.0.0.1', () => console.log(`experimental local CLI UHP listening on 127.0.0.1:${PORT}`)); + server.on('error', error => { console.error(`local CLI UHP could not listen on 127.0.0.1:${PORT}: ${error.message}`); process.exit(1); }); + if (!AUTH_TOKEN) console.warn('WARNING: LOCAL_CLI_UHP_TOKEN is not set; this bridge accepts unauthenticated requests from any local process (set it to require "Authorization: Bearer ").'); + server.listen(PORT, '127.0.0.1', () => console.log(`experimental local CLI UHP listening on 127.0.0.1:${PORT}${AUTH_TOKEN ? ' (bearer token required)' : ''}`)); } diff --git a/investigations/local-cli-uhp/test.mjs b/investigations/local-cli-uhp/test.mjs index 490a794..4ed7834 100644 --- a/investigations/local-cli-uhp/test.mjs +++ b/investigations/local-cli-uhp/test.mjs @@ -4,6 +4,8 @@ import { mkdtemp, writeFile, chmod, readFile, rm, truncate, mkdir, readdir, stat import { tmpdir } from 'node:os'; import { join, dirname, resolve } from 'node:path'; import { spawn } from 'node:child_process'; +import { request as httpRequest } from 'node:http'; +import { createServer as createNetServer } from 'node:net'; import { fileURLToPath } from 'node:url'; await import('tsx/esm/api').then(({ register }) => register()); @@ -11,11 +13,14 @@ const { UhpClient } = await import('../../src/uhp.ts'); const { Controller } = await import('../../src/controller.ts'); const { JsonStore } = await import('../../src/store.ts'); const { snapshotGitCommit } = await import('../../src/git-workspace.ts'); +const { seedBridgeWorkspace, fetchBridgeSnapshot, overlayBridgeWorkspace } = await import('../../src/verified-workspace.ts'); +const { bearerFetch } = await import('../../src/local-bridge.ts'); const { verifyBridgeWorkspace, validateBridgeSnapshot } = await import('./workspace-verifier.mjs'); const { createWorkspaceFixture, applyAllFileCaseChanges } = await import('./workspace-fixture.mjs'); const { assertReviewerBounds, isPrepareOnly, REVIEWER_TASK_BOUNDS, REVIEWER_STREAM_INACTIVITY_TIMEOUT_MS, REVIEWER_BRIDGE_MAX_STEP } = await import('./reviewer-smoke-bounds.mjs'); const here = dirname(fileURLToPath(import.meta.url)); +const BRIDGE_TOKEN = '0123456789abcdef'.repeat(4); const dirs = []; async function fixtureCli(dir, name, body) { const path = join(dir, name); @@ -42,11 +47,17 @@ async function setup(t, options = {}) { if (options.noAgyModel) delete env.AGY_MODEL; if (!options.agyWorkerEffort) delete env.AGY_WORKER_EFFORT; if (options.agyEnabled) Object.assign(env,{AGY_CONFIG_DIR:join(dir,'agy-auth'),...(options.noAgyModel?{}:{AGY_MODEL:options.agyModel ?? 'gemini-3.8-flash-medium'}),AGY_BIN:agy}); - let proc = spawn(process.execPath, [join(here, 'server.mjs')], { env, cwd: dir, stdio: 'ignore' }); + // Tests run without a bridge token unless one is requested, whatever the developer's shell exports. + delete env.LOCAL_CLI_UHP_TOKEN; + if (options.token) env.LOCAL_CLI_UHP_TOKEN = options.token; + const authHeaders = options.token ? { authorization: `Bearer ${options.token}` } : {}; + let output = ''; + const startBridge = () => { const child = spawn(process.execPath, [join(here, 'server.mjs')], { env, cwd: dir, stdio: options.captureOutput ? ['ignore','pipe','pipe'] : 'ignore' }); if (options.captureOutput) for (const stream of [child.stdout, child.stderr]) stream.on('data', chunk => { output += chunk; }); return child; }; + let proc = startBridge(); t.after(async () => { if (proc.exitCode === null) { proc.kill('SIGTERM'); await new Promise(r => proc.once('exit', r)); } await rm(dir, { recursive: true, force: true }); }); const base = `http://127.0.0.1:${port}`; - for (let i=0;i<100;i++) { try { const r=await fetch(`${base}/v1/uhp`); if(r.ok) break; } catch {} await new Promise(r=>setTimeout(r,20)); } - return { base, dir, env, baseCommit: options.baseCommit ?? fixture?.baseCommit, sourceRepo: options.sourceRepo ?? fixture?.repo, countFor: workspaceId=>join(env.LOCAL_CLI_UHP_WORK,workspaceId,'.fixture-cli-count'), restart: async () => { proc.kill('SIGTERM'); await new Promise(r => proc.once('exit', r)); proc = spawn(process.execPath, [join(here, 'server.mjs')], { env, cwd: dir, stdio: 'ignore' }); for(let i=0;i<100;i++){try{if((await fetch(`${base}/v1/uhp`)).ok)break;}catch{} await new Promise(r=>setTimeout(r,20));} } }; + for (let i=0;i<250&&proc.exitCode===null;i++) { try { const r=await fetch(`${base}/v1/uhp`, { headers: authHeaders }); if(r.ok) break; } catch {} await new Promise(r=>setTimeout(r,20)); } + return { base, dir, env, port, token: options.token, authHeaders, output: () => output, baseCommit: options.baseCommit ?? fixture?.baseCommit, sourceRepo: options.sourceRepo ?? fixture?.repo, countFor: workspaceId=>join(env.LOCAL_CLI_UHP_WORK,workspaceId,'.fixture-cli-count'), restart: async () => { proc.kill('SIGTERM'); await new Promise(r => proc.once('exit', r)); proc = startBridge(); for(let i=0;i<250&&proc.exitCode===null;i++){try{if((await fetch(`${base}/v1/uhp`, { headers: authHeaders })).ok)break;}catch{} await new Promise(r=>setTimeout(r,20));} } }; } async function submit(base, harness, model, key, baseCommit, workspaceId, input = 'Say bounded answer') { const seeded=workspaceId ? {workspace_id:workspaceId} : await (await fetch(`${base}/extensions/foreman-workspace/v1/workspaces`,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({base_commit:baseCommit})})).json(); @@ -625,16 +636,17 @@ test('tool-enabled child cannot read or write outside its seeded workspace', asy assert.equal(evidence.validation,'verified_by_foreman_git_comparison'); }); -test('Codex Worker edits only its seeded workspace; Foreman verifies the complete snapshot and validates it', async t => { +// The full Foreman-to-bridge Worker flow, with or without a bridge token (Foreman wires it as server.ts does for a project bridge). +async function codexWorkerFlow(t, token) { const fixture = await createWorkspaceFixture(); t.after(fixture.cleanup); const parent = await mkdtemp(join(tmpdir(),'foreman-codex-boundary-')); t.after(()=>rm(parent,{recursive:true,force:true})); const sentinel=join(parent,'outside-sentinel'); await writeFile(sentinel,'FOREMAN-CODEX-OUTSIDE-SENTINEL'); const codexBody = `import {readFileSync,writeFileSync} from 'node:fs'; const path=${JSON.stringify(sentinel)}; let read='allowed',write='allowed'; try{readFileSync(path,'utf8')}catch{read='denied'} try{writeFileSync(path,'CHANGED')}catch{write='denied'} writeFileSync('README.md',${JSON.stringify('# Fixture\n\nCodex changed the assigned README.\n')}); writeFileSync('.boundary-result.json',JSON.stringify({read,write})); console.log(JSON.stringify({type:'thread.started',thread_id:'codex-worker-thread'})); console.log(JSON.stringify({type:'item.completed',item:{type:'agent_message',text:'Updated README.md in the assigned workspace.'}})); console.log(JSON.stringify({type:'turn.completed',usage:{input_tokens:11,output_tokens:7}}));`; const claudeBody=`let prompt='';process.stdin.setEncoding('utf8');process.stdin.on('data',chunk=>prompt+=chunk);process.stdin.on('end',()=>{const model='claude-actual';const result=prompt.includes('"workerTask"')&&prompt.includes('"targetFiles"')?JSON.stringify({workerTask:'Change README.md with one short sentence.',targetFiles:['README.md']}):'Planner recommends a concise README note.';console.log(JSON.stringify({type:'system',subtype:'init',model,session_id:'judgment-session'}));console.log(JSON.stringify({type:'result',subtype:'success',result,model,session_id:'judgment-session',usage:{input_tokens:5,output_tokens:2}}));});`; - const {base,env}=await setup(t,{sourceRepo:fixture.repo,baseCommit:fixture.baseCommit,codexBody,claudeBody}); - const uhp = new UhpClient({baseUrl:base,timeoutMs:20_000}); + const {base,env}=await setup(t,{sourceRepo:fixture.repo,baseCommit:fixture.baseCommit,codexBody,claudeBody,...(token?{token}:{})}); + const uhp = new UhpClient({baseUrl:base,timeoutMs:20_000,...(token?{token,fetch:bearerFetch(token)}:{})}); const controller = new Controller(new JsonStore(join(parent,'foreman-state.json')),uhp,false,true); - controller.configureVerifiedWorkspace({repoPath:fixture.repo,allowedScope:['README.md','.boundary-result.json'],commands:[{name:'assert verified README content',command:process.execPath,args:['-e',`const fs=require('node:fs');if(fs.readFileSync('README.md','utf8')!=='# Fixture\\n\\nCodex changed the assigned README.\\n')process.exit(1)`]}],bridgeBaseUrl:base,timeoutMs:10_000,maxOutputBytes:2_000}); + controller.configureVerifiedWorkspace({repoPath:fixture.repo,allowedScope:['README.md','.boundary-result.json'],commands:[{name:'assert verified README content',command:process.execPath,args:['-e',`const fs=require('node:fs');if(fs.readFileSync('README.md','utf8')!=='# Fixture\\n\\nCodex changed the assigned README.\\n')process.exit(1)`]}],bridgeBaseUrl:base,timeoutMs:10_000,maxOutputBytes:2_000,...(token?{bridgeToken:token}:{})}); await controller.refreshDiscovery(); const project=await controller.createProject('Codex disposable workspace fixture'); const task=await controller.createTask(project.id,'Edit the assigned README'); @@ -683,7 +695,9 @@ test('Codex Worker edits only its seeded workspace; Foreman verifies the complet assert.equal(await readFile(sentinel,'utf8'),'FOREMAN-CODEX-OUTSIDE-SENTINEL'); const boundary=JSON.parse(await readFile(join(env.LOCAL_CLI_UHP_WORK,stored.workspaceId,'.boundary-result.json'),'utf8')); assert.deepEqual(boundary,{read:'denied',write:'denied'}); -}); +} +test('Codex Worker edits only its seeded workspace; Foreman verifies the complete snapshot and validates it', t => codexWorkerFlow(t)); +test('the same Worker flow works end to end against a token-protected bridge, with the token on every Foreman-to-bridge call', t => codexWorkerFlow(t, BRIDGE_TOKEN)); test('failed boundary probe blocks the tool-enabled CLI before spawn', async t => { const fixture = await createWorkspaceFixture(); t.after(fixture.cleanup); @@ -913,3 +927,126 @@ test('activity SSE emits contiguous ordered events and bounds sanitized summarie assert.ok(summaries.every(summary=>summary.length<=200&&!/[\u0000-\u001f\u007f-\u009f]/.test(summary))); assert.equal(activities.every(item=>item.response.status==='in_progress'),true); }); + +// ---- Bearer-token authentication and Host validation ---- +// fetch() will not let a caller choose the Host header, so these tests speak raw HTTP. +function rawRequest(port, { method = 'GET', path = '/', headers = {}, body } = {}) { + return new Promise((resolveRequest, reject) => { + const req = httpRequest({ host: '127.0.0.1', port, method, path, headers, setHost: false }, res => { + let text = ''; res.setEncoding('utf8'); res.on('data', chunk => { text += chunk; }); res.on('end', () => resolveRequest({ status: res.statusCode, headers: res.headers, text })); + }); + req.on('error', reject); req.end(body); + }); +} +const jsonHeaders = { 'content-type': 'application/json' }; +function protectedRoutes(baseCommit) { + const workspace = 'ws_00000000-0000-0000-0000-000000000000'; + const seed = JSON.stringify({ base_commit: baseCommit }); + const submission = JSON.stringify({ input: 'unauthenticated prompt', model: 'claude-requested', metadata: { harness_id: 'claude-code' }, timeout_seconds: 5, max_step: 1 }); + return [ + ['GET', '/v1/uhp'], ['GET', '/v1/harnesses'], ['GET', '/v1/harnesses/claude-code/models'], ['GET', '/v1/usage'], + ['POST', '/extensions/foreman-workspace/v1/workspaces', seed], + ['GET', `/extensions/foreman-workspace/v1/workspaces/${workspace}/snapshot`], + ['POST', `/extensions/foreman-workspace/v1/workspaces/${workspace}/overlay`, JSON.stringify({ entries: [] })], + ['POST', '/v1/responses', submission, { 'idempotency-key': 'unauthenticated-key' }], + ['GET', '/v1/responses/resp_00000000-0000-0000-0000-000000000000'], + ['POST', '/v1/responses/resp_00000000-0000-0000-0000-000000000000/cancel', '{}'], + ]; +} + +test('with LOCAL_CLI_UHP_TOKEN set, a missing or wrong bearer token gets 401 on every route before any work', async t => { + const { port, env, token, baseCommit } = await setup(t, { token: BRIDGE_TOKEN }); + const wrong = [undefined, 'Bearer', 'Bearer ', 'Bearer wrong', `Bearer ${token.slice(0, -1)}`, `Bearer ${token}x`, `Bearer ${token.toUpperCase()}`, `Basic ${token}`, token, `Bearer ${token} ${token}`]; + for (const [method, path, body, extra] of protectedRoutes(baseCommit)) for (const authorization of wrong) { + const response = await rawRequest(port, { method, path, body, headers: { host: `127.0.0.1:${port}`, ...(body ? jsonHeaders : {}), ...extra, ...(authorization === undefined ? {} : { authorization }) } }); + const label = `${method} ${path} with ${authorization === undefined ? 'no Authorization' : JSON.stringify(authorization.replace(token, ''))}`; + assert.equal(response.status, 401, label); + assert.equal(JSON.parse(response.text).error.code, 'unauthorized', label); + assert.equal(response.headers['www-authenticate'], 'Bearer', label); + assert.ok(!response.text.includes(token), label); + } + // Nothing behind the rejected requests ran: no workspace was seeded and no idempotency intent was persisted. + assert.deepEqual((await readdir(env.LOCAL_CLI_UHP_WORK)).filter(name => name.startsWith('ws_')), []); + await assert.rejects(stat(env.LOCAL_CLI_UHP_STATE), { code: 'ENOENT' }); +}); + +test('the right bearer token unlocks discovery, seeding, snapshots, usage and submissions through Foreman clients', async t => { + const claudeBody = `process.stdin.resume();process.stdin.on('end',()=>{console.log(JSON.stringify({type:'system',subtype:'init',model:'claude-requested',session_id:'token-session'}));console.log(JSON.stringify({type:'result',subtype:'success',result:'authorized answer',model:'claude-requested',session_id:'token-session',usage:{input_tokens:1,output_tokens:1}}));});`; + const { base, authHeaders, token, baseCommit } = await setup(t, { token: BRIDGE_TOKEN, claudeBody }); + const discovery = await fetch(`${base}/v1/uhp`, { headers: { authorization: `bEaReR ${token}` } }); + assert.equal(discovery.status, 200); assert.equal((await discovery.json()).implementation.name, 'local-cli-uhp'); + assert.equal((await fetch(`${base}/v1/usage`, { headers: authHeaders })).status, 200); + // The Foreman workspace helpers: seed, overlay and read the snapshot back, all with the token. + const seeded = await seedBridgeWorkspace(base, baseCommit, { token }); + assert.equal(seeded.baseCommit, baseCommit); + const overlay = await overlayBridgeWorkspace(base, seeded.workspaceId, [{ path: 'README.md', contentBase64: Buffer.from('overlaid\n').toString('base64'), mode: '100644' }], { token }); + assert.equal(overlay.applied, 1); + const snapshot = await fetchBridgeSnapshot(base, seeded.workspaceId, baseCommit, { token }); + assert.equal(snapshot.complete, true); assert.ok(snapshot.entries.length > 0); + assert.equal(Buffer.from(snapshot.entries.find(entry => entry.path === 'README.md').contentBase64, 'base64').toString('utf8'), 'overlaid\n'); + // Without the token, or with a wrong one, the same helpers surface a 401 and never echo the token. + for (const call of [() => seedBridgeWorkspace(base, baseCommit), () => overlayBridgeWorkspace(base, seeded.workspaceId, []), () => fetchBridgeSnapshot(base, seeded.workspaceId, baseCommit), () => seedBridgeWorkspace(base, baseCommit, { token: `${token}0` })]) { + await assert.rejects(call, error => /\(401\)/.test(error.message) && !error.message.includes(token)); + } + // Foreman's client for the bridge (see createProjectRuntime): the token rides on every call, discovery included. + const input = { submissionId: 'sub-token', assignmentId: 'asg-token', runId: 'run-token', roleId: 'planner', taskId: 'task-token', projectId: 'prj-token', prompt: 'Reply once.', config: { harnessId: 'claude-code', model: 'claude-requested', timeoutSeconds: 5 }, idempotencyKey: 'token-submit-key' }; + const authorized = new UhpClient({ baseUrl: base, token, fetch: bearerFetch(token), harnessId: 'claude-code', model: 'claude-requested' }); + assert.equal((await authorized.discover()).version, '2026-09-12'); + assert.equal((await authorized.submit(input)).outputText, 'authorized answer'); + await assert.rejects(new UhpClient({ baseUrl: base, harnessId: 'claude-code', model: 'claude-requested' }).submit(input), /UHP request failed \(401\): unauthorized/); + await assert.rejects(new UhpClient({ baseUrl: base, token: `${token}0`, fetch: bearerFetch(`${token}0`), harnessId: 'claude-code', model: 'claude-requested' }).discover(), /UHP request failed \(401\): unauthorized/); +}); + +for (const withToken of [false, true]) test(`a foreign Host header gets 403 ${withToken ? 'even with the right token' : 'when no token is configured'}; loopback names on the bridge port are accepted`, async t => { + const { port, authHeaders, baseCommit } = await setup(t, { token: withToken ? BRIDGE_TOKEN : undefined }); + const foreign = [`evil.example:${port}`, 'evil.example', '127.0.0.1', 'localhost', '[::1]', `127.0.0.1:${port + 1}`, `localhost:${port}0`, `localhost.evil.example:${port}`, `127.0.0.1.evil.example:${port}`, `0.0.0.0:${port}`, `[::2]:${port}`, `::1:${port}`, `localhost..:${port}`, `user@localhost:${port}`]; + for (const host of foreign) for (const [method, path, body, extra] of protectedRoutes(baseCommit).filter((_, index) => index % 3 === 0)) { + const label = `Host ${JSON.stringify(host)} ${method} ${path}`; + const response = await rawRequest(port, { method, path, body, headers: { host, ...authHeaders, ...(body ? jsonHeaders : {}), ...extra } }); + assert.equal(response.status, 403, label); assert.equal(JSON.parse(response.text).error.code, 'host_forbidden', label); + } + // A foreign Host is refused before the missing or wrong token is even considered. + if (withToken) for (const authorization of [undefined, 'Bearer wrong']) assert.equal((await rawRequest(port, { path: '/v1/uhp', headers: { host: `evil.example:${port}`, ...(authorization ? { authorization } : {}) } })).status, 403); + for (const host of [`127.0.0.1:${port}`, `localhost:${port}`, `LOCALHOST:${port}`, `LocalHost.:${port}`, `127.0.0.1.:${port}`, `[::1]:${port}`]) { + const response = await rawRequest(port, { path: '/v1/uhp', headers: { host, ...authHeaders } }); + assert.equal(response.status, 200, host); assert.equal(JSON.parse(response.text).protocol, 'uhp'); + } +}); + +test('the bridge warns once at startup when unauthenticated, stays quiet when a token is set, and never prints the token', async t => { + const waitForListening = async bridge => { for (let i = 0; i < 500 && !/listening on/.test(bridge.output()); i++) await new Promise(r => setTimeout(r, 20)); assert.match(bridge.output(), /listening on/, 'the bridge never reported that it is listening'); }; + const open = await setup(t, { captureOutput: true }); await waitForListening(open); + assert.equal(open.output().split('\n').filter(line => line.includes('WARNING: LOCAL_CLI_UHP_TOKEN is not set')).length, 1, open.output()); + assert.equal((await fetch(`${open.base}/v1/uhp`)).status, 200); // still serves unauthenticated requests in the documented manual mode + const secured = await setup(t, { token: BRIDGE_TOKEN, captureOutput: true }); await waitForListening(secured); + assert.match(secured.output(), /listening on 127\.0\.0\.1:\d+ \(bearer token required\)/); + assert.doesNotMatch(secured.output(), /WARNING/); assert.ok(!secured.output().includes(BRIDGE_TOKEN)); +}); + +test('an unusable token or an occupied port stops the bridge at startup without printing the token', async () => { + const badToken = 'bad token with spaces'; + const bad = spawn(process.execPath, [join(here, 'server.mjs')], { env: { ...process.env, LOCAL_CLI_UHP_TOKEN: badToken, LOCAL_CLI_UHP_WORK: join(tmpdir(), `local-cli-uhp-bad-token-${process.pid}`) }, stdio: ['ignore', 'pipe', 'pipe'] }); + let badOutput = ''; for (const stream of [bad.stdout, bad.stderr]) stream.on('data', chunk => { badOutput += chunk; }); + assert.equal(await new Promise(resolveCode => bad.once('exit', resolveCode)), 1); + assert.match(badOutput, /LOCAL_CLI_UHP_TOKEN must be 1-512 printable ASCII characters without spaces/); assert.ok(!badOutput.includes(badToken)); + // A restart that cannot rebind its port must exit promptly with a diagnostic so a supervisor can count it as a failure. + const holder = createNetServer(); await new Promise(resolveListen => holder.listen(0, '127.0.0.1', resolveListen)); + const dir = await mkdtemp(join(tmpdir(), 'local-cli-uhp-port-')); dirs.push(dir); + try { + const busy = spawn(process.execPath, [join(here, 'server.mjs')], { env: { ...process.env, LOCAL_CLI_UHP_PORT: String(holder.address().port), LOCAL_CLI_UHP_STATE: join(dir, 'state.json'), LOCAL_CLI_UHP_WORK: join(dir, 'work') }, stdio: ['ignore', 'pipe', 'pipe'] }); + let busyOutput = ''; for (const stream of [busy.stdout, busy.stderr]) stream.on('data', chunk => { busyOutput += chunk; }); + assert.equal(await new Promise(resolveCode => busy.once('exit', resolveCode)), 1); + assert.match(busyOutput, /could not listen on 127\.0\.0\.1:\d+: .*EADDRINUSE/); + } finally { holder.close(); await rm(dir, { recursive: true, force: true }); } +}); + +test('an authenticated request with a malformed URL is rejected without crashing the bridge, and CLI children never see the token', async t => { + const claudeBody = `process.stdin.resume();process.stdin.on('end',()=>{const seen=Object.entries(process.env).some(([name,value])=>/TOKEN/i.test(name)||String(value).includes(${JSON.stringify(BRIDGE_TOKEN)}));console.log(JSON.stringify({type:'system',subtype:'init',model:'claude-requested',session_id:'env-session'}));console.log(JSON.stringify({type:'result',subtype:'success',result:'token_visible='+seen+';env_entries='+(Object.keys(process.env).length>0),model:'claude-requested',session_id:'env-session',usage:{input_tokens:1,output_tokens:1}}));});`; + const { base, port, token, authHeaders } = await setup(t, { token: BRIDGE_TOKEN, claudeBody }); + const malformed = await rawRequest(port, { path: '//[', headers: { host: `127.0.0.1:${port}`, ...authHeaders } }); + assert.equal(malformed.status, 400); assert.equal(JSON.parse(malformed.text).error.code, 'invalid_url'); + assert.equal((await fetch(`${base}/v1/uhp`, { headers: authHeaders })).status, 200); + const client = new UhpClient({ baseUrl: base, token, fetch: bearerFetch(token), harnessId: 'claude-code', model: 'claude-requested' }); + const result = await client.submit({ submissionId: 'sub-env', assignmentId: 'asg-env', runId: 'run-env', roleId: 'planner', taskId: 'task-env', projectId: 'prj-env', prompt: 'Report your environment.', config: { harnessId: 'claude-code', model: 'claude-requested', timeoutSeconds: 5 }, idempotencyKey: 'token-env-key' }); + assert.equal(result.outputText, 'token_visible=false;env_entries=true'); +}); diff --git a/src/config.ts b/src/config.ts index bccf843..93f20db 100644 --- a/src/config.ts +++ b/src/config.ts @@ -16,6 +16,7 @@ export interface ForemanConfig { workerTimeoutMs: number; workspaceSourceRepo?: string; workspaceBridgeUrl?: string; + workspaceBridgeToken?: string; workspaceAllowedScope: string[]; validationCommands: Array<{name:string;command:string;args:string[];cwd?:string;network:boolean}>; validationTimeoutMs: number; @@ -44,12 +45,22 @@ function optionalHttpUrl(name: string): string | undefined { return value.toString().replace(/\/$/, ''); } +/** A bearer credential: printable ASCII without spaces, so it is always a valid header value and never needs to appear in an error. */ +function optionalToken(name: string): string | undefined { + const raw = process.env[name]?.trim(); + if (!raw) return undefined; + if (!/^[\x21-\x7e]{1,512}$/.test(raw)) throw new Error(`${name} must be 1-512 printable ASCII characters without spaces`); + return raw; +} + /** Load an optional local .env and validate the supported service configuration. */ export function loadConfig(): ForemanConfig { if (existsSync('.env')) loadEnvFile('.env'); const uhpHarnessId = process.env.UHP_HARNESS_ID?.trim() || undefined; const uhpModel = process.env.UHP_MODEL?.trim() || undefined; if (Boolean(uhpHarnessId) !== Boolean(uhpModel)) throw new Error('UHP_HARNESS_ID and UHP_MODEL must be configured together'); + const workspaceBridgeToken = optionalToken('FOREMAN_WORKSPACE_BRIDGE_TOKEN'); + if (workspaceBridgeToken && !process.env.FOREMAN_WORKSPACE_BRIDGE_URL?.trim()) throw new Error('FOREMAN_WORKSPACE_BRIDGE_TOKEN requires FOREMAN_WORKSPACE_BRIDGE_URL'); let validationCommands: ForemanConfig['validationCommands'] = []; const rawCommands=process.env.FOREMAN_VALIDATION_COMMANDS?.trim(); if(rawCommands){try{const value=JSON.parse(rawCommands);if(!Array.isArray(value))throw new Error();validationCommands=value.map((item:unknown)=>{if(!item||typeof item!=='object')throw new Error();const x=item as Record;if(typeof x.name!=='string'||!x.name.trim()||typeof x.command!=='string'||!x.command||!Array.isArray(x.args)||x.args.some(arg=>typeof arg!=='string')||(x.cwd!==undefined&&typeof x.cwd!=='string')||(x.network!==undefined&&typeof x.network!=='boolean'))throw new Error();return {name:x.name,command:x.command,args:x.args as string[],...(typeof x.cwd==='string'?{cwd:x.cwd}:{}),network:(x.network as boolean|undefined)??defaultNetworkAccess(x.command,x.args as string[])};});}catch{throw new Error('FOREMAN_VALIDATION_COMMANDS must be a JSON array of {name,command,args,cwd?,network?} where network is true or false');}} @@ -73,6 +84,7 @@ export function loadConfig(): ForemanConfig { workerTimeoutMs: integer('FOREMAN_WORKER_TIMEOUT_MS', 600000, 1000, 900000), workspaceSourceRepo: process.env.FOREMAN_WORKSPACE_SOURCE_REPO?.trim() || undefined, workspaceBridgeUrl: process.env.FOREMAN_WORKSPACE_BRIDGE_URL?.trim() || undefined, + workspaceBridgeToken, workspaceAllowedScope, validationCommands, validationTimeoutMs: integer('FOREMAN_VALIDATION_TIMEOUT_MS', 120000, 1, 600000), diff --git a/src/controller.ts b/src/controller.ts index a3d65a0..88c1c3b 100644 --- a/src/controller.ts +++ b/src/controller.ts @@ -53,7 +53,7 @@ export class Controller { private readonly automaticRuns = new Set(); private readonly taskStarts = new Set(); private readonly projectPlannerTurns = new Set(); - private verifiedWorkspaceConfig?:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig}; + private verifiedWorkspaceConfig?:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string}; private discovery?: { version:string; capabilities:Record; harnesses:Array<{id:string;models?:Array<{id:string;available?:boolean}>}>; models?:Array<{id:string;harnessId?:string;available?:boolean}> }; private discoveryError?: string; private memoryServiceState: 'ready'|'degraded'|'unavailable'|'not_configured'; @@ -74,10 +74,10 @@ export class Controller { return this.taskTimeoutSeconds; } async state(): Promise { const state=await this.store.load();for(const project of state.projects){const projectAssignments=[...(project.plannerAssignments??[]),...project.tasks.flatMap(t=>t.runs.flatMap(r=>r.assignments))];project.usage=aggregate(projectAssignments.map(a=>a.usage));project.usageByHarnessModel=groupUsage(projectAssignments);for(const task of project.tasks){const latest=task.runs.at(-1);if(latest?.controller?.active||latest?.status==='awaiting_approval'||latest?.status==='completed'&&latest.promotion?.status!=='applied')task.status='active';else if(latest?.promotion?.status==='applied')task.status='completed';else if(latest?.status==='failed')task.status='blocked';for(const run of task.runs){run.usage=aggregate(run.assignments.map(a=>a.usage));run.usageByRole=Object.fromEntries([...new Set(run.assignments.map(a=>a.roleId))].map(role=>[role,aggregate(run.assignments.filter(a=>a.roleId===role).map(a=>a.usage))]).filter(([,v])=>!!v));run.usageByHarnessModel=groupUsage(run.assignments);}}}for(const role of state.roles){const assignments=[...state.projects.flatMap(p=>p.tasks.flatMap(t=>t.runs.flatMap(r=>r.assignments))),...state.projects.flatMap(p=>p.plannerAssignments??[])].filter(a=>a.roleId===role.id);role.usage=aggregate(assignments.map(a=>a.usage));role.usageByHarnessModel=groupUsage(assignments);}return state; } - configureVerifiedWorkspace(config:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig}):void { if(!config.repoPath||!config.allowedScope.length||!config.commands.length)throw new Error('Workspace repository, allowed scope, and validation commands are required');this.verifiedWorkspaceConfig=structuredClone(config); } + configureVerifiedWorkspace(config:{repoPath:string;allowedScope:string[];commands:ValidationCommand[];bridgeBaseUrl?:string;timeoutMs?:number;maxOutputBytes?:number;sandbox?:ValidationSandboxConfig;bridgeToken?:string}):void { if(!config.repoPath||!config.allowedScope.length||!config.commands.length)throw new Error('Workspace repository, allowed scope, and validation commands are required');this.verifiedWorkspaceConfig=structuredClone(config); } private async workspaceConfigForRun(runId:string){const base=this.verifiedWorkspaceConfig;if(!base)return undefined;const run=this.findRun(await this.store.load(),runId);return {...base,allowedScope:run.allowedScope??base.allowedScope,commands:(run.validationCommands as ValidationCommand[]|undefined)??base.commands};} async pinWorkerBase(runId:string,baseCommit:string):Promise{const config=this.verifiedWorkspaceConfig;if(!config)throw Object.assign(new Error('Verified workspace workflow is not configured'),{statusCode:503});const snapshot=await snapshotGitCommit(config.repoPath,baseCommit);await this.store.mutate(s=>{const run=this.findRun(s,runId);if(run.pinnedBaseCommit&&run.pinnedBaseCommit.toLowerCase()!==baseCommit.toLowerCase())throw Object.assign(new Error('Run base is already pinned'),{statusCode:409});run.pinnedBaseCommit=baseCommit.toLowerCase();const scope=run.allowedScope??config.allowedScope;run.baseFilePaths=snapshot.entries.map(entry=>entry.path).filter(path=>scopeContainedBy(path,scope)).slice(0,200);s.events.push(event('worker.base_pinned','run',runId,{pinnedBaseCommit:run.pinnedBaseCommit,provenance:'fixture_or_operator'}));});} - async prepareWorkerWorkspace(runId:string,baseCommit:string):Promise{const config=this.verifiedWorkspaceConfig;if(!config?.bridgeBaseUrl)throw Object.assign(new Error('Workspace bridge URL is not configured'),{statusCode:503});const current=this.findRun(await this.store.load(),runId);if(current.workspaceId){if(current.pinnedBaseCommit?.toLowerCase()!==baseCommit.toLowerCase())throw Object.assign(new Error('Run workspace is already pinned to a different base'),{statusCode:409});return current;}await this.pinWorkerBase(runId,baseCommit);const seeded=await seedBridgeWorkspace(config.bridgeBaseUrl,baseCommit);await this.store.mutate(s=>{const run=this.findRun(s,runId);if(run.workspaceId&&run.workspaceId!==seeded.workspaceId)throw Object.assign(new Error('Run already has a different workspace'),{statusCode:409});run.workspaceId=seeded.workspaceId;s.events.push(event('worker.workspace_seeded','run',runId,{workspaceId:run.workspaceId,pinnedBaseCommit:run.pinnedBaseCommit}));});return this.findRun(await this.store.load(),runId);} + async prepareWorkerWorkspace(runId:string,baseCommit:string):Promise{const config=this.verifiedWorkspaceConfig;if(!config?.bridgeBaseUrl)throw Object.assign(new Error('Workspace bridge URL is not configured'),{statusCode:503});const current=this.findRun(await this.store.load(),runId);if(current.workspaceId){if(current.pinnedBaseCommit?.toLowerCase()!==baseCommit.toLowerCase())throw Object.assign(new Error('Run workspace is already pinned to a different base'),{statusCode:409});return current;}await this.pinWorkerBase(runId,baseCommit);const seeded=await seedBridgeWorkspace(config.bridgeBaseUrl,baseCommit,{token:config.bridgeToken});await this.store.mutate(s=>{const run=this.findRun(s,runId);if(run.workspaceId&&run.workspaceId!==seeded.workspaceId)throw Object.assign(new Error('Run already has a different workspace'),{statusCode:409});run.workspaceId=seeded.workspaceId;s.events.push(event('worker.workspace_seeded','run',runId,{workspaceId:run.workspaceId,pinnedBaseCommit:run.pinnedBaseCommit}));});return this.findRun(await this.store.load(),runId);} async serviceStatus() { const validationSandbox=await validationSandboxStatus(this.verifiedWorkspaceConfig?.sandbox);const discovery=this.discovery?{version:this.discovery.version,capabilities:this.discovery.capabilities,harnesses:this.discovery.harnesses.map(h=>({id:h.id,models:[...(h.models??[]),...(this.discovery!.models??[]).filter(m=>(m as any).harnessId===h.id)]}))}:undefined;return { uhp:{status:this.discovery?'ready':this.discoveryError?'degraded':this.uhpConfigured?'unavailable':'not_configured',configured:this.uhpConfigured,discovery,error:this.discoveryError},memory:{status:this.memoryServiceState,configured:this.memoryConfigured},validationSandbox }; } async refreshDiscovery() { if(!this.uhpConfigured||!this.uhp.discover){this.discoveryError=undefined;return;} @@ -156,7 +156,7 @@ export class Controller { if(this.projectScopeId&&projectId!==this.projectScopeId)throw Object.assign(new Error('Project not found'),{statusCode:404});const clean=text.trim();if(!clean)throw new Error('Planner message is required');if(Buffer.byteLength(clean,'utf8')>12_000)throw Object.assign(new Error('Planner message exceeds the 12,000 byte limit'),{statusCode:413}); const initial=await this.store.load(),project=must(initial.projects.find(p=>p.id===projectId),'Project'),role=must(initial.roles.find(r=>r.id==='planner'),'Planner');const referenced=[...new Set(clean.match(/\btsk_[0-9a-f-]{36}\b/gi)??[])];if(referenced.length>1)throw Object.assign(new Error('Refer to one task ID per Planner message when steering active work'),{statusCode:422});const steeringTask=referenced.length?project.tasks.find(t=>t.id===referenced[0]):undefined;if(referenced.length&&!steeringTask)throw Object.assign(new Error('Referenced task ID does not belong to this project'),{statusCode:404});const activeSteeringRun=steeringTask?.runs.find(r=>r.controller?.active);const config=resolveRoleConfig(initial.roles,project,undefined,'planner').config??role.config;validateRoleConfig(role,config); await this.store.mutate(s=>{const p=must(s.projects.find(x=>x.id===projectId),'Project');p.plannerMessages??=[];p.plannerMessages.push({id:id('pmsg'),role:'user',text:clean,createdAt:now()});s.events.push(event('project.planner_message_added','project',projectId,{messageCount:p.plannerMessages.length}));}); - let repoAccess:RepoAccessRecord|undefined,readOnlyWorkspaceId:string|undefined,repoAccessText='You cannot read the repository filesystem directly.',digestRequest:DigestRequest|undefined;{const isAgy=config.harnessId==='antigravity-cli';const _ws=this.verifiedWorkspaceConfig;if(_ws?.bridgeBaseUrl&&!isAgy){try{const headCommit=await this.currentHead();const _cached=this.snapshotWorkspaceCache.get(headCommit);if(_cached){readOnlyWorkspaceId=_cached;repoAccess={mode:'snapshot',commit:headCommit};repoAccessText=`You can read (not modify) a snapshot of the repository at commit ${headCommit} in your working directory using Read, Grep and Glob. Look at the relevant code before proposing tasks; cite files you inspected.`;}else{const seeded=await seedBridgeWorkspace(_ws.bridgeBaseUrl,headCommit);this.snapshotWorkspaceCache.set(headCommit,seeded.workspaceId);readOnlyWorkspaceId=seeded.workspaceId;repoAccess={mode:'snapshot',commit:headCommit};repoAccessText=`You can read (not modify) a snapshot of the repository at commit ${headCommit} in your working directory using Read, Grep and Glob. Look at the relevant code before proposing tasks; cite files you inspected.`;}}catch(err){const reason=err instanceof Error?err.message:String(err);const headCommit=await this.currentHead().catch(()=>'');if(headCommit)digestRequest={repoPath:_ws.repoPath,allowedScope:_ws.allowedScope,commit:headCommit,keywords:extractKeywords(clean),reason,fallback:{mode:'digest',commit:'unavailable',reason}};else repoAccess={mode:'digest',commit:'unavailable',reason};}}else if(isAgy&&_ws){const reason='antigravity-cli does not support read-only workspace tools';const headCommit=await this.currentHead().catch(()=>'');if(headCommit)digestRequest={repoPath:_ws.repoPath,allowedScope:_ws.allowedScope,commit:headCommit,keywords:extractKeywords(clean),reason};}} + let repoAccess:RepoAccessRecord|undefined,readOnlyWorkspaceId:string|undefined,repoAccessText='You cannot read the repository filesystem directly.',digestRequest:DigestRequest|undefined;{const isAgy=config.harnessId==='antigravity-cli';const _ws=this.verifiedWorkspaceConfig;if(_ws?.bridgeBaseUrl&&!isAgy){try{const headCommit=await this.currentHead();const _cached=this.snapshotWorkspaceCache.get(headCommit);if(_cached){readOnlyWorkspaceId=_cached;repoAccess={mode:'snapshot',commit:headCommit};repoAccessText=`You can read (not modify) a snapshot of the repository at commit ${headCommit} in your working directory using Read, Grep and Glob. Look at the relevant code before proposing tasks; cite files you inspected.`;}else{const seeded=await seedBridgeWorkspace(_ws.bridgeBaseUrl,headCommit,{token:_ws.bridgeToken});this.snapshotWorkspaceCache.set(headCommit,seeded.workspaceId);readOnlyWorkspaceId=seeded.workspaceId;repoAccess={mode:'snapshot',commit:headCommit};repoAccessText=`You can read (not modify) a snapshot of the repository at commit ${headCommit} in your working directory using Read, Grep and Glob. Look at the relevant code before proposing tasks; cite files you inspected.`;}}catch(err){const reason=err instanceof Error?err.message:String(err);const headCommit=await this.currentHead().catch(()=>'');if(headCommit)digestRequest={repoPath:_ws.repoPath,allowedScope:_ws.allowedScope,commit:headCommit,keywords:extractKeywords(clean),reason,fallback:{mode:'digest',commit:'unavailable',reason}};else repoAccess={mode:'digest',commit:'unavailable',reason};}}else if(isAgy&&_ws){const reason='antigravity-cli does not support read-only workspace tools';const headCommit=await this.currentHead().catch(()=>'');if(headCommit)digestRequest={repoPath:_ws.repoPath,allowedScope:_ws.allowedScope,commit:headCommit,keywords:extractKeywords(clean),reason};}} const snapshot=await this.store.load(),p=must(snapshot.projects.find(x=>x.id===projectId),'Project'),messages=boundedUtf8((p.plannerMessages??[]).slice(-13,-1).map(m=>`${m.role==='user'?'Human':'Planner'}: ${boundedUtf8(m.text,1200)}`).join('\n\n'),7000),taskIndexFull=boundedUtf8(p.tasks.slice(-20).map(t=>`${t.id} | ${displayTaskStatus(t)} | ${boundedUtf8(t.title,100)} | deps: ${boundedUtf8((t.dependsOn??[]).join(', ')||'none',180)}`).join('\n'),3000),allowedScopeText=boundedUtf8((this.verifiedWorkspaceConfig?.allowedScope??[]).join('\n'),4000),steeringContextFull=steeringTask?`\n\nExplicit task reference ${steeringTask.id}: ${boundedUtf8(steeringTask.title,250)}\nGoal: ${boundedUtf8(steeringTask.goal??steeringTask.title,700)}\nAllowed paths: ${boundedUtf8((steeringTask.suggestedAllowedPaths??[]).join(', '),650)}\nValidation criteria: ${boundedUtf8((steeringTask.validationCriteria??[]).join('; '),800)}\nLatest run: ${steeringTask.runs.at(-1)?.status??'not started'}`:'';const assignmentId=id('asgn'),submissionId=id('sub'),key=id('idem'); const instructionsFor=(access:string)=>`You are the continuing Planner for Foreman project ${p.id} (${boundedUtf8(p.name,300)}). You only advise and propose structured tasks. ${access} You cannot dispatch work, approve changes, or change Git. Keep answers useful and concise. If creating tasks, return one JSON object with reply and tasks. Every task must contain title, goal, suggestedAllowedPaths (an array of repository paths), validationCriteria (an array of non-empty strings), optional dependsOn stable task IDs, and optional ref for references between proposed tasks. Do not use a string for validationCriteria. Use this shape:\n{"reply":"...","tasks":[{"title":"...","goal":"...","suggestedAllowedPaths":["nodes/"],"validationCriteria":["The check passes."]}]}\nFor dependencies, use existing stable task IDs from the index above or refs declared by other tasks in this proposal. Suggested paths must stay inside the configured repository allowed scope.\n\nConfigured repository allowed scope:\n${allowedScopeText||'(not configured)'}\n\nCurrent task index (stable IDs):\n`;const fixedTail='\n\nRecent bounded project conversation:\n\n\nCurrent human message:\n'; let taskIndex=taskIndexFull,steeringContext=steeringContextFull,accessText=digestRequest?PLANNER_DIGEST_OMITTED_TEXT:repoAccessText;const headerFor=()=>`${instructionsFor(accessText)}${taskIndex||'(none)'}${steeringContext}`;let header=headerFor(),fixedLength=Buffer.byteLength(header+fixedTail+clean,'utf8'); @@ -591,7 +591,7 @@ export class Controller { const checksHint=(this.verifiedWorkspaceConfig?.commands??[]).map(c=>`${c.name} (${c.command}${c.args.length?` ${c.args.join(' ')}`:''})` ).join(', ')||'none configured'; const workerConfig=run.roleConfigs.worker; const capabilityHint=workerConfig?.harnessId?`\nWorker capabilities: ${workerCapabilityStatement(workerConfig.harnessId)}\n`:''; - let orchRepoAccess:RepoAccessRecord|undefined,orchReadOnlyWorkspaceId:string|undefined,orchRepoAccessText='No tools, file access, or shell commands are available in this context. Respond with text only.',orchDigestRequest:DigestRequest|undefined;{const isAgy=config.harnessId==='antigravity-cli';const _ws=this.verifiedWorkspaceConfig;if(_ws?.bridgeBaseUrl&&!isAgy&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;try{const _cached=this.snapshotWorkspaceCache.get(pinnedCommit);if(_cached){orchReadOnlyWorkspaceId=_cached;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}else{const seeded=await seedBridgeWorkspace(_ws.bridgeBaseUrl,pinnedCommit);this.snapshotWorkspaceCache.set(pinnedCommit,seeded.workspaceId);orchReadOnlyWorkspaceId=seeded.workspaceId;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}}catch(err){const reason=err instanceof Error?err.message:String(err);orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason,fallback:{mode:'digest',commit:pinnedCommit,reason}};}}else if(isAgy&&_ws&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;const reason='antigravity-cli does not support read-only workspace tools';orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason};}} + let orchRepoAccess:RepoAccessRecord|undefined,orchReadOnlyWorkspaceId:string|undefined,orchRepoAccessText='No tools, file access, or shell commands are available in this context. Respond with text only.',orchDigestRequest:DigestRequest|undefined;{const isAgy=config.harnessId==='antigravity-cli';const _ws=this.verifiedWorkspaceConfig;if(_ws?.bridgeBaseUrl&&!isAgy&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;try{const _cached=this.snapshotWorkspaceCache.get(pinnedCommit);if(_cached){orchReadOnlyWorkspaceId=_cached;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}else{const seeded=await seedBridgeWorkspace(_ws.bridgeBaseUrl,pinnedCommit,{token:_ws.bridgeToken});this.snapshotWorkspaceCache.set(pinnedCommit,seeded.workspaceId);orchReadOnlyWorkspaceId=seeded.workspaceId;orchRepoAccess={mode:'snapshot',commit:pinnedCommit};orchRepoAccessText=`You can read (not modify) a snapshot of the repository at commit ${pinnedCommit} in your working directory using Read, Grep and Glob. Inspect relevant code before proposing the Worker task; cite files you read.`;}}catch(err){const reason=err instanceof Error?err.message:String(err);orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason,fallback:{mode:'digest',commit:pinnedCommit,reason}};}}else if(isAgy&&_ws&&run.pinnedBaseCommit){const pinnedCommit=run.pinnedBaseCommit;const reason='antigravity-cli does not support read-only workspace tools';orchDigestRequest={repoPath:_ws.repoPath,allowedScope:run.allowedScope??_ws.allowedScope,commit:pinnedCommit,keywords:extractKeywords(clean),reason};}} // Foreman-supplied context is trimmed to fit; only the instructions, the operator note and the task title are mandatory. let goal=task.goal??task.title,criteria=(task.validationCriteria??[]).join('; '),planner=plannerContext,paths=pathList,checks=checksHint,scopeText=allowedScope.join(', '),guidanceCount=guidanceLines.length; const guidanceText=()=>`${guidanceLines.slice(0,guidanceCount).join('\n')}${guidanceCount0){await overlayBridgeWorkspace(config.bridgeBaseUrl,seeded.workspaceId,overlayEntries);overlaidPriorAssignmentId=run.workerEvidence.workerAssignmentId;} + if(overlayEntries.length>0){await overlayBridgeWorkspace(config.bridgeBaseUrl,seeded.workspaceId,overlayEntries,{token:config.bridgeToken});overlaidPriorAssignmentId=run.workerEvidence.workerAssignmentId;} } await this.store.mutate(s=>{const current=this.findRun(s,runId),stored=current.workerProposal,kind=stored&&this.workerRetryKind(current,stored),prior=stored?.workerAssignmentId&¤t.assignments.find(a=>a.id===stored.workerAssignmentId);if(!stored||stored.id!==proposalId||!['dispatched','dispatching','pending'].includes(stored.status)||!kind||!prior)throw Object.assign(new Error('Worker retry eligibility changed while the fresh workspace was being seeded'),{statusCode:409});stored.workerAssignmentHistoryIds??=current.assignments.filter(a=>a.roleId==='worker').map(a=>a.id);current.workspaceId=seeded.workspaceId;s.events.push(event('worker.retry_workspace_seeded','run',runId,{proposalId,workspaceId:seeded.workspaceId,pinnedBaseCommit:current.pinnedBaseCommit,priorWorkerAssignmentId:prior.id,retryKind:kind,...(overlaidPriorAssignmentId?{overlaidPriorAssignmentId}:{})}));stored.status='pending';s.events.push(event('orchestrator.worker_retry_authorized','run',runId,{proposalId,priorWorkerAssignmentId:prior.id,pinnedBaseCommit:current.pinnedBaseCommit,workspaceId:seeded.workspaceId,retryKind:kind,...(overlaidPriorAssignmentId?{overlaidPriorAssignmentId}:{})}));}); return await this.dispatchWorkerProposalInternal(runId,proposalId,INTERNAL_WORKER_RETRY); @@ -791,7 +791,7 @@ export class Controller { const recoveredInvocation=cliInvocationFrom(metadata.cli_invocation),isCli=isCliHarness(worker.requestedConfig.harnessId); if(isCli&&(metadata.actual_model_status!==worker.actualModelStatus||canonicalJson(recoveredInvocation)!==canonicalJson(worker.cliInvocation)||metadata.harness_id!==worker.requestedConfig.harnessId||!explicitModelArg(recoveredInvocation?.args,worker.requestedModel??'')))throw Object.assign(new Error('Retrieved CLI response does not match the saved model-reporting and invocation evidence'),{statusCode:409}); await this.store.mutate(s=>{const a=this.findAssignment(s,worker.id).a;if(typeof response.model==='string'){a.actualConfig={harnessId:String(metadata.harness_id??a.actualConfig?.harnessId??a.requestedConfig.harnessId),model:response.model};a.actualModelStatus='observed';}else if(metadata.actual_model_status==='unavailable')a.actualModelStatus='unavailable';if(isCli)a.cliInvocation=recoveredInvocation;a.requestedModel=String(requestedModel);const usage=usageFromResponse(response.usage);if(usage)a.usage={...a.usage,...usage,measured:true};}); - const envelope=await fetchBridgeSnapshot(config.bridgeBaseUrl,run.workspaceId,run.pinnedBaseCommit); + const envelope=await fetchBridgeSnapshot(config.bridgeBaseUrl,run.workspaceId,run.pinnedBaseCommit,{token:config.bridgeToken}); return this.processWorkerOutput(runId,workerAssignmentId,envelope,'bridge_snapshot'); } /** Deterministic fixture path. Never exposed by the production HTTP API. */ diff --git a/src/local-bridge.ts b/src/local-bridge.ts index 8bfdf4b..fc03c36 100644 --- a/src/local-bridge.ts +++ b/src/local-bridge.ts @@ -1,12 +1,25 @@ import { createServer } from 'node:net'; -import { createHash } from 'node:crypto'; +import { createHash, randomBytes } from 'node:crypto'; import { spawn, type ChildProcess } from 'node:child_process'; -import { mkdir, realpath, access, stat } from 'node:fs/promises'; +import { appendFile, mkdir, open, realpath, rename, access, stat } from 'node:fs/promises'; import { constants as fsConstants } from 'node:fs'; import { homedir, tmpdir } from 'node:os'; import { dirname, join, resolve } from 'node:path'; import { fileURLToPath } from 'node:url'; +/** Restart policy for a bridge that exits after it was ready. Every value has a production default; tests inject small ones. */ +export interface LocalBridgeSupervision { + /** Delay before the first restart; it doubles for each further consecutive restart. Default 1 s. */ + restartBaseDelayMs?: number; + /** Upper bound of the restart delay. Default 30 s. */ + restartMaxDelayMs?: number; + /** Give up after this many consecutive runs that each ended within `stableAfterMs` of starting. Default 5. */ + maxFastFailures?: number; + /** A run that lasted at least this long was healthy: its crash restarts the failure count and the backoff. Default 60 s. */ + stableAfterMs?: number; + /** Rotate bridge.log to bridge.log.1 (keeping one old file) at (re)start once it is larger than this. Default 5 MiB. */ + logMaxBytes?: number; +} export interface LocalBridgeOptions { bridgeScript?: string; nodePath?: string; @@ -14,13 +27,71 @@ export interface LocalBridgeOptions { homeDir?: string; tempDir?: string; startupTimeoutMs?: number; + /** Pause between readiness probes while the bridge is starting. Default 100 ms. */ + startupPollMs?: number; fetch?: typeof fetch; spawn?: typeof spawn; allocatePort?: () => Promise; + supervision?: LocalBridgeSupervision; + /** Called when the bridge crashes, is being restarted, comes back, or is given up on. Never receives the token. */ + onHealthChange?: (health: LocalBridgeHealth) => void; +} +/** A ready bridge. `token` is the per-start bearer token; it is deliberately not enumerable, so serialising or logging a status can never leak it. */ +export interface LocalBridgeStatus { repoPath: string; baseUrl: string; readonly token: string; } +export interface LocalBridgeExit { + code: number | null; + signal: string | null; + at: string; + /** How long that run lasted. */ + uptimeMs: number; + /** Set when the run did not exit by itself, e.g. it never became ready and was stopped. */ + reason?: string; +} +export interface LocalBridgeHealth { + /** `ready`: serving. `restarting`: crashed, a restart is scheduled or starting. `unavailable`: stopped, never started, or given up on. */ + state: 'ready' | 'restarting' | 'unavailable'; + /** Where the bridge's stdout and stderr go, once a bridge has been started. */ + logPath?: string; + lastExit?: LocalBridgeExit; + /** Restart attempts since the bridge was started. */ + restarts: number; + nextRestartAt?: string; + /** Delay chosen for the pending restart. */ + restartDelayMs?: number; + message?: string; } -export interface LocalBridgeStatus { repoPath: string; baseUrl: string; } const SAFE_ENV = ['PATH','HOME','USER','LOGNAME','LANG','LC_ALL','TERM','TMPDIR','TMP','TEMP','XDG_RUNTIME_DIR','DBUS_SESSION_BUS_ADDRESS','HTTP_PROXY','HTTPS_PROXY','ALL_PROXY','NO_PROXY','SSL_CERT_FILE','SSL_CERT_DIR','NODE_EXTRA_CA_CERTS']; +const STOP_GRACE_MS = 1_500; + +/** Everything a restart needs to bring the same bridge back: same port, same token, same state. */ +interface Session { + readonly repo: string; + readonly instanceId: string; + readonly script: string; + readonly port: number; + readonly token: string; + readonly baseUrl: string; + readonly env: NodeJS.ProcessEnv; + readonly logPath: string; + readonly status: LocalBridgeStatus; + /** Matches `LocalBridge.generation` while this session is current; any stop() or new start() invalidates it. */ + readonly generation: number; +} +interface Run { child: ChildProcess; startedAt: number; ready: boolean; exit?: { code: number | null; signal: string | null }; error?: Error; } + +/** + * A `fetch` that adds `Authorization: Bearer ` to every request it makes unless the caller already set one. + * UhpClient deliberately sends no credential on UHP discovery (`GET /v1/uhp`), but the bridge protects that route + * too, so Foreman's clients for a token-protected bridge use this as their `fetch`. + */ +export function bearerFetch(token: string, base: typeof fetch = globalThis.fetch): typeof fetch { + return (input, init) => { + const headers = new Headers(init?.headers ?? (input instanceof Request ? input.headers : undefined)); + if (!headers.has('authorization')) headers.set('authorization', `Bearer ${token}`); + return base(input, { ...init, headers }); + }; +} async function freePort(): Promise { return new Promise((resolvePort, reject) => { @@ -44,12 +115,32 @@ async function isGitRepo(path: string): Promise { return repo; } +/** Keep one old log: past `maxBytes` the current log becomes `.1`, replacing any earlier one. */ +async function rotateLog(path: string, maxBytes: number): Promise { + try { if ((await stat(path)).size > maxBytes) await rename(path, `${path}.1`); } + catch { /* no log yet, or rotation is best effort */ } +} +function describeExit(exit: Pick): string { + return exit.signal ? `signal ${exit.signal}` : exit.code !== null ? `exit code ${exit.code}` : exit.reason ?? 'unknown exit'; +} +const alive = (child: ChildProcess): boolean => child.exitCode === null && (child.signalCode ?? null) === null; + export class LocalBridge { private child?: ChildProcess; private active?: LocalBridgeStatus; private activeInstanceId?: string; private starting?: Promise; - private readonly options: Required> & { allocatePort: () => Promise }; + private session?: Session; + private generation = 0; + private logWrites: Promise = Promise.resolve(); + private restartTimer?: ReturnType; + private fastFailures = 0; + private backoffAttempts = 0; + private restarts = 0; + private lastExit?: LocalBridgeExit; + private restartAt?: { at: number; delayMs: number }; + private givenUp?: string; + private readonly options: Required> & { allocatePort: () => Promise; supervision: Required; onHealthChange?: (health: LocalBridgeHealth) => void }; constructor(options: LocalBridgeOptions = {}) { const h = options.homeDir ?? homedir(); @@ -60,14 +151,36 @@ export class LocalBridge { homeDir: h, tempDir: options.tempDir ?? tmpdir(), startupTimeoutMs: options.startupTimeoutMs ?? 8_000, + startupPollMs: options.startupPollMs ?? 100, fetch: options.fetch ?? globalThis.fetch, spawn: options.spawn ?? spawn, allocatePort: options.allocatePort ?? freePort, + supervision: { + restartBaseDelayMs: options.supervision?.restartBaseDelayMs ?? 1_000, + restartMaxDelayMs: options.supervision?.restartMaxDelayMs ?? 30_000, + maxFastFailures: options.supervision?.maxFastFailures ?? 5, + stableAfterMs: options.supervision?.stableAfterMs ?? 60_000, + logMaxBytes: options.supervision?.logMaxBytes ?? 5 * 1024 * 1024, + }, + ...(options.onHealthChange ? { onHealthChange: options.onHealthChange } : {}), }; } + /** The bridge, only while it is ready to serve. During a restart or after giving up this is undefined; see `health`. */ get status(): LocalBridgeStatus | undefined { return this.active; } + get health(): LocalBridgeHealth { + const state = this.active ? 'ready' : this.restartAt || this.restartTimer ? 'restarting' : 'unavailable'; + return { + state, + ...(this.session ? { logPath: this.session.logPath } : {}), + ...(this.lastExit ? { lastExit: { ...this.lastExit } } : {}), + restarts: this.restarts, + ...(state === 'restarting' && this.restartAt ? { nextRestartAt: new Date(this.restartAt.at).toISOString(), restartDelayMs: this.restartAt.delayMs, message: `Local bridge ended (${this.lastExit ? describeExit(this.lastExit) : 'unknown exit'}); restarting in ${this.restartAt.delayMs} ms; see ${this.session?.logPath}` } : {}), + ...(state === 'unavailable' && this.givenUp ? { message: this.givenUp } : {}), + }; + } + /** Compute the key used for a given repoPath + instanceId pair. */ bridgeKey(repoPath: string, instanceId: string): string { return createHash('sha256').update(`${repoPath}\0${instanceId}`).digest('hex').slice(0, 20); @@ -95,6 +208,7 @@ export class LocalBridge { await this.stop(); const script = await realpath(this.options.bridgeScript); const port = await this.options.allocatePort(); + const token = randomBytes(32).toString('hex'); const key = createHash('sha256').update(`${repo}\0${instanceId}`).digest('hex').slice(0, 20); const stateDir = resolve(this.options.dataDir, key); const statePath = join(stateDir, 'uhp-state.json'); @@ -105,6 +219,7 @@ export class LocalBridge { for (const key of SAFE_ENV) if (typeof process.env[key] === 'string') env[key] = process.env[key]; env.HOME = this.options.homeDir; env.LOCAL_CLI_UHP_PORT = String(port); + env.LOCAL_CLI_UHP_TOKEN = token; env.LOCAL_CLI_UHP_STATE = statePath; env.LOCAL_CLI_UHP_WORK = workPath; env.LOCAL_CLI_UHP_SOURCE_REPO = repo; @@ -115,49 +230,164 @@ export class LocalBridge { env.AGY_CONFIG_DIR = join(this.options.homeDir, '.gemini', 'antigravity-cli'); env.AGY_MODEL = 'gemini-3.8-flash-low'; env.AGY_WORKER_EFFORT = 'low'; - const child = this.options.spawn(this.options.nodePath, [script], { cwd: dirname(script), env, stdio: ['ignore','ignore','ignore'], shell: false, windowsHide: true }); - this.child = child; const baseUrl = `http://127.0.0.1:${port}`; - const deadline = Date.now() + this.options.startupTimeoutMs; - let startupError: Error | undefined; - const onError = (error: Error) => { startupError = error; }; - child.once('error', onError); + const status = Object.defineProperty({ repoPath: repo, baseUrl } as LocalBridgeStatus, 'token', { value: token, enumerable: false }); + this.fastFailures = 0; this.backoffAttempts = 0; this.restarts = 0; this.lastExit = undefined; this.givenUp = undefined; + const session: Session = { repo, instanceId, script, port, token, baseUrl, env, logPath: join(stateDir, 'bridge.log'), status, generation: ++this.generation }; + this.session = session; try { - while (Date.now() < deadline) { - if (startupError) throw startupError; - if (child.exitCode !== null) throw new Error(`Local bridge exited during startup (${child.exitCode})`); - try { - const response = await this.options.fetch(`${baseUrl}/v1/uhp`, { signal: AbortSignal.timeout(350) }); - if (response.ok) { - const discovery = await response.json() as { protocol?: string; implementation?: { name?: string } }; - if (discovery.protocol === 'uhp' && discovery.implementation?.name === 'local-cli-uhp') { - const status = { repoPath: repo, baseUrl }; - this.active = status; - this.activeInstanceId = instanceId; - return status; - } - } - } catch { /* server is still starting */ } - await new Promise(resolveDelay => setTimeout(resolveDelay, 100)); - } - throw new Error('Local bridge did not become ready'); + await this.bringUp(session, {}); + return status; } catch (error) { await this.stop(); throw error; - } finally { child.removeListener('error', onError); } + } + } + + /** Spawn the bridge, wait until it answers authenticated discovery, then publish it as the active status. */ + private async bringUp(session: Session, attempt: { run?: Run }): Promise { + const run = attempt.run = await this.launch(session); + await this.waitUntilReady(session, run); + // Log before publishing, so the bridge turns ready and its health is reported with no await in between. + await this.note(session, `bridge ready on ${session.baseUrl} (pid ${run.child.pid ?? 'unknown'})`); + this.assertCurrent(session); + if (run.exit) throw new Error(`Local bridge exited during startup (${describeExit(run.exit)}); see ${session.logPath}`); + run.ready = true; + this.active = session.status; + this.activeInstanceId = session.instanceId; + } + + private async launch(session: Session): Promise { + // Let notes about the previous run land before rotation, so none can end up in the new log or be lost with the old one. + await this.logWrites; + await rotateLog(session.logPath, this.options.supervision.logMaxBytes); + const log = await open(session.logPath, 'a', 0o600); + try { + await log.chmod(0o600); + this.assertCurrent(session); + const child = this.options.spawn(this.options.nodePath, [session.script], { cwd: dirname(session.script), env: session.env, stdio: ['ignore', log.fd, log.fd], shell: false, windowsHide: true }); + const run: Run = { child, startedAt: Date.now(), ready: false }; + this.child = child; + child.once('exit', (code, signal) => { run.exit = { code, signal: signal ?? null }; this.onExit(session, run); }); + // A spawn failure reports only 'error'. Keep a listener for the child's whole life: an unhandled 'error' would crash Foreman. + child.on('error', error => { if (child.pid === undefined && !run.exit) { run.error = error; run.exit = { code: null, signal: null }; this.onExit(session, run); } }); + await this.note(session, `starting bridge on ${session.baseUrl} (pid ${child.pid ?? 'unknown'})`); + return run; + } finally { await log.close().catch(() => undefined); } + } + + private async waitUntilReady(session: Session, run: Run): Promise { + const deadline = Date.now() + this.options.startupTimeoutMs; + while (Date.now() < deadline) { + this.assertCurrent(session); + if (run.error) throw run.error; + if (run.exit) throw new Error(`Local bridge exited during startup (${describeExit(run.exit)}); see ${session.logPath}`); + let rejected: number | undefined; + try { + const response = await this.options.fetch(`${session.baseUrl}/v1/uhp`, { headers: { authorization: `Bearer ${session.token}` }, signal: AbortSignal.timeout(350) }); + if (response.status === 401 || response.status === 403) rejected = response.status; + else if (response.ok) { + const discovery = await response.json() as { protocol?: string; implementation?: { name?: string } }; + if (discovery.protocol === 'uhp' && discovery.implementation?.name === 'local-cli-uhp') return; + } + } catch { /* server is still starting */ } + // The bridge was started with this token, so a rejection means something else owns the port. + if (rejected !== undefined) throw new Error(`Local bridge port ${session.port} answered the readiness check with ${rejected}; another process may own it`); + await new Promise(resolveDelay => setTimeout(resolveDelay, this.options.startupPollMs)); + } + throw new Error(`Local bridge did not become ready; see ${session.logPath}`); + } + + private assertCurrent(session: Session): void { + if (session.generation !== this.generation) throw new Error('Local bridge was stopped during startup'); + } + + /** A child ended. Only an exit after readiness, of the current child, is an unexpected crash; startup failures are reported by the startup poll. */ + private onExit(session: Session, run: Run): void { + if (session.generation !== this.generation || this.child !== run.child || !run.ready) return; + this.child = undefined; + this.active = undefined; + this.afterFailure(session, { code: run.exit?.code ?? null, signal: run.exit?.signal ?? null, at: new Date().toISOString(), uptimeMs: Date.now() - run.startedAt }, false); + } + + /** Count a finished run and either schedule the next restart or give up. */ + private afterFailure(session: Session, exit: LocalBridgeExit, startupFailure: boolean): void { + const policy = this.options.supervision; + this.lastExit = exit; + if (!startupFailure && exit.uptimeMs >= policy.stableAfterMs) { this.fastFailures = 0; this.backoffAttempts = 0; } + else this.fastFailures++; + this.note(session, `bridge ended (${describeExit(exit)}) after ${exit.uptimeMs} ms`); + if (this.fastFailures >= policy.maxFastFailures) { + this.restartAt = undefined; + this.givenUp = `Local bridge stopped after ${this.fastFailures} consecutive failures within ${Math.round(policy.stableAfterMs / 1000)} s of starting (last exit: ${describeExit(exit)}); see ${session.logPath}`; + this.note(session, `giving up: ${this.givenUp}`); + this.notify(); + return; + } + const delayMs = Math.min(policy.restartBaseDelayMs * 2 ** this.backoffAttempts, policy.restartMaxDelayMs); + this.backoffAttempts++; + this.restartAt = { at: Date.now() + delayMs, delayMs }; + this.note(session, `restarting in ${delayMs} ms`); + this.restartTimer = setTimeout(() => { this.restartTimer = undefined; void this.restart(session).catch(() => undefined); }, delayMs); + this.restartTimer.unref?.(); + this.notify(); + } + + /** Bring the same bridge back on the same port with the same token, so every URL and credential already handed out stays valid. */ + private async restart(session: Session): Promise { + if (session.generation !== this.generation) return; + this.restarts++; + const attempt: { run?: Run } = {}; + try { + await this.bringUp(session, attempt); + } catch (error) { + if (session.generation !== this.generation) return; // a deliberate stop() owns the cleanup + const run = attempt.run; + const own = run && !run.error ? run.exit : undefined; // how it ended by itself, before we stop it below + if (run) await this.terminate(run.child); + if (session.generation !== this.generation) return; + this.child = undefined; + // An unbindable port makes the bridge exit at once; a bridge that never answers is stopped. Both are failed attempts. + this.afterFailure(session, { code: own?.code ?? null, signal: own?.signal ?? null, at: new Date().toISOString(), uptimeMs: run ? Date.now() - run.startedAt : 0, ...(own ? {} : { reason: error instanceof Error ? error.message : 'restart failed' }) }, true); + return; + } + this.restartAt = undefined; + this.notify(); } async stop(): Promise { + this.generation++; + if (this.restartTimer) clearTimeout(this.restartTimer); + this.restartTimer = undefined; + this.restartAt = undefined; const child = this.child; this.child = undefined; this.active = undefined; this.activeInstanceId = undefined; - if (!child || child.exitCode !== null || child.killed) return; + if (child) await this.terminate(child); + } + + private async terminate(child: ChildProcess): Promise { + if (!alive(child)) return; const exited = new Promise(resolveExit => child.once('exit', () => resolveExit())); child.kill('SIGTERM'); - const timeout = new Promise(resolveTimeout => setTimeout(resolveTimeout, 1_500)); + let timer: ReturnType | undefined; + const timeout = new Promise(resolveTimeout => { timer = setTimeout(resolveTimeout, STOP_GRACE_MS); }); await Promise.race([exited, timeout]); - if (child.exitCode === null) child.kill('SIGKILL'); + clearTimeout(timer); + if (alive(child)) child.kill('SIGKILL'); + } + + /** Append a Foreman line to the bridge log. Writes are queued so lines keep their order; callers on the startup path await it. */ + private note(session: Session, message: string): Promise { + const line = `[foreman ${new Date().toISOString()}] ${message}\n`; + this.logWrites = this.logWrites.then(() => appendFile(session.logPath, line, { mode: 0o600 })).catch(() => undefined); + return this.logWrites; + } + private notify(): void { + const health = this.health; + if (health.state === 'unavailable' && !this.givenUp) return; + try { this.options.onHealthChange?.(health); } catch { /* a listener must not break supervision */ } } } diff --git a/src/server.ts b/src/server.ts index fb1b4e7..e695860 100644 --- a/src/server.ts +++ b/src/server.ts @@ -9,7 +9,7 @@ import { Controller } from './controller.js'; import { JsonStore } from './store.js'; import { UhpClient } from './uhp.js'; import { HindsightClient } from './hindsight.js'; -import { LocalBridge } from './local-bridge.js'; +import { LocalBridge, bearerFetch } from './local-bridge.js'; import { browseRepositories, repositoryName } from './repository-browser.js'; import { inspectRepository } from './repository-inspector.js'; import { deleteWorkspaceSetup, findSavedProjectForRepository, loadWorkspaceSetup, saveWorkspaceSetup, validateWorkspaceSetup } from './workspace-setup.js'; @@ -24,21 +24,21 @@ if(workspacePolicyConfigured&&(!config.workspaceSourceRepo||!config.workspaceBri const store=new JsonStore(resolve(config.dataDir,'state.json')); const github=new GitHubIntegration(store,config.dataDir); const uhpToken=process.env.UHP_TOKEN; -const uhp=config.uhpBaseUrl ? new UhpClient({baseUrl:config.uhpBaseUrl,...(uhpToken?{token:uhpToken}:{}),harnessId:config.uhpHarnessId,model:config.uhpModel,timeoutMs:Math.max(config.requestTimeoutMs,45_000)}) : { +const uhp=config.uhpBaseUrl ? new UhpClient({baseUrl:config.uhpBaseUrl,...(uhpToken?{token:uhpToken,fetch:bearerFetch(uhpToken)}:{}),harnessId:config.uhpHarnessId,model:config.uhpModel,timeoutMs:Math.max(config.requestTimeoutMs,45_000)}) : { async submit():Promise{throw new Error('UHP is not configured (set UHP_BASE_URL)');}, async cancel():Promise{throw new Error('UHP is not configured (set UHP_BASE_URL)');} }; const hindsight=config.hindsightBaseUrl?new HindsightClient({baseUrl:config.hindsightBaseUrl,token:process.env.HINDSIGHT_TOKEN}):undefined; const controller=new Controller(store,uhp,!!config.hindsightBaseUrl,!!config.uhpBaseUrl,config.uhpHarnessId&&config.uhpModel?{harnessId:config.uhpHarnessId,model:config.uhpModel}:undefined,hindsight,Math.ceil(config.taskTimeoutMs/1000),Math.ceil(config.workerTimeoutMs/1000),300); -if(config.workspaceSourceRepo&&config.workspaceAllowedScope.length&&config.validationCommands.length)controller.configureVerifiedWorkspace({repoPath:config.workspaceSourceRepo,allowedScope:config.workspaceAllowedScope,commands:config.validationCommands,bridgeBaseUrl:config.workspaceBridgeUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox}); +if(config.workspaceSourceRepo&&config.workspaceAllowedScope.length&&config.validationCommands.length)controller.configureVerifiedWorkspace({repoPath:config.workspaceSourceRepo,allowedScope:config.workspaceAllowedScope,commands:config.validationCommands,bridgeBaseUrl:config.workspaceBridgeUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox,bridgeToken:config.workspaceBridgeToken}); const projectControllers=new Map>}>(); const createProjectRuntime=async(projectId:string,workspace:Awaited>)=>{ - const bridge=new LocalBridge({dataDir:resolve(config.dataDir,'local-bridges')}); + const bridge=new LocalBridge({dataDir:resolve(config.dataDir,'local-bridges'),onHealthChange:health=>process.stderr.write(`Local bridge for ${projectId} is ${health.state}${health.message?`: ${health.message}`:''}\n`)}); try { const status=await bridge.start(workspace.repoPath,projectId); - const projectUhp=new UhpClient({baseUrl:status.baseUrl,timeoutMs:Math.max(config.requestTimeoutMs,45_000)}); + const projectUhp=new UhpClient({baseUrl:status.baseUrl,token:status.token,fetch:bearerFetch(status.token),timeoutMs:Math.max(config.requestTimeoutMs,45_000)}); const scoped=new Controller(store,projectUhp,!!config.hindsightBaseUrl,true,undefined,hindsight,Math.ceil(config.taskTimeoutMs/1000),Math.ceil(config.workerTimeoutMs/1000),300,projectId); - scoped.configureVerifiedWorkspace({repoPath:workspace.repoPath,allowedScope:workspace.allowedScope,commands:workspace.validationCommands,bridgeBaseUrl:status.baseUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox}); + scoped.configureVerifiedWorkspace({repoPath:workspace.repoPath,allowedScope:workspace.allowedScope,commands:workspace.validationCommands,bridgeBaseUrl:status.baseUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox,bridgeToken:status.token}); projectControllers.set(projectId,{controller:scoped,bridge,workspace}); return {scoped,bridge,status}; } catch(error) { await bridge.stop(); throw error; } @@ -116,9 +116,9 @@ const server=createServer(async(req,res)=>{ const runtime=projectControllers.get(projectId); const unavailable={status:'unavailable' as const}; const fallback={harnesses:['claude-code','codex-cli','antigravity-cli'].map(harnessId=>({harnessId,status:'unavailable' as const,windows:{fiveHour:unavailable,weekly:unavailable}}))}; - const baseUrl=runtime?.bridge.status?.baseUrl; - if(!baseUrl){json(res,200,fallback);return;} - try{json(res,200,await new UhpClient({baseUrl,timeoutMs:15_000}).usage());return;} + const bridgeStatus=runtime?.bridge.status; + if(!bridgeStatus){json(res,200,{...fallback,bridgeStatus:runtime?.bridge.health.state??'unavailable'});return;} + try{json(res,200,await new UhpClient({baseUrl:bridgeStatus.baseUrl,token:bridgeStatus.token,fetch:bearerFetch(bridgeStatus.token),timeoutMs:15_000}).usage());return;} catch{json(res,200,fallback);return;} } const activeController=await controllerForPath(path,url.searchParams,req.method); @@ -129,9 +129,10 @@ const server=createServer(async(req,res)=>{ const savedState=await store.load(),existingId=await findSavedProjectForRepository(config.dataDir,savedState.projects,workspace.repoPath); if(existingId){ let runtime=projectControllers.get(existingId); + if(runtime&&runtime.bridge.health.state==='unavailable'){await runtime.bridge.stop();projectControllers.delete(existingId);runtime=undefined;} if(runtime){ - const bridgeBaseUrl=runtime.bridge.status?.baseUrl;if(!bridgeBaseUrl)throw Object.assign(new Error('Local repository bridge is unavailable for this project'),{statusCode:503}); - runtime.controller.configureVerifiedWorkspace({repoPath:workspace.repoPath,allowedScope:workspace.allowedScope,commands:workspace.validationCommands,bridgeBaseUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox}); + const bridgeBaseUrl=runtime.bridge.status?.baseUrl;if(!bridgeBaseUrl)throw Object.assign(new Error(`Local repository bridge is ${runtime.bridge.health.state}; retry shortly`),{statusCode:503}); + runtime.controller.configureVerifiedWorkspace({repoPath:workspace.repoPath,allowedScope:workspace.allowedScope,commands:workspace.validationCommands,bridgeBaseUrl,timeoutMs:config.validationTimeoutMs,maxOutputBytes:config.validationMaxOutputBytes,sandbox:config.validationSandbox,bridgeToken:runtime.bridge.status?.token}); const saved=await saveWorkspaceSetup(config.dataDir,existingId,{repoPath:workspace.repoPath,allowedScope:workspace.allowedScope,validationCommands:workspace.validationCommands}); projectControllers.set(existingId,{...runtime,workspace:saved}); }else{ @@ -147,7 +148,7 @@ const server=createServer(async(req,res)=>{ catch(error){projectControllers.delete(projectId);await bridge.stop();throw error;} } const setupMatch=path.match(/^\/api\/projects\/([^/]+)\/workspace-setup$/); - if(req.method==='GET'&&setupMatch){const projectId=decodeURIComponent(setupMatch[1]!);let workspace;try{workspace=await loadWorkspaceSetup(config.dataDir,projectId);}catch{throw Object.assign(new Error('Saved repository is unavailable; reopen the repository to continue'),{statusCode:503});}if(!workspace)throw Object.assign(new Error('Workspace setup not found'),{statusCode:404});json(res,200,{...workspace,bridgeStatus:projectControllers.get(projectId)?.bridge.status?'ready':'unavailable'});return;} + if(req.method==='GET'&&setupMatch){const projectId=decodeURIComponent(setupMatch[1]!);let workspace;try{workspace=await loadWorkspaceSetup(config.dataDir,projectId);}catch{throw Object.assign(new Error('Saved repository is unavailable; reopen the repository to continue'),{statusCode:503});}if(!workspace)throw Object.assign(new Error('Workspace setup not found'),{statusCode:404});const bridgeHealth=projectControllers.get(projectId)?.bridge.health;json(res,200,{...workspace,bridgeStatus:bridgeHealth?.state??'unavailable',...(bridgeHealth?{bridgeHealth}:{})});return;} if(req.method==='GET'&&path==='/api/state'){json(res,200,stateForUi(await activeController.state()));return;} if(req.method==='GET'&&path==='/api/status'){json(res,200,await activeController.serviceStatus());return;} if(req.method==='GET'&&path==='/api/events'){ diff --git a/src/verified-workspace.ts b/src/verified-workspace.ts index 6410367..7a8689f 100644 --- a/src/verified-workspace.ts +++ b/src/verified-workspace.ts @@ -19,14 +19,22 @@ export interface ValidationObservation { name: string; command: string; args: st export interface ControllerValidationEvidence { passed: boolean; checks: ValidationObservation[] } const SHA=/^(?:[a-f0-9]{40}|[a-f0-9]{64})$/i; +/** Per-call bridge options. A bare number is the timeout in ms; `token` is the bridge's bearer token, sent only to the validated loopback bridge URL and never echoed in errors. */ +export interface BridgeRequestOptions { timeoutMs?: number; token?: string } +export interface BridgeSnapshotOptions extends BridgeRequestOptions { maxResponseBytes?: number } +function bridgeTimeout(options:number|BridgeRequestOptions):number{return typeof options==='number'?options:options.timeoutMs??15_000;} +function bridgeHeaders(options:number|BridgeRequestOptions,extra:Record={}):Record{const token=typeof options==='number'?undefined:options.token;return {...(token?{authorization:`Bearer ${token}`}:{}),...extra};} +function bridgeFailure(what:string,status:number):Error{return new Error(`Workspace bridge ${what} (${status})${status===401?': the bridge requires a valid bearer token':status===403?': the bridge refused the request (Host check)':''}`);} + /** Fetch a bridge snapshot from an explicitly loopback-only bridge endpoint. */ -export async function fetchBridgeSnapshot(baseUrl:string,workspaceId:string,expectedBase:string,timeoutMs=15_000,maxResponseBytes=300*1024*1024):Promise{ +export async function fetchBridgeSnapshot(baseUrl:string,workspaceId:string,expectedBase:string,options:number|BridgeSnapshotOptions=15_000,maxResponseBytesArg=300*1024*1024):Promise{ + const timeoutMs=bridgeTimeout(options),maxResponseBytes=typeof options==='object'&&options.maxResponseBytes!==undefined?options.maxResponseBytes:maxResponseBytesArg; const base=new URL(baseUrl); if(!['http:','https:'].includes(base.protocol)||base.username||base.password||base.search||base.hash||!['localhost','127.0.0.1','[::1]','::1'].includes(base.hostname)) throw new Error('Workspace bridge URL must be loopback-only and contain no credentials, query, or fragment'); if(!workspaceId||workspaceId.includes('/')||workspaceId.includes('\\'))throw new Error('Invalid bridge workspace ID'); const endpoint=new URL(`/extensions/foreman-workspace/v1/workspaces/${encodeURIComponent(workspaceId)}/snapshot`,base); - const response=await fetch(endpoint,{signal:AbortSignal.timeout(timeoutMs)}); - if(!response.ok)throw new Error(`Workspace bridge snapshot request failed (${response.status})`); + const response=await fetch(endpoint,{headers:bridgeHeaders(options),signal:AbortSignal.timeout(timeoutMs)}); + if(!response.ok)throw bridgeFailure('snapshot request failed',response.status); const declared=Number(response.headers.get('content-length')??0);if(declared>maxResponseBytes)throw new Error('Workspace bridge response exceeds size limit'); if(!response.body)throw new Error('Workspace bridge response has no body'); const reader=response.body.getReader(),chunks:Uint8Array[]=[];let size=0; @@ -40,25 +48,26 @@ export async function fetchBridgeSnapshot(baseUrl:string,workspaceId:string,expe export interface OverlayEntry { path: string; contentBase64?: string; mode?: string; delete?: true } -export async function overlayBridgeWorkspace(baseUrl:string,workspaceId:string,entries:OverlayEntry[],timeoutMs=15_000):Promise<{workspaceId:string;applied:number}>{ +export async function overlayBridgeWorkspace(baseUrl:string,workspaceId:string,entries:OverlayEntry[],options:number|BridgeRequestOptions=15_000):Promise<{workspaceId:string;applied:number}>{ if(!workspaceId||workspaceId.includes('/')||workspaceId.includes('\\'))throw new Error('Invalid bridge workspace ID'); const endpoint=bridgeEndpoint(baseUrl,`/extensions/foreman-workspace/v1/workspaces/${encodeURIComponent(workspaceId)}/overlay`); - const response=await fetch(endpoint,{method:'POST',headers:{'content-type':'application/json'},body:JSON.stringify({entries}),signal:AbortSignal.timeout(timeoutMs)}); - if(!response.ok)throw new Error(`Workspace bridge overlay request failed (${response.status})`); + const response=await fetch(endpoint,{method:'POST',headers:bridgeHeaders(options,{'content-type':'application/json'}),body:JSON.stringify({entries}),signal:AbortSignal.timeout(bridgeTimeout(options))}); + if(!response.ok)throw bridgeFailure('overlay request failed',response.status); const body=await boundedJson(response,64*1024) as Record; if(typeof body.workspace_id!=='string'||typeof body.applied!=='number')throw new Error('Workspace bridge returned an invalid overlay result'); return {workspaceId:body.workspace_id,applied:body.applied}; } -export async function seedBridgeWorkspace(baseUrl:string,pinnedBaseCommit:string,timeoutMs=15_000):Promise<{workspaceId:string;baseCommit:string}>{ +export async function seedBridgeWorkspace(baseUrl:string,pinnedBaseCommit:string,options:number|BridgeRequestOptions=15_000):Promise<{workspaceId:string;baseCommit:string}>{ if(!SHA.test(pinnedBaseCommit))throw new Error('A full pinned base commit SHA is required'); - const advertisedResponse=await fetch(bridgeEndpoint(baseUrl,'/v1/uhp'),{signal:AbortSignal.timeout(timeoutMs)}); - if(!advertisedResponse.ok)throw new Error(`Workspace bridge discovery failed (${advertisedResponse.status})`); + const timeoutMs=bridgeTimeout(options); + const advertisedResponse=await fetch(bridgeEndpoint(baseUrl,'/v1/uhp'),{headers:bridgeHeaders(options),signal:AbortSignal.timeout(timeoutMs)}); + if(!advertisedResponse.ok)throw bridgeFailure('discovery failed',advertisedResponse.status); const advertised=await boundedJson(advertisedResponse,64*1024) as any; const capability=advertised?.capabilities?.extensions?.foreman_workspace_bridge_v1; if(capability?.version!==1||capability.seed!==true||capability.complete_snapshot!==true||capability.execution_boundary!=='bubblewrap')throw new Error('Bridge does not advertise the complete snapshot and bubblewrap workspace extension'); const endpoint=bridgeEndpoint(baseUrl,'/extensions/foreman-workspace/v1/workspaces'); - const response=await fetch(endpoint,{method:'POST',headers:{'content-type':'application/json'},body:JSON.stringify({base_commit:pinnedBaseCommit}),signal:AbortSignal.timeout(timeoutMs)}); - if(!response.ok)throw new Error(`Workspace bridge seed request failed (${response.status})`); + const response=await fetch(endpoint,{method:'POST',headers:bridgeHeaders(options,{'content-type':'application/json'}),body:JSON.stringify({base_commit:pinnedBaseCommit}),signal:AbortSignal.timeout(timeoutMs)}); + if(!response.ok)throw bridgeFailure('seed request failed',response.status); const body=await boundedJson(response,64*1024) as Record; if(typeof body.workspace_id!=='string'||typeof body.base_commit!=='string'||body.base_commit.toLowerCase()!==pinnedBaseCommit.toLowerCase())throw new Error('Workspace bridge returned an invalid seed result'); return {workspaceId:body.workspace_id,baseCommit:body.base_commit}; diff --git a/tests/bridge-token.test.ts b/tests/bridge-token.test.ts new file mode 100644 index 0000000..0fb4d91 --- /dev/null +++ b/tests/bridge-token.test.ts @@ -0,0 +1,151 @@ +import { mkdtemp, readFile, rm } from 'node:fs/promises'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { createServer, type IncomingMessage, type Server } from 'node:http'; +import { once } from 'node:events'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { Controller, type UhpAdapter } from '../src/controller.js'; +import { JsonStore } from '../src/store.js'; +import { fetchBridgeSnapshot, overlayBridgeWorkspace, seedBridgeWorkspace } from '../src/verified-workspace.js'; + +vi.mock('../src/repo-digest.js', () => ({ + extractKeywords: () => ['Sleeper'], + buildRepoDigest: async (opts: { commit: string }) => ({ text: `## Digest for ${opts.commit.slice(0, 10)}`, commit: opts.commit }), +})); + +const TOKEN = 'f00dfeed'.repeat(8); +const SHA = 'a'.repeat(40); +const BASE = 'b'.repeat(40); +const servers: Server[] = []; +const dirs: string[] = []; +afterEach(async () => { + vi.restoreAllMocks(); + await Promise.all([...servers.splice(0).map(server => new Promise(resolve => server.close(() => resolve()))), ...dirs.splice(0).map(dir => rm(dir, { recursive: true, force: true }))]); +}); + +interface Seen { method: string; url: string; authorization: string | undefined } +/** A bridge that, like the real one with LOCAL_CLI_UHP_TOKEN set, answers 401 to any request without the bearer token. */ +async function tokenBridge(token: string | null = TOKEN): Promise<{ baseUrl: string; seen: Seen[] }> { + const seen: Seen[] = []; + let seeded = 0; + const server = createServer((req: IncomingMessage, res) => { + seen.push({ method: req.method!, url: req.url!, authorization: req.headers.authorization }); + let body = ''; + req.on('data', chunk => { body += chunk; }); + req.on('end', () => { + res.setHeader('content-type', 'application/json'); + if (token && req.headers.authorization !== `Bearer ${token}`) { res.statusCode = 401; res.end('{"error":{"code":"unauthorized"}}'); return; } + if (req.method === 'GET' && req.url === '/v1/uhp') { res.end(JSON.stringify({ capabilities: { extensions: { foreman_workspace_bridge_v1: { version: 1, seed: true, complete_snapshot: true, execution_boundary: 'bubblewrap' } } } })); return; } + if (req.method === 'POST' && req.url === '/extensions/foreman-workspace/v1/workspaces') { res.statusCode = 201; res.end(JSON.stringify({ workspace_id: `ws-${++seeded}`, base_commit: (JSON.parse(body) as { base_commit: string }).base_commit })); return; } + if (req.method === 'POST' && /^\/extensions\/foreman-workspace\/v1\/workspaces\/ws-\d+\/overlay$/.test(req.url!)) { res.end(JSON.stringify({ workspace_id: 'ws-1', applied: (JSON.parse(body) as { entries: unknown[] }).entries.length })); return; } + if (req.method === 'GET' && /^\/extensions\/foreman-workspace\/v1\/workspaces\/ws-\d+\/snapshot$/.test(req.url!)) { res.end(JSON.stringify({ complete: true, base_commit: SHA, entries: [], errors: [] })); return; } + res.statusCode = 404; res.end('{}'); + }); + }); + server.listen(0, '127.0.0.1'); + await once(server, 'listening'); + servers.push(server); + return { baseUrl: `http://127.0.0.1:${(server.address() as import('node:net').AddressInfo).port}`, seen }; +} + +describe('bridge helpers and the bearer token', () => { + it('send the token on discovery, seed, overlay and snapshot requests', async () => { + const { baseUrl, seen } = await tokenBridge(); + expect(await seedBridgeWorkspace(baseUrl, SHA, { token: TOKEN })).toEqual({ workspaceId: 'ws-1', baseCommit: SHA }); + expect(await overlayBridgeWorkspace(baseUrl, 'ws-1', [{ path: 'a.ts', contentBase64: '', mode: '100644' }], { token: TOKEN })).toEqual({ workspaceId: 'ws-1', applied: 1 }); + expect((await fetchBridgeSnapshot(baseUrl, 'ws-1', SHA, { token: TOKEN })).complete).toBe(true); + expect(seen.map(request => `${request.method} ${request.url}`)).toEqual(['GET /v1/uhp', 'POST /extensions/foreman-workspace/v1/workspaces', 'POST /extensions/foreman-workspace/v1/workspaces/ws-1/overlay', 'GET /extensions/foreman-workspace/v1/workspaces/ws-1/snapshot']); + for (const request of seen) expect(request.authorization).toBe(`Bearer ${TOKEN}`); + }); + + it('keep the numeric timeout arguments and send no credential when none is configured', async () => { + const { baseUrl, seen } = await tokenBridge(null); + await seedBridgeWorkspace(baseUrl, SHA, 5_000); + await overlayBridgeWorkspace(baseUrl, 'ws-1', [], 5_000); + await fetchBridgeSnapshot(baseUrl, 'ws-1', SHA, 5_000, 1024 * 1024); + await seedBridgeWorkspace(baseUrl, SHA); + expect(seen.length).toBe(6); + for (const request of seen) expect(request.authorization).toBeUndefined(); + await expect(fetchBridgeSnapshot(baseUrl, 'ws-1', SHA, { maxResponseBytes: 8 })).rejects.toThrow('exceeds size limit'); + }); + + it('report a rejected token as a 401 that says what to fix and never contains the token', async () => { + const { baseUrl } = await tokenBridge(); + for (const options of [undefined, { token: 'not-the-token' }]) { + const failures = await Promise.all([ + seedBridgeWorkspace(baseUrl, SHA, options).catch((error: Error) => error.message), + overlayBridgeWorkspace(baseUrl, 'ws-1', [], options).catch((error: Error) => error.message), + fetchBridgeSnapshot(baseUrl, 'ws-1', SHA, options).catch((error: Error) => error.message), + ]); + expect(failures).toEqual([ + 'Workspace bridge discovery failed (401): the bridge requires a valid bearer token', + 'Workspace bridge overlay request failed (401): the bridge requires a valid bearer token', + 'Workspace bridge snapshot request failed (401): the bridge requires a valid bearer token', + ]); + } + }); + + it('never send the token to a bridge URL that is not loopback', async () => { + const fetchSpy = vi.spyOn(globalThis, 'fetch'); + await expect(seedBridgeWorkspace('https://example.com', SHA, { token: TOKEN })).rejects.toThrow('loopback-only'); + await expect(overlayBridgeWorkspace('https://example.com', 'ws-1', [], { token: TOKEN })).rejects.toThrow('loopback-only'); + await expect(fetchBridgeSnapshot('https://example.com', 'ws-1', SHA, { token: TOKEN })).rejects.toThrow('loopback-only'); + expect(fetchSpy).not.toHaveBeenCalled(); + }); + + it('are given the token at every call site in the controller', async () => { + const source = await readFile(new URL('../src/controller.ts', import.meta.url), 'utf8'); + const calls = [...source.matchAll(/\b(seedBridgeWorkspace|overlayBridgeWorkspace|fetchBridgeSnapshot)\(([^()]*)\)/g)]; + expect(calls.length).toBeGreaterThanOrEqual(6); + for (const [call] of calls) expect(call, 'a bridge helper call that does not pass the bridge token').toMatch(/\{token:[\w.]+\.bridgeToken\}\)$/); + }); +}); + +async function controllerFor(bridgeToken: string | undefined) { + const dir = await mkdtemp(join(tmpdir(), 'foreman-bridge-token-')); + dirs.push(dir); + const store = new JsonStore(join(dir, 'state.json')); + const submissions: Array<{ roleId: string; config: Record }> = []; + const uhp: UhpAdapter = { submit: async input => { submissions.push(input as never); return { externalId: 'ext-1', status: 'completed', outputText: JSON.stringify({ workerTask: 'Update README.md.' }), result: { ok: true } }; }, cancel: async () => ({ status: 'cancelled' }) }; + await store.mutate(s => { for (const role of s.roles) { role.enabled = true; role.availableConfigs = [{ harnessId: 'claude-code', model: 'claude-fixture' }]; role.config = { harnessId: 'claude-code', model: 'claude-fixture' }; } }); + const controller = new Controller(store, uhp); + const { baseUrl, seen } = await tokenBridge(); + controller.configureVerifiedWorkspace({ repoPath: '/fixture/repo', bridgeBaseUrl: baseUrl, allowedScope: ['README.md'], commands: [{ name: 'check', command: 'true', args: [] }], ...(bridgeToken ? { bridgeToken } : {}) }); + return { controller, store, submissions, seen }; +} + +describe('Controller with a token-protected bridge', () => { + it('seeds the Planner and Orchestrator snapshots and prepares the Worker workspace with the configured token', async () => { + const { controller, store, submissions, seen } = await controllerFor(TOKEN); + vi.spyOn(controller as never, 'currentHead').mockResolvedValue(SHA as never); + vi.spyOn(controller as never, 'pinWorkerBase').mockResolvedValue(undefined as never); + const project: any = await controller.createProject('Token project'); + await controller.sendProjectPlannerMessage(project.id, 'Plan some work.'); + expect((await store.load()).projects[0]!.plannerAssignments?.at(-1)?.repoAccess).toMatchObject({ mode: 'snapshot', commit: SHA }); + const task: any = await controller.createTask(project.id, 'Task'); + const run: any = await controller.createRun(task.id); + await controller.prepareWorkerWorkspace(run.id, BASE); + await store.mutate(s => { s.projects[0]!.tasks[0]!.runs[0]!.pinnedBaseCommit = BASE; }); // pinWorkerBase is stubbed above + const workerWorkspace = (await store.load()).projects[0]!.tasks[0]!.runs[0]!.workspaceId; + expect(workerWorkspace).toMatch(/^ws-\d+$/); + await controller.addGuidance(run.id, 'Do the work.'); + await controller.orchestrate(run.id, 'Prepare a bounded task.'); + const orchestratorWorkspace = submissions.find(s => s.roleId === 'orchestrator')?.config.readOnlyWorkspaceId; + expect(orchestratorWorkspace).toMatch(/^ws-\d+$/); + expect(orchestratorWorkspace).not.toBe(workerWorkspace); + expect(seen.filter(request => request.method === 'POST').length).toBe(3); // planner snapshot, Worker workspace, Orchestrator snapshot + for (const request of seen) expect(request.authorization, `${request.method} ${request.url}`).toBe(`Bearer ${TOKEN}`); + }); + + it('falls back to a digest, naming the missing credential but never a token, when the controller has none', async () => { + const { controller, store, seen } = await controllerFor(undefined); + vi.spyOn(controller as never, 'currentHead').mockResolvedValue(SHA as never); + const project: any = await controller.createProject('No token project'); + await controller.sendProjectPlannerMessage(project.id, 'Plan some work.'); + const access = (await store.load()).projects[0]!.plannerAssignments?.at(-1)?.repoAccess; + expect(access?.mode).toBe('digest'); + expect(access?.reason).toMatch(/401.*bearer token/); + expect(JSON.stringify(await store.load())).not.toContain(TOKEN); + for (const request of seen) expect(request.authorization).toBeUndefined(); + }); +}); diff --git a/tests/config-bridge-token.test.ts b/tests/config-bridge-token.test.ts new file mode 100644 index 0000000..89ba71d --- /dev/null +++ b/tests/config-bridge-token.test.ts @@ -0,0 +1,36 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { loadConfig } from '../src/config.js'; + +const NAMES = ['FOREMAN_WORKSPACE_BRIDGE_TOKEN', 'FOREMAN_WORKSPACE_BRIDGE_URL']; +let saved: Record; +beforeEach(() => { saved = Object.fromEntries(NAMES.map(name => [name, process.env[name]])); for (const name of NAMES) delete process.env[name]; }); +afterEach(() => { for (const name of NAMES) { if (saved[name] === undefined) delete process.env[name]; else process.env[name] = saved[name]; } }); + +describe('FOREMAN_WORKSPACE_BRIDGE_TOKEN', () => { + it('is optional and absent by default', () => { + expect(loadConfig().workspaceBridgeToken).toBeUndefined(); + process.env.FOREMAN_WORKSPACE_BRIDGE_TOKEN = ' '; + expect(loadConfig().workspaceBridgeToken).toBeUndefined(); + }); + + it('is read, trimmed, alongside the external bridge URL', () => { + process.env.FOREMAN_WORKSPACE_BRIDGE_URL = 'http://127.0.0.1:8787'; + process.env.FOREMAN_WORKSPACE_BRIDGE_TOKEN = ' s3cret-Token_1 '; + expect(loadConfig()).toMatchObject({ workspaceBridgeUrl: 'http://127.0.0.1:8787', workspaceBridgeToken: 's3cret-Token_1' }); + }); + + it('requires the bridge URL it belongs to', () => { + process.env.FOREMAN_WORKSPACE_BRIDGE_TOKEN = 'orphan-token'; + expect(() => loadConfig()).toThrow('FOREMAN_WORKSPACE_BRIDGE_TOKEN requires FOREMAN_WORKSPACE_BRIDGE_URL'); + }); + + it('rejects a value that cannot be a header token, without echoing it', () => { + process.env.FOREMAN_WORKSPACE_BRIDGE_URL = 'http://127.0.0.1:8787'; + for (const bad of ['has a space', 'tab\there', 'café', 'x'.repeat(513)]) { + process.env.FOREMAN_WORKSPACE_BRIDGE_TOKEN = bad; + let message = ''; + try { loadConfig(); } catch (error) { message = (error as Error).message; } + expect(message).toBe('FOREMAN_WORKSPACE_BRIDGE_TOKEN must be 1-512 printable ASCII characters without spaces'); + } + }); +}); diff --git a/tests/local-bridge.test.ts b/tests/local-bridge.test.ts index ad89356..d889faf 100644 --- a/tests/local-bridge.test.ts +++ b/tests/local-bridge.test.ts @@ -1,14 +1,22 @@ import { EventEmitter } from 'node:events'; -import { mkdtemp, mkdir, rm } from 'node:fs/promises'; +import { readFileSync } from 'node:fs'; +import { chmod, mkdtemp, mkdir, readdir, rm, stat, writeFile } from 'node:fs/promises'; import { tmpdir, homedir } from 'node:os'; import { join } from 'node:path'; import { afterEach, describe, expect, it, vi } from 'vitest'; -import { LocalBridge } from '../src/local-bridge.js'; +import { LocalBridge, bearerFetch, type LocalBridgeHealth, type LocalBridgeOptions } from '../src/local-bridge.js'; +let nextPid = 40_000; class FakeChild extends EventEmitter { exitCode: number | null = null; + signalCode: string | null = null; killed = false; + pid = nextPid++; + /** Whether the fake bridge answers HTTP; a bridge that cannot bind its port never does. */ + listening = true; kill(signal = 'SIGTERM') { this.killed = true; if (signal === 'SIGKILL' || signal === 'SIGTERM') { this.exitCode = 0; queueMicrotask(() => this.emit('exit', 0, signal)); } return true; } + /** The process dies by itself, as after a crash or `kill -9`. */ + crash(code: number | null = 1, signal: string | null = null) { this.exitCode = signal ? null : code; this.signalCode = signal; this.listening = false; this.emit('exit', code, signal); } } const roots: string[] = []; @@ -107,3 +115,364 @@ describe('LocalBridge', () => { expect(wd).toBe(join(tempRoot, 'foreman-local-bridge-work', key)); }); }); + +const DISCOVERY = JSON.stringify({ protocol: 'uhp', implementation: { name: 'local-cli-uhp' } }); + +/** Poll with real time (setImmediate is never faked) until an asynchronous chain of real fs and timer work has landed. */ +async function until(predicate: () => boolean, label: string, timeoutMs = 3_000): Promise { + const deadline = Date.now() + timeoutMs; // callers only use this with real timers + while (!predicate()) { if (Date.now() > deadline) throw new Error(`timed out waiting for ${label}`); await new Promise(resolveTick => setTimeout(resolveTick, 2)); } +} +/** Same, but for tests that fake timers: only setImmediate and a wall-clock deadline are used. */ +async function untilFaked(predicate: () => boolean, label: string): Promise { + const deadline = performance.now() + 3_000; + while (!predicate()) { if (performance.now() > deadline) throw new Error(`timed out waiting for ${label}`); await new Promise(resolveTick => setImmediate(resolveTick)); } +} + +/** + * A LocalBridge wired to fake processes. `behavior` decides what each spawned process does: `ready` (default) answers + * discovery, `exit` dies at once (as when its port cannot be rebound), `silent` stays alive but never answers. + */ +async function supervised(options: Partial & { behavior?: (spawnNumber: number) => 'ready' | 'exit' | 'silent' } = {}) { + const f = await fixture(); + const children: FakeChild[] = []; + const spawnOptions: Array<{ env: NodeJS.ProcessEnv; stdio: unknown[] }> = []; + const fetchCalls: Array<{ url: string; authorization?: string }> = []; + const health: LocalBridgeHealth[] = []; + const allocatePort = vi.fn(async () => 49500); + const spawnFake = vi.fn((_command: string, _args: readonly string[], spawnOpts: { env: NodeJS.ProcessEnv; stdio: unknown[] }) => { + const child = new FakeChild(); + children.push(child); spawnOptions.push(spawnOpts); + const behavior = options.behavior?.(children.length) ?? 'ready'; + if (behavior === 'exit') { child.listening = false; queueMicrotask(() => child.crash(1)); } + else if (behavior === 'silent') child.listening = false; + return child as never; + }); + const fetchFake = vi.fn(async (url: string | URL | Request, init?: RequestInit) => { + fetchCalls.push({ url: String(url), authorization: (init?.headers as Record | undefined)?.authorization }); + const child = children.at(-1); + if (!child || !child.listening || child.exitCode !== null || child.signalCode) throw new TypeError('fetch failed'); + return new Response(DISCOVERY, { status: 200 }); + }); + const { behavior: _behavior, ...bridgeOptions } = options; + const manager = new LocalBridge({ bridgeScript: f.bridge, dataDir: join(f.root, 'data'), homeDir: f.home, tempDir: join(f.root, 'tmp'), spawn: spawnFake as never, fetch: fetchFake as never, allocatePort, startupPollMs: 2, onHealthChange: h => health.push(h), ...bridgeOptions }); + return { f, manager, children, spawnOptions, spawnFake, fetchCalls, health, allocatePort, logPath: () => join(manager.stateDirForInstance(f.repo, f.repo), 'bridge.log') }; +} + +describe('LocalBridge authentication', () => { + it('generates a fresh 32-byte token per start, gives it only to the child, and polls readiness with it', async () => { + process.env.LOCAL_CLI_UHP_TOKEN = 'ambient-token-that-must-not-be-forwarded'; + try { + const a = await supervised(); + const first = await a.manager.start(a.f.repo); + expect(first.token).toMatch(/^[a-f0-9]{64}$/); + expect(a.spawnOptions[0]!.env.LOCAL_CLI_UHP_TOKEN).toBe(first.token); + expect(a.spawnOptions[0]!.env.LOCAL_CLI_UHP_TOKEN).not.toBe(process.env.LOCAL_CLI_UHP_TOKEN); + expect(a.fetchCalls.length).toBeGreaterThan(0); + for (const call of a.fetchCalls) expect(call).toEqual({ url: `${first.baseUrl}/v1/uhp`, authorization: `Bearer ${first.token}` }); + // The token is readable but never enumerable: serialising or logging a status cannot leak it. + expect(Object.keys(first)).toEqual(['repoPath', 'baseUrl']); + expect(JSON.stringify(first)).not.toContain(first.token); + expect(JSON.stringify(a.manager.health)).not.toContain(first.token); + expect(a.manager.status?.token).toBe(first.token); + await a.manager.stop(); + // A new start is a new credential. + const restarted = await a.manager.start(a.f.repo); + expect(restarted.token).toMatch(/^[a-f0-9]{64}$/); + expect(restarted.token).not.toBe(first.token); + expect(a.spawnOptions[1]!.env.LOCAL_CLI_UHP_TOKEN).toBe(restarted.token); + await a.manager.stop(); + // Two managers never share one. + const b = await supervised(); + expect((await b.manager.start(b.f.repo)).token).not.toBe(first.token); + await b.manager.stop(); + } finally { delete process.env.LOCAL_CLI_UHP_TOKEN; } + }); + + it('fails fast, without echoing the token, when the port answers the authenticated readiness check with 401', async () => { + const f = await fixture(); + const child = new FakeChild(); + const fetchFake = vi.fn(async () => new Response('{"error":{"code":"unauthorized"}}', { status: 401 })); + const manager = new LocalBridge({ bridgeScript: f.bridge, dataDir: join(f.root, 'data'), homeDir: f.home, tempDir: join(f.root, 'tmp'), spawn: (() => child) as never, fetch: fetchFake as never, allocatePort: async () => 49501, startupTimeoutMs: 5_000 }); + const started = Date.now(); + const error = await manager.start(f.repo).catch(e => e as Error); + expect(error).toBeInstanceOf(Error); + expect((error as Error).message).toMatch(/answered the readiness check with 401/); + expect(Date.now() - started).toBeLessThan(2_000); + expect(fetchFake).toHaveBeenCalledTimes(1); + expect(child.killed).toBe(true); + const sent = (fetchFake.mock.calls[0] as unknown as [string, { headers: { authorization: string } }])[1].headers.authorization.replace('Bearer ', ''); + expect((error as Error).message).not.toContain(sent); + }); +}); + +describe('bearerFetch', () => { + it('adds the bearer token to every request, keeps what the caller sent, and never overrides an explicit Authorization', async () => { + const seen: Array<{ url: string; headers: Headers; method?: string; body?: unknown }> = []; + const base = (async (input: string | URL | Request, init?: RequestInit) => { seen.push({ url: String(input), headers: new Headers(init?.headers), method: init?.method, body: init?.body }); return new Response('{}'); }) as typeof fetch; + const authed = bearerFetch('secret-token', base); + await authed(new URL('http://127.0.0.1:1/v1/uhp'), { signal: AbortSignal.timeout(1_000) }); + await authed('http://127.0.0.1:1/v1/responses', { method: 'POST', headers: { 'Content-Type': 'application/json', 'Idempotency-Key': 'k' }, body: '{"a":1}' }); + await authed('http://127.0.0.1:1/x', { headers: { Authorization: 'Bearer explicit' } }); + await authed(new Request('http://127.0.0.1:1/y', { headers: { 'UHP-Version': 'v' } })); + expect(seen.map(call => call.headers.get('authorization'))).toEqual(['Bearer secret-token', 'Bearer secret-token', 'Bearer explicit', 'Bearer secret-token']); + expect(seen[1]).toMatchObject({ method: 'POST', body: '{"a":1}' }); + expect(seen[1]!.headers.get('content-type')).toBe('application/json'); + expect(seen[1]!.headers.get('idempotency-key')).toBe('k'); + expect(seen[3]!.headers.get('uhp-version')).toBe('v'); + }); +}); + +describe('LocalBridge supervision', () => { + it('restarts a crashed bridge on the same port with the same token and reports its health', async () => { + const s = await supervised({ supervision: { restartBaseDelayMs: 5, restartMaxDelayMs: 40 } }); + const status = await s.manager.start(s.f.repo); + expect(s.manager.health).toMatchObject({ state: 'ready', restarts: 0 }); + s.children[0]!.crash(null, 'SIGKILL'); + // Readiness is cleared at once, so nobody is handed the dead process's status. + expect(s.manager.status).toBeUndefined(); + expect(s.manager.health).toMatchObject({ state: 'restarting', restartDelayMs: 5, restarts: 0, lastExit: { code: null, signal: 'SIGKILL' } }); + expect(s.manager.health.message).toMatch(/signal SIGKILL.*restarting in 5 ms.*bridge\.log/); + await until(() => s.manager.health.state === 'ready', 'the restarted bridge'); + expect(s.spawnFake).toHaveBeenCalledTimes(2); + expect(s.allocatePort).toHaveBeenCalledTimes(1); + for (const key of ['LOCAL_CLI_UHP_PORT', 'LOCAL_CLI_UHP_TOKEN', 'LOCAL_CLI_UHP_STATE', 'LOCAL_CLI_UHP_WORK', 'LOCAL_CLI_UHP_SOURCE_REPO'] as const) expect(s.spawnOptions[1]!.env[key]).toBe(s.spawnOptions[0]!.env[key]); + // Controllers still hold the original URL and credential, and both keep working. + expect(s.manager.status).toBe(status); + expect(s.manager.status?.baseUrl).toBe(status.baseUrl); + expect(s.manager.status?.token).toBe(status.token); + expect(s.fetchCalls.at(-1)).toEqual({ url: `${status.baseUrl}/v1/uhp`, authorization: `Bearer ${status.token}` }); + expect(s.manager.health).toMatchObject({ state: 'ready', restarts: 1, lastExit: { code: null, signal: 'SIGKILL' } }); + expect(s.health.map(h => h.state)).toEqual(['restarting', 'ready']); + expect(JSON.stringify(s.health)).not.toContain(status.token); + await s.manager.stop(); + }); + + it('backs off exponentially from 1 s and caps at 30 s (fake timers, default policy)', async () => { + vi.useFakeTimers({ toFake: ['setTimeout', 'clearTimeout', 'Date'] }); + try { + // The first restart succeeds; every later one dies at once, so the delays keep growing. + const s = await supervised({ behavior: n => n === 1 ? 'ready' : 'exit', supervision: { maxFastFailures: 8 } }); + await s.manager.start(s.f.repo); + const crashedAt = Date.now(); + s.children[0]!.crash(1); + expect(s.manager.health).toMatchObject({ state: 'restarting', restartDelayMs: 1_000 }); + expect(s.manager.health.nextRestartAt).toBe(new Date(crashedAt + 1_000).toISOString()); + // Each restart is scheduled after the previous failure: exactly `delay` ms later, never a millisecond sooner. + let spawned = 1; + for (const delay of [1_000, 2_000, 4_000, 8_000, 16_000, 30_000, 30_000]) { + await vi.advanceTimersByTimeAsync(delay - 1); + await new Promise(resolveTick => setImmediate(resolveTick)); + expect(s.spawnFake, `no restart before ${delay} ms`).toHaveBeenCalledTimes(spawned); + await vi.advanceTimersByTimeAsync(1); + spawned++; + // The restarted process dies at once (its port cannot be rebound), which is reported and schedules the next restart. + await untilFaked(() => s.spawnFake.mock.calls.length === spawned && s.health.length === spawned, `failure of spawn ${spawned}`); + } + expect(s.health.filter(h => h.state === 'restarting').map(h => h.restartDelayMs)).toEqual([1_000, 2_000, 4_000, 8_000, 16_000, 30_000, 30_000]); + // The eighth consecutive fast failure exhausts the policy. + expect(s.manager.health.state).toBe('unavailable'); + expect(s.health.at(-1)?.state).toBe('unavailable'); + await s.manager.stop(); + } finally { vi.useRealTimers(); } + }); + + it('gives up after 5 consecutive fast failures, including restarts that cannot rebind the port, and reports why', async () => { + const s = await supervised({ behavior: n => n === 1 ? 'ready' : 'exit', supervision: { restartBaseDelayMs: 1, restartMaxDelayMs: 4 } }); + const status = await s.manager.start(s.f.repo); + s.children[0]!.crash(null, 'SIGKILL'); + await until(() => s.manager.health.state === 'unavailable', 'the bridge to be given up on'); + // The crash plus four failed restarts: five consecutive runs within 60 s of starting. + expect(s.spawnFake).toHaveBeenCalledTimes(5); + expect(s.manager.status).toBeUndefined(); + const health = s.manager.health; + expect(health).toMatchObject({ state: 'unavailable', restarts: 4, lastExit: { code: 1, signal: null }, logPath: s.logPath() }); + expect(health.message).toMatch(/5 consecutive failures within 60 s of starting \(last exit: exit code 1\); see .*bridge\.log/); + expect(JSON.stringify(health)).not.toContain(status.token); + expect(s.health.at(-1)?.state).toBe('unavailable'); + await new Promise(resolveWait => setTimeout(resolveWait, 60)); + expect(s.spawnFake).toHaveBeenCalledTimes(5); // no further attempts + await s.manager.stop(); + }); + + it('stops a restarted bridge that never answers and counts it as a failed attempt', async () => { + const s = await supervised({ behavior: n => n === 1 ? 'ready' : 'silent', startupTimeoutMs: 40, supervision: { restartBaseDelayMs: 1, maxFastFailures: 2 } }); + await s.manager.start(s.f.repo); + s.children[0]!.crash(1); + await until(() => s.manager.health.state === 'unavailable', 'the silent restart to be given up on'); + expect(s.children[1]!.killed).toBe(true); + expect(s.manager.health.lastExit?.reason).toMatch(/did not become ready; see .*bridge\.log/); + expect(s.spawnFake).toHaveBeenCalledTimes(2); + }); + + it('counts only runs that ended within the stability window: a long healthy run resets the failure count', async () => { + // maxFastFailures 2. Without the reset the first crash (after 200 ms) would count, and the bridge would be given up after two crashes. + const s = await supervised({ supervision: { restartBaseDelayMs: 1, stableAfterMs: 150, maxFastFailures: 2 } }); + await s.manager.start(s.f.repo); + await new Promise(resolveWait => setTimeout(resolveWait, 200)); + s.children[0]!.crash(1); // healthy for >= 150 ms: restarts the count + await until(() => s.children.length === 2 && s.manager.health.state === 'ready', 'the first restart'); + s.children[1]!.crash(1); // fast failure 1 of 2 + await until(() => s.children.length === 3 && s.manager.health.state === 'ready', 'the second restart'); + expect(s.manager.health).toMatchObject({ state: 'ready', restarts: 2 }); + s.children[2]!.crash(1); // fast failure 2 of 2 + await until(() => s.manager.health.state === 'unavailable', 'the bridge to be given up on'); + expect(s.spawnFake).toHaveBeenCalledTimes(3); + await s.manager.stop(); + }); + + it('never restarts after a deliberate stop, whether the bridge is ready, waiting to restart, or mid-restart', async () => { + // Ready. + const ready = await supervised({ supervision: { restartBaseDelayMs: 5 } }); + await ready.manager.start(ready.f.repo); + await ready.manager.stop(); + expect(ready.children[0]!.killed).toBe(true); + await new Promise(resolveWait => setTimeout(resolveWait, 40)); + expect(ready.spawnFake).toHaveBeenCalledTimes(1); + expect(ready.manager.health.state).toBe('unavailable'); + expect(ready.health).toEqual([]); // a deliberate stop is not reported as a failure + + // Waiting for its restart timer. + const waiting = await supervised({ supervision: { restartBaseDelayMs: 30 } }); + await waiting.manager.start(waiting.f.repo); + waiting.children[0]!.crash(1); + expect(waiting.manager.health.state).toBe('restarting'); + await waiting.manager.stop(); + await new Promise(resolveWait => setTimeout(resolveWait, 90)); + expect(waiting.spawnFake).toHaveBeenCalledTimes(1); + expect(waiting.manager.health.state).toBe('unavailable'); + + // In the middle of a restart whose bridge is not up yet. + const mid = await supervised({ behavior: n => n === 1 ? 'ready' : 'silent', startupTimeoutMs: 5_000, supervision: { restartBaseDelayMs: 1 } }); + await mid.manager.start(mid.f.repo); + mid.children[0]!.crash(1); + await until(() => mid.children.length === 2, 'the restart to spawn'); + await mid.manager.stop(); + expect(mid.children[1]!.killed).toBe(true); + await new Promise(resolveWait => setTimeout(resolveWait, 60)); + expect(mid.spawnFake).toHaveBeenCalledTimes(2); + expect(mid.manager.health.state).toBe('unavailable'); + }); + + it('does not treat the exit of a process it stopped for a new repository as a crash', async () => { + const s = await supervised({ supervision: { restartBaseDelayMs: 1 } }); + const second = join(s.f.root, 'other'); + await mkdir(join(second, '.git'), { recursive: true }); + await s.manager.start(s.f.repo, 'one'); + await s.manager.start(second, 'two'); + await new Promise(resolveWait => setTimeout(resolveWait, 40)); + expect(s.spawnFake).toHaveBeenCalledTimes(2); + expect(s.manager.status?.repoPath).toBe(second); + await s.manager.stop(); + }); +}); + +describe('LocalBridge log file', () => { + it('sends stdout and stderr of every run to a private log in the state directory', async () => { + const s = await supervised({ supervision: { restartBaseDelayMs: 2 } }); + await s.manager.start(s.f.repo); + const log = s.logPath(); + expect(log).toBe(join(s.manager.stateDirForInstance(s.f.repo, s.f.repo), 'bridge.log')); + expect(s.manager.health.logPath).toBe(log); + expect((await stat(log)).mode & 0o777).toBe(0o600); + const [stdin, stdout, stderr] = s.spawnOptions[0]!.stdio; + expect(stdin).toBe('ignore'); + expect(typeof stdout).toBe('number'); + expect(stderr).toBe(stdout); + s.children[0]!.crash(null, 'SIGKILL'); + await until(() => s.manager.health.state === 'ready', 'the restart'); + await s.manager.stop(); + await until(() => /restarting in 2 ms/.test(readLog(log)) && (readLog(log).match(/bridge ready/g)?.length ?? 0) === 2, 'Foreman notes in the log'); + const text = readLog(log); + expect(text).toMatch(/starting bridge on http:\/\/127\.0\.0\.1:49500/); + expect(text).toMatch(/bridge ended \(signal SIGKILL\)/); + expect(text).not.toContain(s.spawnOptions[0]!.env.LOCAL_CLI_UHP_TOKEN); + expect((await stat(log)).mode & 0o777).toBe(0o600); + }); + + it('tightens the mode of an existing log and appends to it', async () => { + const s = await supervised(); + const dir = s.manager.stateDirForInstance(s.f.repo, s.f.repo); + await mkdir(dir, { recursive: true }); + await writeFile(join(dir, 'bridge.log'), 'earlier run\n'); + await chmod(join(dir, 'bridge.log'), 0o644); + await s.manager.start(s.f.repo); + expect((await stat(join(dir, 'bridge.log'))).mode & 0o777).toBe(0o600); + expect(readLog(join(dir, 'bridge.log'))).toMatch(/^earlier run\n/); + await s.manager.stop(); + }); + + it('rotates a log over the limit to bridge.log.1 at start, keeping a single old file', async () => { + const s = await supervised({ supervision: { logMaxBytes: 32 } }); + const dir = s.manager.stateDirForInstance(s.f.repo, s.f.repo); + await mkdir(dir, { recursive: true }); + await writeFile(join(dir, 'bridge.log'), 'x'.repeat(100)); + await writeFile(join(dir, 'bridge.log.1'), 'an older rotation that must be discarded'); + await s.manager.start(s.f.repo); + expect(readLog(join(dir, 'bridge.log.1'))).toBe('x'.repeat(100)); + expect(readLog(join(dir, 'bridge.log'))).not.toContain('xxxx'); + expect(readLog(join(dir, 'bridge.log'))).toMatch(/starting bridge/); + expect((await stat(join(dir, 'bridge.log'))).mode & 0o777).toBe(0o600); + expect((await readdir(dir)).filter(name => name.startsWith('bridge.log')).sort()).toEqual(['bridge.log', 'bridge.log.1']); + await s.manager.stop(); + }); + + it('leaves a log under the limit alone', async () => { + const s = await supervised({ supervision: { logMaxBytes: 1_000 } }); + const dir = s.manager.stateDirForInstance(s.f.repo, s.f.repo); + await mkdir(dir, { recursive: true }); + await writeFile(join(dir, 'bridge.log'), 'small\n'); + await s.manager.start(s.f.repo); + expect((await readdir(dir)).filter(name => name.startsWith('bridge.log'))).toEqual(['bridge.log']); + expect(readLog(join(dir, 'bridge.log'))).toMatch(/^small\n/); + await s.manager.stop(); + }); +}); + +function readLog(path: string): string { try { return readFileSync(path, 'utf8'); } catch { return ''; } } + +describe('LocalBridge with a real child process', () => { + it('logs both streams, answers only the token, and comes back on the same port and token after kill -9', async () => { + const f = await fixture(); + const script = join(f.root, 'real-bridge.mjs'); + await writeFile(script, ` +import { createServer } from 'node:http'; +const token = process.env.LOCAL_CLI_UHP_TOKEN, port = Number(process.env.LOCAL_CLI_UHP_PORT); +console.log('bridge stdout pid=' + process.pid); console.error('bridge stderr pid=' + process.pid); +createServer((req, res) => { + if (req.headers.authorization !== 'Bearer ' + token) { res.writeHead(401); res.end('{}'); return; } + res.writeHead(200, { 'content-type': 'application/json' }); + res.end(JSON.stringify({ protocol: 'uhp', implementation: { name: 'local-cli-uhp' }, pid: process.pid, token })); +}).listen(port, '127.0.0.1'); +`); + const health: LocalBridgeHealth[] = []; + const manager = new LocalBridge({ bridgeScript: script, dataDir: join(f.root, 'data'), homeDir: f.home, tempDir: join(f.root, 'tmp'), supervision: { restartBaseDelayMs: 20 }, onHealthChange: h => health.push(h) }); + try { + const status = await manager.start(f.repo); + const ask = async (token?: string) => fetch(`${status.baseUrl}/v1/uhp`, { headers: token ? { authorization: `Bearer ${token}` } : {} }); + expect((await ask()).status).toBe(401); + const before = await (await ask(status.token)).json() as { pid: number; token: string }; + expect(before.token).toBe(status.token); + process.kill(before.pid, 'SIGKILL'); + await until(() => manager.health.state === 'restarting', 'the crash to be noticed'); + expect(manager.status).toBeUndefined(); + expect(manager.health.lastExit).toMatchObject({ code: null, signal: 'SIGKILL' }); + await until(() => manager.health.state === 'ready', 'the restart', 8_000); + const after = await (await fetch(`${status.baseUrl}/v1/uhp`, { headers: { authorization: `Bearer ${status.token}` } })).json() as { pid: number; token: string }; + expect(after.pid).not.toBe(before.pid); + expect(after.token).toBe(status.token); + expect(manager.status?.baseUrl).toBe(status.baseUrl); + const log = readLog(manager.health.logPath!); + expect(log).toContain(`bridge stdout pid=${before.pid}`); + expect(log).toContain(`bridge stderr pid=${before.pid}`); + expect(log).toContain(`bridge stdout pid=${after.pid}`); + expect(log).toContain('bridge ended (signal SIGKILL)'); + expect(log).not.toContain(status.token); + expect(health.map(h => h.state)).toEqual(['restarting', 'ready']); + await manager.stop(); + await new Promise(resolveWait => setTimeout(resolveWait, 150)); + expect(manager.health.state).toBe('unavailable'); + await expect(ask(status.token)).rejects.toThrow(); // the port is closed and nothing came back + } finally { await manager.stop(); } + }); +});