diff --git a/docker-compose.yaml b/docker-compose.yaml index ab37fd24..1d7d28cf 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -15,6 +15,12 @@ services: - CODEAPI_BRIDGE_DYNAMIC_WORKERS=${CODEAPI_BRIDGE_DYNAMIC_WORKERS:-true} - CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS=${CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS:-1} - CODEAPI_BRIDGE_WORKER_ID=${CODEAPI_BRIDGE_WORKER_ID:-} + - CODEAPI_BRIDGE_RECOVERY_SERVER_ID=${CODEAPI_BRIDGE_RECOVERY_SERVER_ID:-} + - CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS=${CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS:-0} + - CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS=${CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS:-60} + - CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE:-12} + - CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE:-30} + - CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE:-240} - CODEAPI_AUTH_PROVIDER=${CODEAPI_AUTH_PROVIDER:-} - CODEAPI_ALLOW_AUTH_PROVIDER_NONE=${CODEAPI_ALLOW_AUTH_PROVIDER_NONE:-} - CODEAPI_JWT_ISSUER=${CODEAPI_JWT_ISSUER:-} @@ -64,6 +70,12 @@ services: - CODEAPI_BRIDGE_DYNAMIC_WORKERS=${CODEAPI_BRIDGE_DYNAMIC_WORKERS:-true} - CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS=${CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS:-1} - CODEAPI_BRIDGE_WORKER_ID=${CODEAPI_BRIDGE_WORKER_ID:-} + - CODEAPI_BRIDGE_RECOVERY_SERVER_ID=${CODEAPI_BRIDGE_RECOVERY_SERVER_ID:-} + - CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS=${CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS:-0} + - CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS=${CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS:-60} + - CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE:-12} + - CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE:-30} + - CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE=${CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE:-240} - CODEAPI_AUTH_PROVIDER=${CODEAPI_AUTH_PROVIDER:-} - CODEAPI_JWT_SINGLE_TENANT_ID=${CODEAPI_JWT_SINGLE_TENANT_ID:-} - CODEAPI_TENANT_ISOLATION_STRICT=${CODEAPI_TENANT_ISOLATION_STRICT:-} @@ -213,9 +225,11 @@ services: redis: image: redis:7-alpine container_name: redis - command: redis-server --requirepass localdev + command: redis-server --requirepass localdev --appendonly yes ports: - ${CODEAPI_REDIS_PORT:-16379}:6379 + volumes: + - redis_data:/data minio: image: quay.io/minio/minio @@ -232,3 +246,4 @@ services: volumes: minio_data: + redis_data: diff --git a/docs/adr/001-stateful-code-environments.md b/docs/adr/001-stateful-code-environments.md index cfa44bf4..7992ded8 100644 --- a/docs/adr/001-stateful-code-environments.md +++ b/docs/adr/001-stateful-code-environments.md @@ -52,8 +52,11 @@ worker replacement; the UI and operator documentation must not imply otherwise. - The VM requires no inbound internet listener. - Code API, not the worker, authenticates LibreChat users and normalizes work. - A stolen short-lived credential is insufficient without the worker private - key; a stolen private key is insufficient after credential expiry or - revocation. + key. In the original pairing-only model, a stolen private key is insufficient + after credential expiry or revocation. With optional durable machine + enrollment and signed credential recovery, the private key itself remains + a revocable long-lived credential: access-credential expiry alone does not + protect against theft of that key. Revocation invalidates both. - Pairing codes and credentials are stored by digest where lookup permits. - One configured worker has at most one active fenced assignment. - Sandbox isolation and default-deny egress remain the mandatory default; diff --git a/docs/byom-worker-admission.md b/docs/byom-worker-admission.md index 2fcd9f9f..0df8f17c 100644 --- a/docs/byom-worker-admission.md +++ b/docs/byom-worker-admission.md @@ -5,8 +5,13 @@ The limit is 32 admitted requests per worker, including the active request. When the limit is reached, the workspace endpoint returns HTTP 429 with `WORKER_QUEUE_FULL`. A different worker has an independent admission queue. -The Code API workspace HTTP endpoint allows up to the smaller of `JOB_TIMEOUT` -and five minutes for admission while the caller stays connected. After admission +The Code API workspace HTTP endpoint allows 30 seconds for admission when the +request has no `X-LibreChat-Workspace-Queue-Wait-Ms` header. A caller can +advertise a positive integer allowance in milliseconds through that header, +up to five minutes and any configured server queue ceiling. An invalid or +out-of-range value is rejected before dispatch. This queue budget is +independent of `JOB_TIMEOUT`; a shorter client or proxy deadline still ends +the wait. After admission and worker validation, a separate execution deadline starts. Commands receive their requested timeout (30 seconds by default, up to five minutes), capped by the operator's `JOB_TIMEOUT`, plus five seconds to settle the result. Other @@ -29,22 +34,27 @@ Existing workers still execute one assignment at a time. Parallel execution acro workspaces requires separate lease claims and isolated native sandbox contexts; this admission change does not advertise that capability. -Clients and reverse proxies must allow queue time plus execution/settlement time -and five seconds for HTTP delivery. With the default five-minute `JOB_TIMEOUT`, -that is at least 335 seconds for non-command tools, 340 seconds for default -commands, and 610 seconds for five-minute commands. With a smaller `JOB_TIMEOUT`, -use `min(JOB_TIMEOUT, 300s)` for the queue, plus `min(JOB_TIMEOUT, 30s)` for other +Clients and reverse proxies must allow the admitted queue budget plus +execution/settlement time and five seconds for HTTP delivery. With the default +five-minute `JOB_TIMEOUT` and no queue header, that is at least 65 seconds for +non-command tools, 70 seconds for default commands, and 340 seconds for +five-minute commands. At the maximum advertised five-minute queue allowance, +those totals become 335, 340, and 610 seconds respectively. With a smaller +`JOB_TIMEOUT`, use the advertised allowance (or 30 seconds without a header), +bounded by the server queue ceiling, plus `min(JOB_TIMEOUT, 30s)` for other operations or `min(JOB_TIMEOUT, requested command timeout) + 5s` for commands, -plus five seconds for delivery. +plus five seconds for delivery. The caller should advertise only the queue time +left after reserving execution, settlement, and delivery under its own HTTP +deadline; Code API does not receive that absolute deadline. -At the time of this change, LibreChat's `getWorkspaceToolTimeoutMs` still budgets -only 30 seconds for a single admission attempt (65/70/340 seconds in total). -Its `maxQueueWaitMs` is a retry horizon after a typed capacity rejection, **not** -a per-attempt HTTP timeout. Updating Code API alone therefore does not guarantee -the full wait. An earlier client, tool, or proxy timeout disconnects the request; -if work was already admitted, a mutation may have run and must not be blindly -retried. Match LibreChat's per-attempt timeout and each intermediary to the new -budget before relying on it. Existing workers do not need an update. +LibreChat's `maxQueueWaitMs` is a retry horizon after a typed capacity +rejection, **not** a per-attempt HTTP timeout. Without its opt-in +`maxRequestTimeoutMs`, LibreChat keeps a 30-second admission allowance per +attempt. Enabling a longer client budget requires LibreChat's header support on +every API replica and a timed canary through each intermediary; changing Code +API alone does not guarantee the full wait. An earlier client, tool, or proxy +timeout disconnects the request; if work was already admitted, a mutation may +have run and must not be blindly retried. Existing workers do not need an update. Focused regression coverage lives in `service/src/bridge/admission.test.ts`, `service/src/bridge/worker-admission.test.ts`, diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index c5bc573c..464e4763 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -31,6 +31,42 @@ CODEAPI_BRIDGE_TOKEN= CODEAPI_BRIDGE_AUTH_MODE=paired ``` +To opt in to durable machine authorization on every Code API replica, set a +single stable public **Code API** origin (not the LibreChat URL): + +```dotenv +CODEAPI_BRIDGE_RECOVERY_SERVER_ID=https://code.example.com +# 0 (default): enrolled machine keys remain authorized until revoked. +# CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS=0 +# CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS=60 +# CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE=12 +# CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE=30 +# CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE=240 +``` + +Omitting the server ID retains the existing pairing and refresh behavior and +hides the recovery routes. Deploy the compatible Code API version to **all** +replicas before setting this value and enrolling workers again. Older Code API +replicas can still pair or refresh a worker but do not write durable enrollment; +they must not serve device login or recovery requests. Only a pairing redeemed +after this option is enabled has a recoverable key. Updating Code API alone +does not make old workers reconnect automatically: the CLI must also implement +this recovery protocol in the later worker release. + +Store Redis state durably across restarts. The primary `docker-compose.yaml` +now uses Redis AOF and a named `/data` volume; preserve that volume when +recreating the stack. If upgrading a running stack with an in-memory Redis, +migrate its state before recreating the container: mounting an empty volume +does **not** preserve active assignments, fences, or earlier revocations. Other +deployments must provide equivalent durable Redis (for example, a managed +persistent Redis service and backups). Revocation and +machine enrollment share that state across replicas; do not configure eviction +of authorization keys. If enrollment state is missing, credentials minted under +that enrollment fail closed, and the worker must be explicitly enrolled again. +Restoring a backup from *before* a revocation can revive trust; reconcile +revocations after recovery from backup. Use a distinct server ID for each Code +API deployment and keep it stable when the endpoint changes behind a proxy. + Use `strict` instead of `affinity` if every request must include a runtime session hint. In hardened mode, startup requires the bridge token to be at least 32 bytes. `PTC_MODE=blocking` is rejected; replay mode is required because a @@ -226,7 +262,40 @@ execution. atomically on their first redemption attempt. - Worker credentials expire after fifteen minutes and are bound to an Ed25519 public key. Exact-request signatures include the HTTP method, path, body - digest, timestamp, nonce, and credential. + digest, timestamp, nonce, and credential. With recovery enabled, redeeming a + pairing also persists a separate machine authorization and its public key in + Redis without a TTL by default; an operator can instead set a bounded + enrollment lifetime. +- `POST /v1/bridge/workers/:workerId/credentials/challenge` does not require + an administrator token or an existing access credential, but **does** require + the enrolled key. Its JSON body contains `protocolVersion: 1`, + `operation: "credential.challenge"`, the configured `serverId`, the matching + `workerId`, a fresh UTC ISO `timestamp`, a random 32-byte base64url `nonce`, + and `signature` computed with `signBridgeRecoveryStart(privateKey, fields)` + from `@librechat/code/identity`. Code API verifies the signed fields and + consumes the nonce once before charging the machine's shared challenge + budget; a fabricated request cannot exhaust another worker's budget. +- The response is a short-lived, single-use challenge with the server ID, + worker ID, enrollment generation, operation and expiry. Sign those fields + with `signBridgeRecovery(privateKey, challenge)` and send the fields plus + `signature` to `POST .../credentials/recover` to obtain a new short-lived + credential. Invalid proofs are limited per high-entropy challenge; only + successfully verified, unused proofs consume the machine's shared recovery + budget. Separately, both recovery endpoints limit all incoming requests per + connection peer *before* key verification, including well-formed JSON with + malformed or forged proofs; forged headers and worker IDs cannot bypass + that limit or consume the signed machine budget. All limits live in shared + Redis; HTTP 429 means back off. When a reverse proxy connects to Code API, + its clients share that peer's limit. Restrict direct backend access and apply + client-IP and global + abuse limits at the trusted ingress to keep one proxy peer from becoming a + shared bottleneck; do not trust an arbitrary `X-Forwarded-For` on Code API. +- Recovery and revocation are atomic Redis transitions across API replicas. + A missing, revoked, expired or superseded enrollment never creates new + credentials. Recovery only restores transport authentication. It does not + clear assignment fences, worker or workspace quarantine, or uncertain + execution state. The worker private key is a durable, revocable credential; + expiry of an access credential alone does **not** protect against key theft. - Accepted proof nonces cannot be replayed, credentials rotate before expiry, and an administrator can revoke the active worker identity immediately. - Assignment leases bind to a stable paired identity rather than an individual @@ -239,12 +308,14 @@ execution. worker. The lower API or worker slot ceiling wins, and assignments sharing the same workspace isolation key remain serialized while independent conversation worktrees may run concurrently. -- Workspace tool admission waits for capacity up to the smaller of `JOB_TIMEOUT` - and five minutes while the HTTP caller remains connected. Disconnects cancel - waiting, and admitted work receives a separate execution budget. A shorter - client or proxy timeout can end the wait sooner; Code API does not receive an - absolute caller deadline. See [BYOM worker admission](../byom-worker-admission.md) - for the caller and proxy timeout requirements. +- Workspace tool admission waits for capacity up to 30 seconds without a + `X-LibreChat-Workspace-Queue-Wait-Ms` header. A caller can advertise a longer + per-request allowance, bounded by five minutes and any server queue ceiling. + Disconnects cancel waiting, and admitted work receives a separate execution + budget capped by `JOB_TIMEOUT`. A shorter client or proxy timeout can end the + wait sooner; Code API does not receive an absolute caller deadline. See + [BYOM worker admission](../byom-worker-admission.md) for the total-request + and proxy timeout requirements. - Dynamic workers are fenced to their server-issued tenant before assignment. - Each assignment has an absolute deadline, generation, and random lease token. - Settlements with the wrong worker, generation, token, or expired deadline are diff --git a/docs/remote-bridge/worker-runbook.md b/docs/remote-bridge/worker-runbook.md index 493eb000..aeeef4c4 100644 --- a/docs/remote-bridge/worker-runbook.md +++ b/docs/remote-bridge/worker-runbook.md @@ -488,11 +488,16 @@ it still advertises named environments. ### Expired bridge credential -A running worker refreshes its short-lived credential automatically. If a -machine is offline long enough that refresh can no longer authenticate, issue -a fresh one-time pairing for the same worker ID and redeem it with a newly -generated keypair. Reusing the worker ID preserves the LibreChat environment -record and its agent assignments; creating a new ID creates a new environment. +A running worker refreshes its short-lived credential automatically. With +Code API durable enrollment enabled, a worker that still has its enrolled +private key can request a short-lived challenge and recover a new access +credential without manual re-pairing. The current CLI does **not** yet invoke +that endpoint automatically; update it when worker reconnect support ships. +Until then, or if enrollment is missing or revoked, use the one-time operator +pairing fallback. A new pairing replaces the Code API worker identity and may +require LibreChat environment reauthorization; reusing a worker ID alone does +not guarantee preservation of its LibreChat environment or agent assignments. +Never clear quarantine or workspace fences as part of credential recovery. ### Failed environment setup or uncertain mutation diff --git a/packages/code/README.md b/packages/code/README.md index 87a6fe61..95d86485 100644 --- a/packages/code/README.md +++ b/packages/code/README.md @@ -791,17 +791,22 @@ Legacy requests without a conversation identity continue to use the selected source root. Older Code API deployments do not negotiate the capability, so the worker omits it until every request path understands the isolation boundary. -On an updated Code API, admission waits up to the smaller of `JOB_TIMEOUT` and -five minutes while the HTTP caller remains connected; older Code API versions -waited at most 30 seconds. A `WORKSPACE_QUEUE_TIMEOUT` response (HTTP 503, +On an updated Code API, admission waits up to 30 seconds without the +`X-LibreChat-Workspace-Queue-Wait-Ms` request header. A caller may advertise a +positive integer millisecond allowance up to five minutes, capped by any server +queue ceiling. This allowance is separate from the `JOB_TIMEOUT` execution +budget and cannot outlast a shorter client or proxy timeout. A +`WORKSPACE_QUEUE_TIMEOUT` response (HTTP 503, `Retry-After: 1`) means the operation was not assigned or started; wait for capacity before submitting it again. This is distinct from `ASSIGNMENT_EXPIRED` or a transport timeout after dispatch, where execution may have occurred and mutations must not be blindly retried. No automatic retry is added by this policy. -Align the client's per-attempt timeout and any proxy with the queue **plus** -execution budget before relying on the longer wait. See the +Align the client's per-attempt timeout and every proxy with the queue **plus** +execution, settlement, and delivery budget before relying on a longer wait. +LibreChat keeps the 30-second admission allowance unless its longer total HTTP +budget is explicitly enabled and the live ingress path is verified. See the [BYOM worker admission guide](../../docs/byom-worker-admission.md) for the -current client limitation and the timeout calculations. +timeout calculations. Keep the existing URL, pairing/identity, and network policy configuration. The primary root keeps its configured workspace ID (default `primary`). Repeat diff --git a/packages/code/package.json b/packages/code/package.json index b00ca807..d3cf714b 100644 --- a/packages/code/package.json +++ b/packages/code/package.json @@ -15,6 +15,10 @@ "types": "./dist/protocol.d.ts", "import": "./dist/protocol.js" }, + "./identity": { + "types": "./dist/identity.d.ts", + "import": "./dist/identity.js" + }, "./worker": { "types": "./dist/worker.d.ts", "import": "./dist/worker.js" diff --git a/packages/code/src/identity.test.ts b/packages/code/src/identity.test.ts index 4e92b91f..bef3b673 100644 --- a/packages/code/src/identity.test.ts +++ b/packages/code/src/identity.test.ts @@ -1,9 +1,14 @@ import assert from 'node:assert/strict'; +import { randomBytes } from 'node:crypto'; import test from 'node:test'; import { createBridgeIdentity, + signBridgeRecovery, + signBridgeRecoveryStart, signBridgeRequest, + verifyBridgeRecovery, + verifyBridgeRecoveryStart, verifyBridgeRequest, } from './identity.js'; @@ -33,3 +38,41 @@ test('worker identity proves possession for the exact HTTP request', () => { false, ); }); + +test('signed recovery starts bind operation, server, worker, time and nonce', () => { + const identity = createBridgeIdentity(); + const start = { + operation: 'credential.challenge' as const, + serverId: 'https://code.example.test', + workerId: 'vm-1', + timestamp: new Date().toISOString(), + nonce: randomBytes(32).toString('base64url'), + }; + const signature = signBridgeRecoveryStart(identity.privateKey, start); + assert.equal(verifyBridgeRecoveryStart(identity.publicKey, start, signature), true); + for (const modified of [ + { ...start, operation: 'credential.recover' as 'credential.challenge' }, + { ...start, serverId: 'https://other.example.test' }, + { ...start, workerId: 'vm-2' }, + { ...start, timestamp: new Date(Date.now() + 60_000).toISOString() }, + { ...start, nonce: randomBytes(32).toString('base64url') }, + ]) { + assert.equal(verifyBridgeRecoveryStart(identity.publicKey, modified, signature), false); + } + + const challenge = { + operation: 'credential.recover' as const, + serverId: start.serverId, + workerId: start.workerId, + enrollmentGeneration: randomBytes(18).toString('base64url'), + challenge: randomBytes(32).toString('base64url'), + expiresAt: new Date(Date.now() + 60_000).toISOString(), + }; + assert.equal(verifyBridgeRecovery(identity.publicKey, challenge, signature), false); + assert.equal( + verifyBridgeRecoveryStart( + identity.publicKey, start, signBridgeRecovery(identity.privateKey, challenge), + ), + false, + ); +}); diff --git a/packages/code/src/identity.ts b/packages/code/src/identity.ts index ab11c865..3556520c 100644 --- a/packages/code/src/identity.ts +++ b/packages/code/src/identity.ts @@ -19,6 +19,25 @@ export interface BridgeRequestProofInput { body: string; } +/** Signed without an access credential before requesting a server recovery challenge. */ +export interface BridgeRecoveryStartProofInput { + operation: 'credential.challenge'; + serverId: string; + workerId: string; + timestamp: string; + nonce: string; +} + +/** Signed separately from access-credential requests, so neither proof can be reused for the other. */ +export interface BridgeRecoveryProofInput { + operation: 'credential.recover'; + serverId: string; + workerId: string; + enrollmentGeneration: string; + challenge: string; + expiresAt: string; +} + export function createBridgeIdentity(): BridgeIdentity { const { publicKey, privateKey } = generateKeyPairSync('ed25519', { publicKeyEncoding: { type: 'spki', format: 'pem' }, @@ -64,3 +83,78 @@ export function verifyBridgeRequest( return false; } } + +function canonicalBridgeRecoveryStart(input: BridgeRecoveryStartProofInput): string { + return [ + 'librechat-code:bridge-recovery-start:v1', + input.operation, + input.serverId, + input.workerId, + input.timestamp, + input.nonce, + ].join('\n'); +} + +export function signBridgeRecoveryStart( + privateKey: string, + input: BridgeRecoveryStartProofInput, +): string { + return sign(null, Buffer.from(canonicalBridgeRecoveryStart(input)), privateKey).toString( + 'base64url', + ); +} + +export function verifyBridgeRecoveryStart( + publicKey: string, + input: BridgeRecoveryStartProofInput, + signature: string, +): boolean { + try { + return verify( + null, + Buffer.from(canonicalBridgeRecoveryStart(input)), + publicKey, + Buffer.from(signature, 'base64url'), + ); + } catch { + return false; + } +} + +function canonicalBridgeRecovery(input: BridgeRecoveryProofInput): string { + return [ + 'librechat-code:bridge-recovery:v1', + input.operation, + input.serverId, + input.workerId, + input.enrollmentGeneration, + input.challenge, + input.expiresAt, + ].join('\n'); +} + +export function signBridgeRecovery( + privateKey: string, + input: BridgeRecoveryProofInput, +): string { + return sign(null, Buffer.from(canonicalBridgeRecovery(input)), privateKey).toString( + 'base64url', + ); +} + +export function verifyBridgeRecovery( + publicKey: string, + input: BridgeRecoveryProofInput, + signature: string, +): boolean { + try { + return verify( + null, + Buffer.from(canonicalBridgeRecovery(input)), + publicKey, + Buffer.from(signature, 'base64url'), + ); + } catch { + return false; + } +} diff --git a/packages/code/src/protocol.ts b/packages/code/src/protocol.ts index c678a1af..7657fc9c 100644 --- a/packages/code/src/protocol.ts +++ b/packages/code/src/protocol.ts @@ -759,6 +759,32 @@ export interface BridgeWorkerCredentialResponse { expiresAt: string; } +/** The enrolled machine signs this request before Code API issues a challenge. */ +export interface BridgeRecoveryChallengeRequest { + protocolVersion: BridgeProtocolVersion; + operation: 'credential.challenge'; + serverId: string; + workerId: string; + timestamp: string; + nonce: string; + signature: string; +} + +/** A short-lived, single-use challenge for an already enrolled machine key. */ +export interface BridgeRecoveryChallengeResponse { + protocolVersion: BridgeProtocolVersion; + operation: 'credential.recover'; + serverId: string; + workerId: string; + enrollmentGeneration: string; + challenge: string; + expiresAt: string; +} + +export interface BridgeRecoveryRequest extends BridgeRecoveryChallengeResponse { + signature: string; +} + export interface BridgeSandboxRequest { body: TBody; headers: Record; diff --git a/service/src/bridge/index.ts b/service/src/bridge/index.ts index 16857fad..b04de55f 100644 --- a/service/src/bridge/index.ts +++ b/service/src/bridge/index.ts @@ -11,7 +11,23 @@ export const bridgeStore = new RedisBridgeStore( undefined, env.BRIDGE_MAX_WORKSPACE_LEASE_SLOTS, ); -export const bridgePairings = new RedisBridgePairingStore(connection); +export const bridgePairings = new RedisBridgePairingStore( + connection, + undefined, + undefined, + undefined, + undefined, + env.BRIDGE_RECOVERY_SERVER_ID + ? { + serverId: env.BRIDGE_RECOVERY_SERVER_ID, + enrollmentTtlSeconds: env.BRIDGE_ENROLLMENT_TTL_SECONDS, + challengeTtlSeconds: env.BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS, + maxChallengesPerMinute: env.BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE, + maxAttemptsPerMinute: env.BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE, + maxUntrustedRequestsPerMinute: env.BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE, + } + : undefined, +); export default createBridgeRouter({ enabled: isBridgeEnabled(), diff --git a/service/src/bridge/pairing.ts b/service/src/bridge/pairing.ts index 9eb6d76c..b8a32162 100644 --- a/service/src/bridge/pairing.ts +++ b/service/src/bridge/pairing.ts @@ -6,13 +6,26 @@ import { import type Redis from 'ioredis'; -import { verifyBridgeRequest } from '../../../packages/code/src/identity'; +import { + verifyBridgeRecovery, + verifyBridgeRecoveryStart, + verifyBridgeRequest, +} from '../../../packages/code/src/identity'; +import type { + BridgeRecoveryProofInput, + BridgeRecoveryStartProofInput, +} from '../../../packages/code/src/identity'; const PREFIX = 'codeapi:bridge:v1'; const DEFAULT_PAIRING_TTL_SECONDS = 10 * 60; const DEFAULT_CREDENTIAL_TTL_SECONDS = 15 * 60; const PROOF_NONCE_TTL_SECONDS = 2 * 60; const PROOF_CLOCK_SKEW_MS = 60_000; +const DEFAULT_CHALLENGE_TTL_SECONDS = 60; +const DEFAULT_CHALLENGES_PER_MINUTE = 12; +const DEFAULT_RECOVERY_ATTEMPTS_PER_MINUTE = 30; +const DEFAULT_UNTRUSTED_RECOVERY_REQUESTS_PER_MINUTE = 240; +const RECOVERY_START_NONCE_TTL_SECONDS = 3 * 60; const LEGACY_SCAN_CLAIM_TTL_MS = 5_000; const LEGACY_SCAN_POLL_INTERVAL_MS = 25; const LEGACY_SCAN_PENDING = 'pending'; @@ -41,6 +54,7 @@ elseif (generation or '0') ~= ARGV[2] then redis.call('DEL', KEYS[1]) return 0 end +if ARGV[7] ~= '' and (generation or '0') ~= ARGV[9] then return 0 end redis.call('DEL', KEYS[1]) if redis.call('GET', KEYS[5]) == KEYS[1] then redis.call('DEL', KEYS[5]) @@ -48,9 +62,26 @@ end redis.call('SET', KEYS[3], ARGV[3], 'EX', ARGV[4]) redis.call('SET', KEYS[4], ARGV[5], 'EX', ARGV[4]) redis.call('SET', KEYS[6], ARGV[6], 'EX', ARGV[4]) +if ARGV[7] ~= '' then + if tonumber(ARGV[8]) > 0 then + redis.call('SET', KEYS[7], ARGV[7], 'EX', ARGV[8]) + else + redis.call('SET', KEYS[7], ARGV[7]) + end + redis.call('SET', KEYS[8], ARGV[10]) +else + redis.call('DEL', KEYS[7], KEYS[8]) +end return 1 `; const ROTATE_CREDENTIAL_SCRIPT = ` +if ARGV[6] ~= '' then + if redis.call('GET', KEYS[5]) ~= ARGV[6] or (redis.call('GET', KEYS[6]) or '0') ~= ARGV[7] or redis.call('GET', KEYS[7]) ~= ARGV[8] then + return 0 + end +elseif redis.call('EXISTS', KEYS[7]) == 1 then + return 0 +end local activeDigest = redis.call('GET', KEYS[1]) local previous = redis.call('GET', KEYS[2]) if not activeDigest or not previous then @@ -82,12 +113,45 @@ if credential then redis.call('DEL', KEYS[3]) redis.call('DEL', ARGV[1] .. credential) end -redis.call('DEL', KEYS[1], KEYS[3], KEYS[4], KEYS[5], KEYS[6], KEYS[7]) +redis.call('DEL', KEYS[1], KEYS[3], KEYS[4], KEYS[5], KEYS[6], KEYS[7], KEYS[8]) if activeIncarnation then redis.call('SET', ARGV[2] .. activeIncarnation .. ':fenced', '1') end return 1 `; +const CREATE_RECOVERY_CHALLENGE_SCRIPT = ` +if redis.call('GET', KEYS[1]) ~= ARGV[1] or (redis.call('GET', KEYS[4]) or '0') ~= ARGV[4] or redis.call('GET', KEYS[5]) ~= ARGV[6] then + return 0 +end +if redis.call('EXISTS', KEYS[6]) == 1 then return -2 end +local count = redis.call('INCR', KEYS[2]) +if count == 1 then redis.call('EXPIRE', KEYS[2], 60) end +if count > tonumber(ARGV[2]) then return -1 end +if redis.call('SET', KEYS[3], ARGV[3], 'EX', ARGV[5], 'NX') ~= 'OK' then return 0 end +redis.call('SET', KEYS[6], '1', 'EX', ARGV[7]) +return 1 +`; +const LIMIT_RECOVERY_ATTEMPTS_SCRIPT = ` +local count = redis.call('INCR', KEYS[1]) +if count == 1 then redis.call('EXPIRE', KEYS[1], 60) end +if count > tonumber(ARGV[1]) then return 0 end +return 1 +`; +const COMPLETE_RECOVERY_CHALLENGE_SCRIPT = ` +if redis.call('GET', KEYS[1]) ~= ARGV[1] or redis.call('GET', KEYS[2]) ~= ARGV[2] or (redis.call('GET', KEYS[6]) or '0') ~= ARGV[3] or redis.call('GET', KEYS[7]) ~= ARGV[8] then + return 0 +end +local stableIdentity = redis.call('GET', KEYS[5]) +if stableIdentity and stableIdentity ~= ARGV[4] then return 0 end +local count = redis.call('INCR', KEYS[8]) +if count == 1 then redis.call('EXPIRE', KEYS[8], 60) end +if count > tonumber(ARGV[9]) then return -1 end +redis.call('DEL', KEYS[2]) +redis.call('SET', KEYS[3], ARGV[5], 'EX', ARGV[6]) +redis.call('SET', KEYS[4], ARGV[7], 'EX', ARGV[6]) +redis.call('SET', KEYS[5], ARGV[4], 'EX', ARGV[6]) +return 1 +`; const RELEASE_LEGACY_SCAN_CLAIM_SCRIPT = ` if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) @@ -147,8 +211,34 @@ interface StoredCredential { publicKey: string; expiresAt: string; binding?: BridgeWorkerBinding; + /** Present only for credentials backed by durable machine enrollment. */ + enrollmentGeneration?: string; +} + +interface StoredEnrollment { + workerId: string; + serverId: string; + identityId: string; + generation: string; + pairingGeneration: number; + publicKey: string; + binding?: BridgeWorkerBinding; +} + +export interface BridgeRecoveryOptions { + /** Stable HTTPS origin of this Code API deployment, identical on every replica. */ + serverId: string; + /** Zero (the default) keeps enrollment until explicit revocation. */ + enrollmentTtlSeconds?: number; + challengeTtlSeconds?: number; + maxChallengesPerMinute?: number; + maxAttemptsPerMinute?: number; + /** Shared per-connection-peer cap, separate from machine-signed budgets. */ + maxUntrustedRequestsPerMinute?: number; } +export type BridgeRecoveryChallenge = BridgeRecoveryProofInput; + export interface BridgePairing { workerId: string; code: string; @@ -168,7 +258,10 @@ export class BridgePairingError extends Error { | 'PUBLIC_KEY_INVALID' | 'CREDENTIAL_INVALID' | 'PROOF_INVALID' - | 'PROOF_REPLAYED', + | 'PROOF_REPLAYED' + | 'ENROLLMENT_INVALID' + | 'CHALLENGE_INVALID' + | 'RECOVERY_RATE_LIMITED', message: string, ) { super(message); @@ -196,6 +289,39 @@ function workerStableIdentityKey(workerId: string): string { return `${PREFIX}:stable-identity:${workerId}`; } +function workerEnrollmentKey(workerId: string): string { + return `${PREFIX}:enrollment:${workerId}`; +} + +function workerEnrollmentRequiredKey(workerId: string): string { + return `${PREFIX}:enrollment-required:${workerId}`; +} + +function recoveryChallengeKey(challenge: string): string { + return `${PREFIX}:recovery:challenge:${digest(challenge)}`; +} + +function recoveryStartNonceKey(workerId: string, nonce: string): string { + return `${PREFIX}:recovery:nonce:${workerId}:${digest(nonce)}`; +} + +function recoveryChallengeAttemptKey(challenge: string): string { + return `${PREFIX}:recovery:attempt:${digest(challenge)}`; +} + +function untrustedRecoveryRateKey(peer: string, operation: 'challenge' | 'recover'): string { + // Never partition by a caller-supplied worker ID or an untrusted forwarded IP. + return `${PREFIX}:recovery:untrusted:${operation}:${digest(peer)}`; +} + +function recoveryRateKey( + workerId: string, + enrollmentGeneration: string, + operation: 'start' | 'complete', +): string { + return `${PREFIX}:recovery:rate:${operation}:${workerId}:${enrollmentGeneration}`; +} + function workerPairingGenerationKey(workerId: string): string { return `${PREFIX}:pairing-generation:${workerId}`; } @@ -236,8 +362,120 @@ export class RedisBridgePairingStore { private readonly credentialTtlSeconds = DEFAULT_CREDENTIAL_TTL_SECONDS, private readonly legacyScanClaimTtlMs = LEGACY_SCAN_CLAIM_TTL_MS, private readonly rollbackEpoch = - process.env.CODEAPI_BRIDGE_PAIRING_ROLLBACK_EPOCH?.trim() ?? '', - ) {} + process.env.CODEAPI_BRIDGE_PAIRING_ROLLBACK_EPOCH?.trim() ?? '', + private readonly recovery?: BridgeRecoveryOptions, + ) { + if (recovery == null) return; + let server: URL; + try { + server = new URL(recovery.serverId); + } catch { + throw new Error('Bridge recovery requires a stable HTTPS server origin'); + } + if ( + server.protocol !== 'https:' || + server.username !== '' || + server.password !== '' || + server.pathname !== '/' || + server.search !== '' || + server.hash !== '' || + server.origin !== recovery.serverId + ) { + throw new Error('Bridge recovery requires a stable HTTPS server origin'); + } + for (const [value, max] of [ + [recovery.enrollmentTtlSeconds ?? 0, 10 * 365 * 24 * 3600], + [recovery.challengeTtlSeconds ?? DEFAULT_CHALLENGE_TTL_SECONDS, 300], + [recovery.maxChallengesPerMinute ?? DEFAULT_CHALLENGES_PER_MINUTE, 120], + [recovery.maxAttemptsPerMinute ?? DEFAULT_RECOVERY_ATTEMPTS_PER_MINUTE, 120], + [recovery.maxUntrustedRequestsPerMinute ?? DEFAULT_UNTRUSTED_RECOVERY_REQUESTS_PER_MINUTE, 1200], + ]) { + if (!Number.isSafeInteger(value) || value < 0 || value > max) { + throw new RangeError('Invalid bridge recovery lifetime or rate limit'); + } + } + if ( + (recovery.challengeTtlSeconds ?? DEFAULT_CHALLENGE_TTL_SECONDS) === 0 || + (recovery.maxChallengesPerMinute ?? DEFAULT_CHALLENGES_PER_MINUTE) === 0 || + (recovery.maxAttemptsPerMinute ?? DEFAULT_RECOVERY_ATTEMPTS_PER_MINUTE) === 0 || + (recovery.maxUntrustedRequestsPerMinute ?? DEFAULT_UNTRUSTED_RECOVERY_REQUESTS_PER_MINUTE) === 0 + ) { + throw new RangeError('Bridge recovery challenge lifetime and rate limits must be positive'); + } + } + + get recoveryEnabled(): boolean { + return this.recovery != null; + } + + async limitUntrustedRecovery( + peer: string, + operation: 'challenge' | 'recover', + ): Promise { + if (this.recovery == null) { + throw new BridgePairingError('ENROLLMENT_INVALID', 'Machine recovery is disabled'); + } + const allowed = await this.redis.eval( + LIMIT_RECOVERY_ATTEMPTS_SCRIPT, + 1, + untrustedRecoveryRateKey(peer, operation), + String(this.recovery.maxUntrustedRequestsPerMinute ?? DEFAULT_UNTRUSTED_RECOVERY_REQUESTS_PER_MINUTE), + ); + if (allowed !== 1) { + throw new BridgePairingError('RECOVERY_RATE_LIMITED', 'Too many recovery requests from this peer'); + } + } + + private checkedEnrollment( + workerId: string, + raw: string | null, + pairingGeneration: string | null, + requiredGeneration: string | null, + ): StoredEnrollment { + let parsed: unknown; + try { + parsed = raw == null ? undefined : JSON.parse(raw); + } catch { + parsed = undefined; + } + const enrollment = ( + typeof parsed === 'object' && parsed != null && !Array.isArray(parsed) + ? parsed + : {} + ) as Partial; + if ( + this.recovery == null || + enrollment.workerId !== workerId || + enrollment.serverId !== this.recovery.serverId || + typeof enrollment.identityId !== 'string' || + !/^[A-Za-z0-9_-]{24}$/.test(enrollment.identityId) || + typeof enrollment.generation !== 'string' || + !/^[A-Za-z0-9_-]{24}$/.test(enrollment.generation) || + enrollment.generation !== requiredGeneration || + typeof enrollment.publicKey !== 'string' || + !validEd25519PublicKey(enrollment.publicKey) || + !Number.isSafeInteger(enrollment.pairingGeneration) || + String(enrollment.pairingGeneration) !== (pairingGeneration ?? '0') + ) { + throw new BridgePairingError('ENROLLMENT_INVALID', 'Machine enrollment is unavailable or revoked'); + } + return enrollment as StoredEnrollment; + } + + private checkEnrolledCredential( + credential: StoredCredential, + enrollment: StoredEnrollment, + ): void { + if ( + credential.workerId !== enrollment.workerId || + credential.identityId !== enrollment.identityId || + credential.publicKey !== enrollment.publicKey || + credential.enrollmentGeneration !== enrollment.generation || + JSON.stringify(credential.binding ?? null) !== JSON.stringify(enrollment.binding ?? null) + ) { + throw new BridgePairingError('CREDENTIAL_INVALID', 'Worker credential does not match its machine enrollment'); + } + } async issue( workerId: string, @@ -314,28 +552,49 @@ export class RedisBridgePairingStore { Date.now() + this.credentialTtlSeconds * 1000, ).toISOString(); const identityId = randomBytes(18).toString('base64url'); + const pairingGeneration = pairing.generation ?? Number( + (await this.redis.get(workerPairingGenerationKey(args.workerId))) ?? '0', + ); + const enrollment: StoredEnrollment | undefined = this.recovery == null + ? undefined + : { + workerId: args.workerId, + serverId: this.recovery.serverId, + identityId, + generation: randomBytes(18).toString('base64url'), + pairingGeneration, + publicKey: args.publicKey, + binding: pairing.binding, + }; const stored: StoredCredential = { workerId: args.workerId, identityId, publicKey: args.publicKey, expiresAt, binding: pairing.binding, + ...(enrollment == null ? {} : { enrollmentGeneration: enrollment.generation }), }; const accepted = await this.redis.eval( REDEEM_PAIRING_SCRIPT, - 6, + 8, codeKey, workerPairingGenerationKey(pairing.workerId), credentialDigestKey(credentialDigest), workerIdentityKey(args.workerId), workerPairingIndexKey(args.workerId), workerStableIdentityKey(args.workerId), + workerEnrollmentKey(args.workerId), + workerEnrollmentRequiredKey(args.workerId), raw, pairing.generation == null ? '' : String(pairing.generation), JSON.stringify(stored), String(this.credentialTtlSeconds), credentialDigest, identityId, + enrollment == null ? '' : JSON.stringify(enrollment), + String(this.recovery?.enrollmentTtlSeconds ?? 0), + String(pairingGeneration), + enrollment?.generation ?? '', ); if (accepted !== 1) { throw new BridgePairingError( @@ -374,10 +633,12 @@ export class RedisBridgePairingStore { ); } const credentialDigest = digest(args.credential); - const [raw, activeDigest, pairingGeneration] = await this.redis.mget( + const [raw, activeDigest, pairingGeneration, enrollmentRaw, requiredGeneration] = await this.redis.mget( credentialDigestKey(credentialDigest), workerIdentityKey(args.workerId), workerPairingGenerationKey(args.workerId), + workerEnrollmentKey(args.workerId), + workerEnrollmentRequiredKey(args.workerId), ); if (raw == null || activeDigest == null) { throw new BridgePairingError( @@ -386,6 +647,12 @@ export class RedisBridgePairingStore { ); } const stored = JSON.parse(raw) as StoredCredential; + if (stored.enrollmentGeneration != null || requiredGeneration != null || enrollmentRaw != null) { + this.checkEnrolledCredential( + stored, + this.checkedEnrollment(args.workerId, enrollmentRaw, pairingGeneration, requiredGeneration), + ); + } if (activeDigest !== credentialDigest) { const activeRaw = await this.redis.get( credentialDigestKey(activeDigest), @@ -439,6 +706,177 @@ export class RedisBridgePairingStore { }; } + async createRecoveryChallenge( + workerId: string, + request: BridgeRecoveryStartProofInput, + signature: string, + ): Promise { + if (this.recovery == null) { + throw new BridgePairingError('ENROLLMENT_INVALID', 'Machine recovery is disabled'); + } + const [raw, pairingGeneration, requiredGeneration] = await this.redis.mget( + workerEnrollmentKey(workerId), + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), + ); + const enrollment = this.checkedEnrollment(workerId, raw, pairingGeneration, requiredGeneration); + const proofTime = Date.parse(request.timestamp); + if ( + String(request.operation) !== 'credential.challenge' || + request.serverId !== enrollment.serverId || + request.workerId !== workerId || + !/^[A-Za-z0-9_-]{43}$/.test(request.nonce) || + !Number.isFinite(proofTime) || + Math.abs(Date.now() - proofTime) > PROOF_CLOCK_SKEW_MS || + !verifyBridgeRecoveryStart(enrollment.publicKey, request, signature) + ) { + throw new BridgePairingError('PROOF_INVALID', 'Machine recovery challenge proof is invalid'); + } + const challenge: BridgeRecoveryChallenge = { + operation: 'credential.recover', + serverId: enrollment.serverId, + workerId, + enrollmentGeneration: enrollment.generation, + challenge: randomBytes(32).toString('base64url'), + expiresAt: new Date( + Date.now() + (this.recovery.challengeTtlSeconds ?? DEFAULT_CHALLENGE_TTL_SECONDS) * 1000, + ).toISOString(), + }; + const created = await this.redis.eval( + CREATE_RECOVERY_CHALLENGE_SCRIPT, + 6, + workerEnrollmentKey(workerId), + recoveryRateKey(workerId, enrollment.generation, 'start'), + recoveryChallengeKey(challenge.challenge), + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), + recoveryStartNonceKey(workerId, request.nonce), + raw!, + String(this.recovery.maxChallengesPerMinute ?? DEFAULT_CHALLENGES_PER_MINUTE), + JSON.stringify(challenge), + String(enrollment.pairingGeneration), + String(this.recovery.challengeTtlSeconds ?? DEFAULT_CHALLENGE_TTL_SECONDS), + enrollment.generation, + String(RECOVERY_START_NONCE_TTL_SECONDS), + ); + if (created === -2) { + throw new BridgePairingError('PROOF_REPLAYED', 'Machine recovery challenge proof was already used'); + } + if (created === -1) { + throw new BridgePairingError('RECOVERY_RATE_LIMITED', 'Too many machine recovery challenges'); + } + if (created !== 1) { + throw new BridgePairingError('ENROLLMENT_INVALID', 'Machine enrollment is unavailable or revoked'); + } + return challenge; + } + + async recoverCredential( + workerId: string, + proof: BridgeRecoveryChallenge, + signature: string, + ): Promise { + if (this.recovery == null) { + throw new BridgePairingError('ENROLLMENT_INVALID', 'Machine recovery is disabled'); + } + const challengeKey = recoveryChallengeKey(proof.challenge); + const [enrollmentRaw, challengeRaw, pairingGeneration, requiredGeneration] = await this.redis.mget( + workerEnrollmentKey(workerId), + challengeKey, + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), + ); + const enrollment = this.checkedEnrollment( + workerId, enrollmentRaw, pairingGeneration, requiredGeneration, + ); + let parsed: unknown; + try { + parsed = challengeRaw == null ? undefined : JSON.parse(challengeRaw); + } catch { + parsed = undefined; + } + const saved = ( + typeof parsed === 'object' && parsed != null && !Array.isArray(parsed) + ? parsed + : {} + ) as Partial; + if ( + saved.operation !== 'credential.recover' || + saved.serverId !== enrollment.serverId || + saved.workerId !== workerId || + saved.enrollmentGeneration !== enrollment.generation || + saved.challenge !== proof.challenge || + saved.expiresAt !== proof.expiresAt || + String(proof.operation) !== saved.operation || + saved.serverId !== proof.serverId || + saved.workerId !== proof.workerId || + saved.enrollmentGeneration !== proof.enrollmentGeneration || + typeof saved.expiresAt !== 'string' || + !Number.isFinite(Date.parse(saved.expiresAt)) || + Date.parse(saved.expiresAt) <= Date.now() + ) { + throw new BridgePairingError('CHALLENGE_INVALID', 'Machine recovery challenge is invalid or expired'); + } + // Invalid signatures can only exhaust the high-entropy challenge they know, + // never the enrolled worker's shared quota. The worker-wide counter is + // charged inside completion, after the key proof and replay checks. + const attempts = await this.redis.eval( + LIMIT_RECOVERY_ATTEMPTS_SCRIPT, + 1, + recoveryChallengeAttemptKey(proof.challenge), + String(this.recovery.maxAttemptsPerMinute ?? DEFAULT_RECOVERY_ATTEMPTS_PER_MINUTE), + ); + if (attempts !== 1) { + throw new BridgePairingError('RECOVERY_RATE_LIMITED', 'Too many machine recovery attempts'); + } + if (!verifyBridgeRecovery(enrollment.publicKey, saved as BridgeRecoveryChallenge, signature)) { + throw new BridgePairingError('PROOF_INVALID', 'Machine recovery signature is invalid'); + } + + const credential = randomBytes(32).toString('base64url'); + const credentialDigest = digest(credential); + const expiresAt = new Date(Date.now() + this.credentialTtlSeconds * 1000).toISOString(); + const stored: StoredCredential = { + workerId, + identityId: enrollment.identityId, + publicKey: enrollment.publicKey, + expiresAt, + binding: enrollment.binding, + enrollmentGeneration: enrollment.generation, + }; + // The same Redis decision consumes the proof and issues the credential. + // A revoke or replacement on another replica wins by invalidating the + // enrollment/generation comparison, with no window to recreate trust. + const issued = await this.redis.eval( + COMPLETE_RECOVERY_CHALLENGE_SCRIPT, + 8, + workerEnrollmentKey(workerId), + challengeKey, + credentialDigestKey(credentialDigest), + workerIdentityKey(workerId), + workerStableIdentityKey(workerId), + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), + recoveryRateKey(workerId, enrollment.generation, 'complete'), + enrollmentRaw!, + challengeRaw!, + String(enrollment.pairingGeneration), + enrollment.identityId, + JSON.stringify(stored), + String(this.credentialTtlSeconds), + credentialDigest, + enrollment.generation, + String(this.recovery.maxAttemptsPerMinute ?? DEFAULT_RECOVERY_ATTEMPTS_PER_MINUTE), + ); + if (issued === -1) { + throw new BridgePairingError('RECOVERY_RATE_LIMITED', 'Too many machine recoveries'); + } + if (issued !== 1) { + throw new BridgePairingError('CHALLENGE_INVALID', 'Machine recovery challenge is invalid or expired'); + } + return { workerId, credential, expiresAt }; + } + async revoke(workerId: string): Promise { await this.removeLegacyPairings(workerId); // Fence redemption and consume the currently indexed code atomically. An @@ -446,7 +884,7 @@ export class RedisBridgePairingStore { // that linearizes afterward installs a distinct generation and code. await this.redis.eval( REVOKE_PAIRING_SCRIPT, - 7, + 8, workerPairingIndexKey(workerId), workerPairingGenerationKey(workerId), workerIdentityKey(workerId), @@ -454,6 +892,7 @@ export class RedisBridgePairingStore { `${PREFIX}:worker:${encodeURIComponent(workerId)}`, `${PREFIX}:worker:${encodeURIComponent(workerId)}:incarnation`, `${PREFIX}:worker:${encodeURIComponent(workerId)}:ready`, + workerEnrollmentKey(workerId), `${PREFIX}:credential:`, `${PREFIX}:worker:${encodeURIComponent(workerId)}:incarnation:`, ); @@ -652,12 +1091,26 @@ export class RedisBridgePairingStore { ); } const previous = JSON.parse(previousRaw) as StoredCredential; + let enrollment: StoredEnrollment | undefined; + let enrollmentRaw: string | null = null; + const [raw, generation, requiredGeneration] = await this.redis.mget( + workerEnrollmentKey(workerId), + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), + ); + if (previous.enrollmentGeneration != null || requiredGeneration != null || raw != null) { + enrollment = this.checkedEnrollment(workerId, raw, generation, requiredGeneration); + this.checkEnrolledCredential(previous, enrollment); + enrollmentRaw = raw; + } return await this.issueCredential( workerId, previous.publicKey, previousDigest, previous.binding, previous.identityId ?? null, + enrollment, + enrollmentRaw, ); } @@ -667,6 +1120,8 @@ export class RedisBridgePairingStore { previousDigest?: string, binding?: BridgeWorkerBinding, identityId?: string | null, + enrollment?: StoredEnrollment, + enrollmentRaw?: string | null, ): Promise { const credential = randomBytes(32).toString('base64url'); const credentialDigest = digest(credential); @@ -683,20 +1138,27 @@ export class RedisBridgePairingStore { publicKey, expiresAt, binding, + ...(enrollment == null ? {} : { enrollmentGeneration: enrollment.generation }), }; if (previousDigest !== undefined) { const rotated = await this.redis.eval( ROTATE_CREDENTIAL_SCRIPT, - 4, + 7, workerIdentityKey(workerId), credentialDigestKey(previousDigest), credentialDigestKey(credentialDigest), workerStableIdentityKey(workerId), + workerEnrollmentKey(workerId), + workerPairingGenerationKey(workerId), + workerEnrollmentRequiredKey(workerId), previousDigest, credentialDigest, JSON.stringify(stored), String(this.credentialTtlSeconds), stableIdentityId ?? '', + enrollmentRaw ?? '', + String(enrollment?.pairingGeneration ?? ''), + enrollment?.generation ?? '', ); if (rotated !== 1) { throw new BridgePairingError( diff --git a/service/src/bridge/recovery-router.test.ts b/service/src/bridge/recovery-router.test.ts new file mode 100644 index 00000000..28ea4b1e --- /dev/null +++ b/service/src/bridge/recovery-router.test.ts @@ -0,0 +1,213 @@ +import { randomBytes } from 'crypto'; +import { createServer, type Server } from 'http'; + +import { afterEach, describe, expect, test } from 'bun:test'; +import express, { json } from 'express'; +import RedisMock from 'ioredis-mock'; + +import type Redis from 'ioredis'; +import type { BridgeRecoveryChallengeRequest, BridgeRecoveryChallengeResponse } from '../../../packages/code/src/protocol'; + +import { + createBridgeIdentity, + signBridgeRecovery, + signBridgeRecoveryStart, +} from '../../../packages/code/src/identity'; +import { BRIDGE_PROTOCOL_VERSION } from '../../../packages/code/src/protocol'; +import { RedisBridgePairingStore } from './pairing'; +import { createBridgeRouter } from './router'; +import { RedisBridgeStore } from './store'; + +const redis = new RedisMock() as unknown as Redis; +const workerId = 'http-recovery-worker'; +const serverId = 'https://code.example.test'; +const servers: Server[] = []; + +function signedStart(privateKey: string): BridgeRecoveryChallengeRequest { + const request = { + operation: 'credential.challenge' as const, + serverId, + workerId, + timestamp: new Date().toISOString(), + nonce: randomBytes(32).toString('base64url'), + }; + return { + protocolVersion: BRIDGE_PROTOCOL_VERSION, + ...request, + signature: signBridgeRecoveryStart(privateKey, request), + }; +} + +async function startRouter(pairings: RedisBridgePairingStore): Promise { + const app = express(); + app.set('trust proxy', 1); + app.use(json()); + app.use('/v1/bridge', createBridgeRouter({ + store: new RedisBridgeStore(redis), + pairings, + authMode: 'paired', + adminToken: 'operator-only', + configuredWorkerId: workerId, + })); + const server = createServer(app); + servers.push(server); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (address == null || typeof address === 'string') throw new Error('Expected TCP listener'); + return `http://127.0.0.1:${address.port}/v1/bridge/workers/${workerId}/credentials`; +} + +function post( + url: string, + body: object, + headers: Record = {}, +): Promise { + return fetch(url, { + method: 'POST', + headers: { 'Content-Type': 'application/json', ...headers }, + body: JSON.stringify(body), + }); +} + +afterEach(async () => { + await Promise.all(servers.splice(0).map(server => new Promise((resolve, reject) => { + server.close(error => error ? reject(error) : resolve()); + }))); + await redis.flushall(); +}); + +describe('machine credential recovery HTTP API', () => { + test('does not expose recovery until the server identity is enabled', async () => { + const baseUrl = await startRouter(new RedisBridgePairingStore(redis)); + const response = await post(`${baseUrl}/challenge`, { protocolVersion: BRIDGE_PROTOCOL_VERSION }); + expect(response.status).toBe(404); + }); + + test('recovers using the stored key without administrator authentication and never leaks the binding', async () => { + const pairings = new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId, + }); + const identity = createBridgeIdentity(); + const pairing = await pairings.issue(workerId, { + tenantId: 'tenant-one', principal: { type: 'user', id: 'owner-one' }, + }); + await pairings.redeem({ workerId, code: pairing.code, publicKey: identity.publicKey }); + const baseUrl = await startRouter(pairings); + const unsigned = await post(`${baseUrl}/challenge`, { + protocolVersion: BRIDGE_PROTOCOL_VERSION, + }); + expect(unsigned.status).toBe(400); + const outsider = createBridgeIdentity(); + const forged = await post(`${baseUrl}/challenge`, signedStart(outsider.privateKey)); + expect(forged.status).toBe(401); + await expect(forged.json()).resolves.toMatchObject({ code: 'PROOF_INVALID' }); + + const challengeResponse = await post(`${baseUrl}/challenge`, signedStart(identity.privateKey)); + expect(challengeResponse.status).toBe(200); + const challenge = (await challengeResponse.json()) as BridgeRecoveryChallengeResponse; + expect(challenge).toMatchObject({ + protocolVersion: BRIDGE_PROTOCOL_VERSION, + serverId, + workerId, + operation: 'credential.recover', + }); + expect(JSON.stringify(challenge)).not.toMatch(/tenant-one|owner-one|privateKey/); + + const signature = signBridgeRecovery(identity.privateKey, challenge); + const rejected = await post(`${baseUrl}/recover`, { + ...challenge, + serverId: 'https://unrelated.example.test', + signature, + }); + expect(rejected.status).toBe(401); + await expect(rejected.json()).resolves.toMatchObject({ code: 'CHALLENGE_INVALID' }); + + const recovered = await post(`${baseUrl}/recover`, { ...challenge, signature }); + expect(recovered.status).toBe(200); + const credential = (await recovered.json()) as { workerId: string; credential: string }; + expect(credential.workerId).toBe(workerId); + expect(credential.credential.length).toBeGreaterThanOrEqual(32); + const replay = await post(`${baseUrl}/recover`, { ...challenge, signature }); + expect(replay.status).toBe(401); + await expect(replay.json()).resolves.toMatchObject({ code: 'CHALLENGE_INVALID' }); + }); + + test('limits forged starts before signature work across replicas despite spoofed proxy headers', async () => { + const first = new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId, maxUntrustedRequestsPerMinute: 2, + }); + const second = new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId, maxUntrustedRequestsPerMinute: 2, + }); + const identity = createBridgeIdentity(); + const attacker = createBridgeIdentity(); + const pairing = await first.issue(workerId); + await first.redeem({ workerId, code: pairing.code, publicKey: identity.publicKey }); + const firstUrl = await startRouter(first); + const secondUrl = await startRouter(second); + let signatureChecks = 0; + for (const store of [first, second]) { + const original = store.createRecoveryChallenge.bind(store); + store.createRecoveryChallenge = async ( + ...args + ): ReturnType => { + signatureChecks += 1; + return original(...args); + }; + } + + const forged = async (url: string, suffix: number): Promise => post( + `${url}/challenge`, signedStart(attacker.privateKey), + { 'X-Forwarded-For': `198.51.100.${suffix}` }, + ); + const firstResponse = await forged(firstUrl, 1); + const secondResponse = await forged(secondUrl, 2); + const blocked = await forged( + firstUrl.replace(`/workers/${workerId}/`, '/workers/attacker-selected-worker/'), + 3, + ); + expect([firstResponse.status, secondResponse.status, blocked.status]).toEqual([401, 401, 429]); + await expect(blocked.json()).resolves.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + expect(Number(blocked.headers.get('retry-after'))).toBeGreaterThan(0); + expect(signatureChecks).toBe(2); + expect(await redis.keys('codeapi:bridge:v1:recovery:rate:start:*')).toEqual([]); + + // Bogus completions use another untrusted bucket, not the worker's signed quota. + const fake = { + protocolVersion: BRIDGE_PROTOCOL_VERSION, + operation: 'credential.recover', + serverId, + workerId, + enrollmentGeneration: 'a'.repeat(24), + challenge: randomBytes(32).toString('base64url'), + expiresAt: new Date(Date.now() + 60_000).toISOString(), + signature: 'a'.repeat(86), + }; + const rejected = await post(`${secondUrl}/recover`, fake); + expect(rejected.status).toBe(401); + expect(await redis.keys('codeapi:bridge:v1:recovery:rate:complete:*')).toEqual([]); + const recovered = await post(`${secondUrl}/recover`, fake); + expect(recovered.status).toBe(401); + const blockedCompletion = await post(`${secondUrl}/recover`, fake); + expect(blockedCompletion.status).toBe(429); + await expect(blockedCompletion.json()).resolves.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + }); + + test('limits recovery challenges across two API routers sharing Redis', async () => { + const first = new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId: 'https://code.example.test', maxChallengesPerMinute: 1, + }); + const second = new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId: 'https://code.example.test', maxChallengesPerMinute: 1, + }); + const identity = createBridgeIdentity(); + const pairing = await first.issue(workerId); + await first.redeem({ workerId, code: pairing.code, publicKey: identity.publicKey }); + const baseUrl = await startRouter(second); + const initial = signedStart(identity.privateKey); + await first.createRecoveryChallenge(workerId, initial, initial.signature); + const denied = await post(`${baseUrl}/challenge`, signedStart(identity.privateKey)); + expect(denied.status).toBe(429); + await expect(denied.json()).resolves.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + }); +}); diff --git a/service/src/bridge/recovery.test.ts b/service/src/bridge/recovery.test.ts new file mode 100644 index 00000000..443e7277 --- /dev/null +++ b/service/src/bridge/recovery.test.ts @@ -0,0 +1,518 @@ +import { afterEach, describe, expect, test } from 'bun:test'; +import { createHash, randomBytes } from 'crypto'; +import RedisMock from 'ioredis-mock'; + +import type Redis from 'ioredis'; +import type { + BridgeRecoveryProofInput, + BridgeRecoveryStartProofInput, +} from '../../../packages/code/src/identity'; +import type { + BridgeRecoveryChallenge, + BridgeRecoveryOptions, + BridgeWorkerBinding, + BridgeWorkerCredential, +} from './pairing'; + +import { + createBridgeIdentity, + signBridgeRecovery, + signBridgeRecoveryStart, + signBridgeRequest, +} from '../../../packages/code/src/identity'; +import { RedisBridgePairingStore } from './pairing'; +import { RedisBridgeStore } from './store'; + +const redis = new RedisMock() as unknown as Redis; +const serverId = 'https://code.example.test'; +const workerId = 'durable-worker'; +const binding: BridgeWorkerBinding = { + tenantId: 'tenant-one', + principal: { type: 'user', id: 'owner-one' }, +}; + +function recoverableStore(options: Partial = {}): RedisBridgePairingStore { + return new RedisBridgePairingStore(redis, 600, 300, 5_000, '', { + serverId, + ...options, + }); +} + +function authorizedRequest( + privateKey: string, + credential: string, + nonce: string, +): Parameters[0] { + const proof = { + credential, + method: 'POST', + path: '/v1/bridge/workers/register', + timestamp: new Date().toISOString(), + nonce, + body: JSON.stringify({ protocolVersion: 1, workerId }), + }; + return { + ...proof, + workerId, + signature: signBridgeRequest(privateKey, proof), + }; +} + +async function enroll( + store: RedisBridgePairingStore, + publicKey: string, +): Promise { + const pairing = await store.issue(workerId, binding); + return store.redeem({ workerId, code: pairing.code, publicKey }); +} + +function recoveryStart( + privateKey: string, + nonce = randomBytes(32).toString('base64url'), + overrides: Partial = {}, +): { proof: BridgeRecoveryStartProofInput; signature: string } { + const proof: BridgeRecoveryStartProofInput = { + operation: 'credential.challenge', + serverId, + workerId, + timestamp: new Date().toISOString(), + nonce, + ...overrides, + }; + return { proof, signature: signBridgeRecoveryStart(privateKey, proof) }; +} + +async function challengeFor( + store: RedisBridgePairingStore, + privateKey: string, +): Promise { + const { proof, signature } = recoveryStart(privateKey); + return store.createRecoveryChallenge(workerId, proof, signature); +} + +async function recover( + store: RedisBridgePairingStore, + privateKey: string, +): Promise<{ challenge: BridgeRecoveryChallenge; credential: BridgeWorkerCredential }> { + const challenge = await challengeFor(store, privateKey); + const credential = await store.recoverCredential( + workerId, + challenge, + signBridgeRecovery(privateKey, challenge), + ); + return { challenge, credential }; +} + +afterEach(async () => { + await redis.flushall(); +}); + +describe('durable bridge enrollment', () => { + test('keeps legacy pairing and refresh compatible until recovery is enabled', async () => { + const legacy = new RedisBridgePairingStore(redis); + const identity = createBridgeIdentity(); + const issued = await enroll(legacy, identity.publicKey); + expect(await redis.get(`codeapi:bridge:v1:enrollment:${workerId}`)).toBeNull(); + + const replica = recoverableStore(); + const { proof, signature } = recoveryStart(identity.privateKey); + await expect(replica.createRecoveryChallenge(workerId, proof, signature)) + .rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + const rotated = await replica.rotate(workerId); + await expect( + replica.authorize(authorizedRequest(identity.privateKey, rotated.credential, 'legacy-proof')), + ).resolves.toMatchObject({ workerId, binding }); + expect(issued.credential).not.toBe(rotated.credential); + }); + + test('recovers after access expiry and restart with the same identity and binding', async () => { + const identity = createBridgeIdentity(); + const firstReplica = recoverableStore(); + const issued = await enroll(firstReplica, identity.publicKey); + const original = await firstReplica.authorize( + authorizedRequest(identity.privateKey, issued.credential, 'before-outage'), + ); + const digest = createHash('sha256').update(issued.credential).digest('hex'); + await redis.del( + `codeapi:bridge:v1:credential:${digest}`, + `codeapi:bridge:v1:identity:${workerId}`, + `codeapi:bridge:v1:stable-identity:${workerId}`, + ); + await expect(firstReplica.rotate(workerId)).rejects.toMatchObject({ + code: 'CREDENTIAL_INVALID', + }); + + const restartedReplica = recoverableStore(); + const { challenge, credential } = await recover(restartedReplica, identity.privateKey); + expect(challenge).toMatchObject({ + operation: 'credential.recover', + serverId, + workerId, + }); + const restored = await firstReplica.authorize( + authorizedRequest(identity.privateKey, credential.credential, 'after-outage'), + ); + expect(restored).toMatchObject({ identityId: original.identityId, binding }); + expect(credential.expiresAt).toBeString(); + const rotated = await restartedReplica.rotate(workerId, restored.credentialId); + await expect(firstReplica.authorize( + authorizedRequest(identity.privateKey, rotated.credential, 'after-rotation'), + )).resolves.toMatchObject({ identityId: original.identityId, binding }); + }); + + test('requires the enrolled key and binds proofs to server, worker, generation, operation, and expiry', async () => { + const enrolled = createBridgeIdentity(); + const outsider = createBridgeIdentity(); + const store = recoverableStore(); + await enroll(store, enrolled.publicKey); + const challenge = await challengeFor(store, enrolled.privateKey); + + await expect(store.recoverCredential( + workerId, challenge, signBridgeRecovery(outsider.privateKey, challenge), + )).rejects.toMatchObject({ code: 'PROOF_INVALID' }); + for (const modified of [ + { ...challenge, serverId: 'https://other.example.test' }, + { ...challenge, workerId: 'another-worker' }, + { ...challenge, enrollmentGeneration: '0'.repeat(24) }, + { ...challenge, operation: 'credential.recover-other' as 'credential.recover' }, + { ...challenge, expiresAt: new Date(Date.now() + 120_000).toISOString() }, + ]) { + await expect(store.recoverCredential( + workerId, + modified, + signBridgeRecovery(enrolled.privateKey, modified as BridgeRecoveryProofInput), + )).rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + } + const signature = signBridgeRecovery(enrolled.privateKey, challenge); + await store.recoverCredential(workerId, challenge, signature); + await expect(store.recoverCredential(workerId, challenge, signature)) + .rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + }); + + test('rejects a signed challenge after its short-lived Redis window expires', async () => { + const identity = createBridgeIdentity(); + const store = recoverableStore({ challengeTtlSeconds: 1 }); + await enroll(store, identity.publicKey); + const challenge = await challengeFor(store, identity.privateKey); + await new Promise((resolve) => setTimeout(resolve, 1_100)); + await expect(store.recoverCredential( + workerId, challenge, signBridgeRecovery(identity.privateKey, challenge), + )).rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + }); + + test('retries a lost recovery response without changing the enrolled identity', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore(); + const originalCredential = await enroll(first, identity.publicKey); + const original = await first.authorize( + authorizedRequest(identity.privateKey, originalCredential.credential, 'before-lost-response'), + ); + await recover(first, identity.privateKey); // The caller lost this credential response. + const second = recoverableStore(); + const retried = await recover(second, identity.privateKey); + const auth = await second.authorize( + authorizedRequest(identity.privateKey, retried.credential.credential, 'after-lost-response'), + ); + expect(auth.identityId).toBe(original.identityId); + expect(auth.binding).toEqual(binding); + const enrolled = JSON.parse((await redis.get(`codeapi:bridge:v1:enrollment:${workerId}`))!) as { + generation: string; + }; + expect(await redis.get(`codeapi:bridge:v1:enrollment-required:${workerId}`)) + .toBe(enrolled.generation); + }); + + test('rejects expired and missing enrollment even when an access credential remains live', async () => { + const identity = createBridgeIdentity(); + const store = recoverableStore({ enrollmentTtlSeconds: 1 }); + const issued = await enroll(store, identity.publicKey); + await new Promise((resolve) => setTimeout(resolve, 1_100)); + await expect(challengeFor(store, identity.privateKey)) + .rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + await expect(store.rotate(workerId)) + .rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + await expect(store.authorize( + authorizedRequest(identity.privateKey, issued.credential, 'expired-enrollment'), + )).rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + }); + + test('never treats an unmarked access token as legacy after authorization state is lost', async () => { + const identity = createBridgeIdentity(); + const store = recoverableStore(); + const issued = await enroll(store, identity.publicKey); + const credentialKey = `codeapi:bridge:v1:credential:${createHash('sha256').update(issued.credential).digest('hex')}`; + const raw = JSON.parse((await redis.get(credentialKey))!) as { enrollmentGeneration?: string }; + delete raw.enrollmentGeneration; // Simulate a refresh by a pre-recovery replica. + await redis.set(credentialKey, JSON.stringify(raw), 'EX', 300); + await redis.del(`codeapi:bridge:v1:enrollment:${workerId}`); + + expect(await redis.get(`codeapi:bridge:v1:enrollment-required:${workerId}`)).not.toBeNull(); + await expect(store.authorize( + authorizedRequest(identity.privateKey, issued.credential, 'lost-authorization'), + )).rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + await expect(store.rotate(workerId)).rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + await expect(challengeFor(store, identity.privateKey)) + .rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + await redis.set(`codeapi:bridge:v1:enrollment:${workerId}`, 'null'); + await expect(challengeFor(store, identity.privateKey)) + .rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + }); + + test('re-pairing explicitly supersedes a previous key and binds to the new principal', async () => { + const first = createBridgeIdentity(); + const second = createBridgeIdentity(); + const store = recoverableStore(); + const initial = await enroll(store, first.publicKey); + const pending = await challengeFor(store, first.privateKey); + const replacement = await store.issue(workerId, { + tenantId: 'tenant-two', principal: { type: 'user', id: 'owner-two' }, + }); + await store.redeem({ workerId, code: replacement.code, publicKey: second.publicKey }); + + await expect(store.recoverCredential( + workerId, pending, signBridgeRecovery(first.privateKey, pending), + )).rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + await expect(store.authorize( + authorizedRequest(first.privateKey, initial.credential, 'superseded-key'), + )).rejects.toMatchObject({ code: 'CREDENTIAL_INVALID' }); + const { credential } = await recover(store, second.privateKey); + await expect(store.authorize( + authorizedRequest(second.privateKey, credential.credential, 'new-owner'), + )).resolves.toMatchObject({ + binding: { tenantId: 'tenant-two', principal: { type: 'user', id: 'owner-two' } }, + }); + }); + + test('concurrent replicas consume a recovery challenge and count it only once', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore({ maxAttemptsPerMinute: 2 }); + const second = recoverableStore({ maxAttemptsPerMinute: 2 }); + await enroll(first, identity.publicKey); + const challenge = await challengeFor(first, identity.privateKey); + const signature = signBridgeRecovery(identity.privateKey, challenge); + const attempts = await Promise.allSettled([ + first.recoverCredential(workerId, challenge, signature), + second.recoverCredential(workerId, challenge, signature), + ]); + const issued = attempts.filter( + (result): result is PromiseFulfilledResult => + result.status === 'fulfilled', + ); + const rejected = attempts.filter( + (result): result is PromiseRejectedResult => result.status === 'rejected', + ); + expect(issued).toHaveLength(1); + expect(rejected).toHaveLength(1); + expect(rejected[0]).toMatchObject({ reason: { code: 'CHALLENGE_INVALID' } }); + const digest = createHash('sha256').update(issued[0].value.credential).digest('hex'); + expect(await redis.get(`codeapi:bridge:v1:identity:${workerId}`)).toBe(digest); + expect(await redis.get(`codeapi:bridge:v1:recovery:rate:complete:${workerId}:${challenge.enrollmentGeneration}`)) + .toBe('1'); + }); + + test('revocation beats a signed recovery pending on a different replica', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore(); + const second = recoverableStore(); + await enroll(first, identity.publicKey); + const challenge = await challengeFor(first, identity.privateKey); + const originalEval = redis.eval.bind(redis); + let release!: () => void; + let enter!: () => void; + const paused = new Promise((resolve) => { enter = resolve; }); + const resume = new Promise((resolve) => { release = resolve; }); + redis.eval = (async (script: string, ...args: unknown[]) => { + if (script.includes('local stableIdentity = redis.call')) { + enter(); + await resume; + } + return (originalEval as (...evalArgs: unknown[]) => Promise)(script, ...args); + }) as Redis['eval']; + try { + const pending = first.recoverCredential( + workerId, challenge, signBridgeRecovery(identity.privateKey, challenge), + ); + await paused; + await second.revoke(workerId); + release(); + await expect(pending).rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + expect(await redis.get(`codeapi:bridge:v1:identity:${workerId}`)).toBeNull(); + expect(await redis.get(`codeapi:bridge:v1:enrollment:${workerId}`)).toBeNull(); + } finally { + redis.eval = originalEval as Redis['eval']; + release(); + } + }); + + test('isolates untrusted recovery limits by socket peer and endpoint from signed machine quotas', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore({ + maxUntrustedRequestsPerMinute: 1, maxChallengesPerMinute: 1, maxAttemptsPerMinute: 1, + }); + const second = recoverableStore({ + maxUntrustedRequestsPerMinute: 1, maxChallengesPerMinute: 1, maxAttemptsPerMinute: 1, + }); + await enroll(first, identity.publicKey); + await first.limitUntrustedRecovery('192.0.2.1', 'challenge'); + await expect(second.limitUntrustedRecovery('192.0.2.1', 'challenge')) + .rejects.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + await expect(second.limitUntrustedRecovery('192.0.2.2', 'challenge')) + .resolves.toBeUndefined(); + await expect(second.limitUntrustedRecovery('192.0.2.1', 'recover')) + .resolves.toBeUndefined(); + expect((await redis.keys('codeapi:bridge:v1:recovery:rate:start:*')).length).toBe(0); + const { credential } = await recover(second, identity.privateKey); + expect(credential.workerId).toBe(workerId); + }); + + test('only a signed, fresh, unused start request consumes the worker challenge budget', async () => { + const identity = createBridgeIdentity(); + const outsider = createBridgeIdentity(); + const first = recoverableStore({ maxChallengesPerMinute: 2 }); + const second = recoverableStore({ maxChallengesPerMinute: 2 }); + await enroll(first, identity.publicKey); + + for (let index = 0; index < 20; index += 1) { + const forged = recoveryStart(outsider.privateKey); + await expect(first.createRecoveryChallenge(workerId, forged.proof, forged.signature)) + .rejects.toMatchObject({ code: 'PROOF_INVALID' }); + } + const otherServer = recoveryStart(identity.privateKey, undefined, { + serverId: 'https://unrelated.example.test', + }); + await expect(first.createRecoveryChallenge(workerId, otherServer.proof, otherServer.signature)) + .rejects.toMatchObject({ code: 'PROOF_INVALID' }); + const stale = recoveryStart(identity.privateKey, undefined, { + timestamp: new Date(Date.now() - 5 * 60_000).toISOString(), + }); + await expect(first.createRecoveryChallenge(workerId, stale.proof, stale.signature)) + .rejects.toMatchObject({ code: 'PROOF_INVALID' }); + + const valid = recoveryStart(identity.privateKey); + await first.createRecoveryChallenge(workerId, valid.proof, valid.signature); + await expect(second.createRecoveryChallenge(workerId, valid.proof, valid.signature)) + .rejects.toMatchObject({ code: 'PROOF_REPLAYED' }); + await challengeFor(second, identity.privateKey); + await expect(challengeFor(first, identity.privateKey)) + .rejects.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + }); + + test('fabricated completions never spend the worker budget or block a fresh signed recovery', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore({ maxAttemptsPerMinute: 1 }); + const second = recoverableStore({ maxAttemptsPerMinute: 1 }); + await enroll(first, identity.publicKey); + const real = await challengeFor(first, identity.privateKey); + + for (let index = 0; index < 35; index += 1) { + const fake = { ...real, challenge: randomBytes(32).toString('base64url') }; + await expect(first.recoverCredential(workerId, fake, 'forged')) + .rejects.toMatchObject({ code: 'CHALLENGE_INVALID' }); + } + await first.recoverCredential(workerId, real, signBridgeRecovery(identity.privateKey, real)); + const next = await challengeFor(second, identity.privateKey); + await expect(second.recoverCredential( + workerId, next, signBridgeRecovery(identity.privateKey, next), + )).rejects.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + }); + + test('invalid signatures exhaust only their own challenge, not the enrolled worker', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore({ maxAttemptsPerMinute: 1 }); + const second = recoverableStore({ maxAttemptsPerMinute: 1 }); + await enroll(first, identity.publicKey); + const attacked = await challengeFor(first, identity.privateKey); + await expect(first.recoverCredential(workerId, attacked, 'forged')) + .rejects.toMatchObject({ code: 'PROOF_INVALID' }); + await expect(second.recoverCredential( + workerId, attacked, signBridgeRecovery(identity.privateKey, attacked), + )).rejects.toMatchObject({ code: 'RECOVERY_RATE_LIMITED' }); + const fresh = await challengeFor(second, identity.privateKey); + await expect(first.recoverCredential( + workerId, fresh, signBridgeRecovery(identity.privateKey, fresh), + )).resolves.toMatchObject({ workerId }); + }); + + test('key replacement starts new signed challenge and recovery budgets', async () => { + const oldKey = createBridgeIdentity(); + const newKey = createBridgeIdentity(); + const store = recoverableStore({ maxChallengesPerMinute: 1, maxAttemptsPerMinute: 1 }); + await enroll(store, oldKey.publicKey); + const oldChallenge = await challengeFor(store, oldKey.privateKey); + await store.recoverCredential( + workerId, oldChallenge, signBridgeRecovery(oldKey.privateKey, oldChallenge), + ); + const replacement = await store.issue(workerId, binding); + await store.redeem({ workerId, code: replacement.code, publicKey: newKey.publicKey }); + + await expect(challengeFor(store, oldKey.privateKey)) + .rejects.toMatchObject({ code: 'PROOF_INVALID' }); + const newChallenge = await challengeFor(store, newKey.privateKey); + await expect(store.recoverCredential( + workerId, newChallenge, signBridgeRecovery(newKey.privateKey, newChallenge), + )).resolves.toMatchObject({ workerId }); + }); + + test('revocation fences a previously verified start proof before a challenge is written', async () => { + const identity = createBridgeIdentity(); + const first = recoverableStore(); + const second = recoverableStore(); + await enroll(first, identity.publicKey); + const { proof, signature } = recoveryStart(identity.privateKey); + const originalEval = redis.eval.bind(redis); + let release!: () => void; + let enter!: () => void; + const paused = new Promise((resolve) => { enter = resolve; }); + const resume = new Promise((resolve) => { release = resolve; }); + redis.eval = (async (script: string, ...args: unknown[]) => { + if (script.includes('KEYS[6]) == 1 then return -2')) { + enter(); + await resume; + } + return (originalEval as (...evalArgs: unknown[]) => Promise)(script, ...args); + }) as Redis['eval']; + try { + const pending = first.createRecoveryChallenge(workerId, proof, signature); + await paused; + await second.revoke(workerId); + release(); + await expect(pending).rejects.toMatchObject({ code: 'ENROLLMENT_INVALID' }); + expect(await redis.keys('codeapi:bridge:v1:recovery:challenge:*')).toEqual([]); + } finally { + redis.eval = originalEval as Redis['eval']; + release(); + } + }); + + test('recovering credentials never clears worker or workspace quarantine', async () => { + const identity = createBridgeIdentity(); + const store = recoverableStore(); + await enroll(store, identity.publicKey); + const workerQuarantine = `codeapi:bridge:v1:worker:${workerId}:incarnation:incarnation-00000001:quarantined`; + const workspaceQuarantine = `codeapi:bridge:v1:worker:${workerId}:workspace:${createHash('sha256').update('session-one').digest('hex')}:quarantined`; + await redis.set(workerQuarantine, '1'); + await redis.set(workspaceQuarantine, '1'); + const { credential } = await recover(store, identity.privateKey); + expect(await redis.get(workerQuarantine)).toBe('1'); + expect(await redis.get(workspaceQuarantine)).toBe('1'); + const auth = await store.authorize( + authorizedRequest(identity.privateKey, credential.credential, 'quarantined-machine'), + ); + await expect(new RedisBridgeStore(redis).register({ + protocolVersion: 1, + workerId, + incarnationId: 'incarnation-00000001', + capabilities: { statefulWorkspace: true, sandboxProfile: 'nsjail', runtimes: ['bash'] }, + }, auth)).rejects.toMatchObject({ code: 'WORKER_QUARANTINED' }); + }); + + test('refuses non-HTTPS or non-origin deployment identity and unbounded recovery policy', () => { + for (const invalid of ['http://code.example.test', 'https://code.example.test/path', 'https://user@code.example.test']) { + expect(() => recoverableStore({ serverId: invalid })).toThrow(); + } + expect(() => recoverableStore({ maxAttemptsPerMinute: 0 })).toThrow(); + expect(() => recoverableStore({ maxUntrustedRequestsPerMinute: 0 })).toThrow(); + expect(() => recoverableStore({ maxUntrustedRequestsPerMinute: 1201 })).toThrow(); + expect(() => recoverableStore({ challengeTtlSeconds: 301 })).toThrow(); + }); +}); diff --git a/service/src/bridge/router.ts b/service/src/bridge/router.ts index 5b319733..b3d549b1 100644 --- a/service/src/bridge/router.ts +++ b/service/src/bridge/router.ts @@ -310,6 +310,114 @@ export function createBridgeRouter(options: BridgeRouterOptions): Router { } })); + const untrustedRecoveryLimit = (operation: 'challenge' | 'recover'): RequestHandler => + (req, res, next) => { + if (options.authMode !== 'paired' || !options.pairings.recoveryEnabled) { + next(); + return; + } + // req.ip trusts X-Forwarded-For in our server. Use the connection peer so + // untrusted headers and arbitrary worker IDs cannot create new buckets. + void options.pairings.limitUntrustedRecovery(req.socket.remoteAddress ?? '', operation) + .then(() => next(), (error: unknown) => { + if (error instanceof BridgePairingError && error.code === 'RECOVERY_RATE_LIMITED') { + res.set('Retry-After', '60').status(429).json({ error: error.message, code: error.code }); + return; + } + next(error); + }); + }; + + router.post('/workers/:workerId/credentials/challenge', untrustedRecoveryLimit('challenge'), asyncRoute(async (req, res) => { + if (options.authMode !== 'paired' || !options.pairings.recoveryEnabled) { + res.status(404).json({ error: 'Machine recovery is disabled' }); + return; + } + const workerId = req.params.workerId; + const body = isRecord(req.body) ? req.body : {}; + if ( + !validWorkerId(workerId) || !configuredWorker(workerId) || + body.protocolVersion !== BRIDGE_PROTOCOL_VERSION || + body.operation !== 'credential.challenge' || + typeof body.serverId !== 'string' || body.serverId.length > 256 || + body.workerId !== workerId || + typeof body.timestamp !== 'string' || body.timestamp.length > 64 || + typeof body.nonce !== 'string' || !/^[A-Za-z0-9_-]{43}$/.test(body.nonce) || + typeof body.signature !== 'string' || !/^[A-Za-z0-9_-]{86}$/.test(body.signature) + ) { + res.status(400).json({ error: 'Invalid machine recovery challenge request' }); + return; + } + try { + const challenge = await options.pairings.createRecoveryChallenge( + workerId, + { + operation: 'credential.challenge', + serverId: body.serverId, + workerId, + timestamp: body.timestamp, + nonce: body.nonce, + }, + body.signature, + ); + res.json({ protocolVersion: BRIDGE_PROTOCOL_VERSION, ...challenge }); + } catch (error) { + if (error instanceof BridgePairingError) { + res.status(error.code === 'RECOVERY_RATE_LIMITED' ? 429 : 401) + .json({ error: error.message, code: error.code }); + return; + } + throw error; + } + })); + + router.post('/workers/:workerId/credentials/recover', untrustedRecoveryLimit('recover'), asyncRoute(async (req, res) => { + if (options.authMode !== 'paired' || !options.pairings.recoveryEnabled) { + res.status(404).json({ error: 'Machine recovery is disabled' }); + return; + } + const workerId = req.params.workerId; + const body = isRecord(req.body) ? req.body : {}; + if ( + !validWorkerId(workerId) || !configuredWorker(workerId) || + body.protocolVersion !== BRIDGE_PROTOCOL_VERSION || + body.operation !== 'credential.recover' || + typeof body.serverId !== 'string' || body.serverId.length > 256 || + typeof body.enrollmentGeneration !== 'string' || + !/^[A-Za-z0-9_-]{24}$/.test(body.enrollmentGeneration) || + typeof body.challenge !== 'string' || + !/^[A-Za-z0-9_-]{43}$/.test(body.challenge) || + typeof body.expiresAt !== 'string' || body.expiresAt.length > 64 || + typeof body.signature !== 'string' || + !/^[A-Za-z0-9_-]{86}$/.test(body.signature) + ) { + res.status(400).json({ error: 'Invalid machine recovery proof' }); + return; + } + try { + const credential = await options.pairings.recoverCredential( + workerId, + { + operation: 'credential.recover', + serverId: body.serverId, + workerId, + enrollmentGeneration: body.enrollmentGeneration, + challenge: body.challenge, + expiresAt: body.expiresAt, + }, + body.signature, + ); + res.json({ protocolVersion: BRIDGE_PROTOCOL_VERSION, ...credential }); + } catch (error) { + if (error instanceof BridgePairingError) { + res.status(error.code === 'RECOVERY_RATE_LIMITED' ? 429 : 401) + .json({ error: error.message, code: error.code }); + return; + } + throw error; + } + })); + router.post( '/workers/:workerId/revoke', adminAuth, @@ -768,6 +876,5 @@ router.post( }), ); - return router; } diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index 3d9d7e1c..0ca75566 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -65,6 +65,26 @@ export class BridgeStoreError extends Error { } } +function classifyPreEnqueueExpiry( + error: unknown, + args: { workspaceRequest?: WorkspaceToolRequest; workspaceId?: string; signal: AbortSignal }, +): unknown { + // A queue deadline reached before the assignment enqueue attempt is a + // definite non-execution, including expiry during the initial Redis reads. + if ( + (args.workspaceRequest != null || args.workspaceId != null) && + !args.signal.aborted && + error instanceof BridgeStoreError && + error.code === 'ASSIGNMENT_EXPIRED' + ) { + return new BridgeStoreError( + 'WORKSPACE_QUEUE_TIMEOUT', + 'Workspace capacity was unavailable before the queue deadline. The operation was not started. Wait for active work to finish or select an independent workspace on a machine with available capacity.', + ); + } + return error; +} + interface StoredAssignment extends CodeBridgeAssignment { leaseTokenHash: string; workerIdentityId?: string; @@ -778,12 +798,17 @@ export class RedisBridgeStore { 'Invalid workspace execution budget', ); } - this.assertDispatchActive(args.signal, args.deadlineAtMs); - const dispatchable = await this.dispatchCommand( - () => this.dispatchableRegistration(args.workerId), - args, - 'Bridge worker registration read', - ); + let dispatchable: Awaited>; + try { + this.assertDispatchActive(args.signal, args.deadlineAtMs); + dispatchable = await this.dispatchCommand( + () => this.dispatchableRegistration(args.workerId), + args, + 'Bridge worker registration read', + ); + } catch (error) { + throw classifyPreEnqueueExpiry(error, args); + } if (dispatchable == null) { throw new BridgeStoreError( 'WORKER_OFFLINE', @@ -844,17 +869,20 @@ export class RedisBridgeStore { `Bridge worker ${args.workerId} does not advertise programmatic execution for the selected workspace`, ); } - if ( - args.runtimeSessionId !== undefined && - (await this.dispatchCommand( - () => - this.redis.exists( - workspaceQuarantineKey(args.workerId, args.runtimeSessionId ?? ''), - ), - args, - 'Bridge workspace fence read', - )) === 1 - ) { + let quarantined = false; + const runtimeSessionId = args.runtimeSessionId; + if (runtimeSessionId !== undefined) { + try { + quarantined = (await this.dispatchCommand( + () => this.redis.exists(workspaceQuarantineKey(args.workerId, runtimeSessionId)), + args, + 'Bridge workspace fence read', + )) === 1; + } catch (error) { + throw classifyPreEnqueueExpiry(error, args); + } + } + if (quarantined) { throw new BridgeStoreError( 'WORKSPACE_QUARANTINED', 'Bridge workspace is quarantined after an incomplete result commit', @@ -1181,15 +1209,7 @@ export class RedisBridgeStore { } } catch (error) { // Once enqueue starts, even a lost Redis response may hide execution. - if ( - admission != null && !enqueueAttempted && !args.signal.aborted && - error instanceof BridgeStoreError && error.code === 'ASSIGNMENT_EXPIRED' - ) { - throw new BridgeStoreError( - 'WORKSPACE_QUEUE_TIMEOUT', - 'Workspace capacity was unavailable before the queue deadline. The operation was not started. Wait for active work to finish or select an independent workspace on a machine with available capacity.', - ); - } + if (!enqueueAttempted) throw classifyPreEnqueueExpiry(error, args); throw error; } finally { if (admission != null) { diff --git a/service/src/bridge/worker-admission.test.ts b/service/src/bridge/worker-admission.test.ts index 98363c34..c77ac19e 100644 --- a/service/src/bridge/worker-admission.test.ts +++ b/service/src/bridge/worker-admission.test.ts @@ -124,6 +124,32 @@ test('an expired queued call never reaches the worker and does not strand later await third; }); +test('an already-expired workspace deadline is a definite queue timeout before registration', async () => { + await register(); + await expect(dispatch('expired-before-read', new AbortController(), -1)).rejects.toMatchObject({ + code: 'WORKSPACE_QUEUE_TIMEOUT', + }); + expect(await redis.zcard(`codeapi:bridge:v1:worker:${workerId}:admission`)).toBe(0); + expect(await store.lease(workerId, incarnationId, 20)).toBeUndefined(); +}); + +test('expiry during the registration read never becomes an ambiguous assignment error', async () => { + await register(); + const registrationRead = spyOn(redis, 'mget').mockImplementation(async () => { + await new Promise(resolve => setTimeout(resolve, 25)); + throw new Error('registration read outlived its queue budget'); + }); + try { + await expect(dispatch('expired-during-read', new AbortController(), 1)).rejects.toMatchObject({ + code: 'WORKSPACE_QUEUE_TIMEOUT', + }); + } finally { + registrationRead.mockRestore(); + } + expect(await redis.zcard(`codeapi:bridge:v1:worker:${workerId}:admission`)).toBe(0); + expect(await store.lease(workerId, incarnationId, 20)).toBeUndefined(); +}); + test('a queued request is rejected if the worker withdraws its capability', async () => { await register(); const first = dispatch('first'); diff --git a/service/src/config.ts b/service/src/config.ts index 55570dfe..29ae854d 100644 --- a/service/src/config.ts +++ b/service/src/config.ts @@ -426,6 +426,22 @@ export const env = { BRIDGE_AUTH_MODE: bridgeAuthMode, /** Enrollment and lease credential shared only with the configured worker. */ BRIDGE_TOKEN: process.env.CODEAPI_BRIDGE_TOKEN ?? '', + /** Opt-in stable HTTPS origin of this Code API deployment (shared by all replicas). */ + BRIDGE_RECOVERY_SERVER_ID: process.env.CODEAPI_BRIDGE_RECOVERY_SERVER_ID ?? '', + /** Zero preserves machine authorization until explicit revocation. */ + BRIDGE_ENROLLMENT_TTL_SECONDS: Number(process.env.CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS ?? 0), + BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS: Number( + process.env.CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS ?? 60, + ), + BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE: Number( + process.env.CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE ?? 12, + ), + BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE: Number( + process.env.CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE ?? 30, + ), + BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE: Number( + process.env.CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE ?? 240, + ), /** * Runtime session affinity for stateful sandbox backends. * - `stateless` (default): no runtime sessions; `runtime_session_hint` ignored. diff --git a/service/src/workspace-tools/router.test.ts b/service/src/workspace-tools/router.test.ts index 18d24a76..d3cb032e 100644 --- a/service/src/workspace-tools/router.test.ts +++ b/service/src/workspace-tools/router.test.ts @@ -67,16 +67,17 @@ test('binds instance admission to the authenticated tenant and user while preser expect(dispatched[1]?.workspaceInstanceId).toBeUndefined(); }); -test.each<[WorkspaceToolRequest, number, number?, number?]>([ - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 90_000 }, 95_000, 125_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 300_000 }, 305_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 6000, 1000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, 600_000, undefined], - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, 5000], -])('separates the admission deadline from execution budget for %j', async (request, expectedExecution, ceiling, queueTimeoutMs) => { +test.each<[WorkspaceToolRequest, number, number?, number?, number?]>([ + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 90_000 }, 95_000, 125_000, undefined, 90_000], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 300_000 }, 305_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 6000, 1000, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, 600_000, undefined, undefined], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, 5000, 90_000], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined, 300_000], +])('separates the admission deadline from execution budget for %j', async (request, expectedExecution, ceiling, queueTimeoutMs, advertisedQueueWaitMs) => { const app = express(); app.use(json()); app.use((req, _res, next) => { @@ -102,12 +103,15 @@ test.each<[WorkspaceToolRequest, number, number?, number?]>([ const address = server.address(); if (address == null || typeof address === 'string') throw new Error('Missing listener'); const response = await fetch(`http://127.0.0.1:${address.port}/workspace-tools/execute`, { - method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(request), + method: 'POST', headers: { + 'Content-Type': 'application/json', + ...(advertisedQueueWaitMs === undefined ? {} : { 'X-LibreChat-Workspace-Queue-Wait-Ms': String(advertisedQueueWaitMs) }), + }, body: JSON.stringify(request), }); await response.json(); expect(executionBudget).toBe(expectedExecution); if (request.operation === 'execute_command') expect(commandTimeout).toBe(expectedExecution - 5000); - const expectedQueueBudget = queueTimeoutMs ?? Math.min(ceiling ?? 30_000, 300_000); + const expectedQueueBudget = Math.min(advertisedQueueWaitMs ?? 30_000, queueTimeoutMs ?? 300_000); expect(queueRemaining).toBeGreaterThan(expectedQueueBudget - 1000); expect(queueRemaining).toBeLessThanOrEqual(expectedQueueBudget); }); @@ -122,6 +126,32 @@ test.each([0, 300_001, Number.POSITIVE_INFINITY])('rejects an unbounded queue ov })).toThrow('Workspace queue timeout must be between 1 and 300000 milliseconds'); }); +test.each(['0', '-1', '300001', '1.5', '01', '1, 2', '999999999999999999999'])('rejects invalid per-request queue allowance %j before dispatch', async (queueWait) => { + const app = express(); + app.use(json()); + app.use((req, _res, next) => { + applyPrincipal(req, { userId: 'user-1', tenantId: 'tenant-1', principalSource: 'librechat_jwt', codeWorkerId: 'user-worker' }); + next(); + }); + let dispatched = false; + app.use(createWorkspaceToolsRouter({ + backend: 'remote-bridge', configuredWorkerId: 'user-worker', dynamicWorkers: false, + store: { async dispatchWorkspaceTool() { dispatched = true; throw new Error('must not dispatch'); } }, + })); + server = createServer(app); + await new Promise(resolve => server!.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (address == null || typeof address === 'string') throw new Error('Missing listener'); + const response = await fetch(`http://127.0.0.1:${address.port}/workspace-tools/execute`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-LibreChat-Workspace-Queue-Wait-Ms': queueWait }, + body: JSON.stringify({ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }), + }); + expect(response.status).toBe(400); + expect(await response.json()).toMatchObject({ code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + expect(dispatched).toBe(false); +}); + test('rejects new workspace dispatches while the service is shutting down', async () => { let dispatched = false; const app = express(); @@ -419,7 +449,7 @@ test.each([ status: expectedStatus, errorCode, outcome: 'completed', - deadlineBudgetMs: 330_000, + deadlineBudgetMs: 60_000, dispatchDurationMs: expect.any(Number), }), ); @@ -549,7 +579,8 @@ test('logs a disconnected dispatch once without inventing HTTP 200', async () => closeConnection(); await expect(response).rejects.toThrow(); await closed.promise; - expect(queueRemaining).toBeGreaterThan(120_000); + expect(queueRemaining).toBeGreaterThan(29_000); + expect(queueRemaining).toBeLessThanOrEqual(30_000); expect(dispatchAborted).toBe(true); expect(logSpy).not.toHaveBeenCalled(); settlementGate.resolve(); diff --git a/service/src/workspace-tools/router.ts b/service/src/workspace-tools/router.ts index 51e47f26..78dbd1e7 100644 --- a/service/src/workspace-tools/router.ts +++ b/service/src/workspace-tools/router.ts @@ -21,6 +21,8 @@ import { import { principalWorkspaceInstanceId } from '../bridge/workspace-instance'; const MAX_WORKSPACE_QUEUE_WAIT_MS = 5 * 60_000; +const DEFAULT_WORKSPACE_QUEUE_WAIT_MS = 30_000; +const WORKSPACE_QUEUE_WAIT_HEADER = 'X-LibreChat-Workspace-Queue-Wait-Ms'; interface WorkspaceToolsRouterOptions { store: Pick; @@ -57,12 +59,10 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) )) { throw new RangeError('Workspace execution timeout must be a positive safe integer'); } - // The HTTP disconnect cancels waiting; this bounds admission while the caller remains connected. - const queueBudgetMs = options.queueTimeoutMs ?? Math.min( - options.timeoutMs ?? 30_000, - MAX_WORKSPACE_QUEUE_WAIT_MS, - ); - if (!Number.isSafeInteger(queueBudgetMs) || queueBudgetMs < 1 || queueBudgetMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { + // A configured queue timeout is a ceiling. Legacy callers get 30 seconds even + // when their execution timeout is longer; only the per-request header opts in. + const queueCeilingMs = options.queueTimeoutMs ?? MAX_WORKSPACE_QUEUE_WAIT_MS; + if (!Number.isSafeInteger(queueCeilingMs) || queueCeilingMs < 1 || queueCeilingMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { throw new RangeError('Workspace queue timeout must be between 1 and 300000 milliseconds'); } const router = Router(); @@ -88,6 +88,21 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) }); return; } + const advertisedQueueWait = req.header(WORKSPACE_QUEUE_WAIT_HEADER); + if (advertisedQueueWait !== undefined && !/^[1-9]\d*$/.test(advertisedQueueWait)) { + outcome.errorCode = 'INVALID_WORKSPACE_QUEUE_WAIT'; + res.status(400).json({ error: 'Invalid workspace queue wait', code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + return; + } + const requestedQueueWaitMs = advertisedQueueWait === undefined + ? DEFAULT_WORKSPACE_QUEUE_WAIT_MS + : Number(advertisedQueueWait); + if (!Number.isSafeInteger(requestedQueueWaitMs) || requestedQueueWaitMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { + outcome.errorCode = 'INVALID_WORKSPACE_QUEUE_WAIT'; + res.status(400).json({ error: 'Invalid workspace queue wait', code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + return; + } + const queueBudgetMs = Math.min(requestedQueueWaitMs, queueCeilingMs); outcome.operation = req.body.operation; const principalRequest: WorkspaceToolRequest = req.body.workspaceInstanceId == null ? req.body diff --git a/tests/compose-bridge-config.cjs b/tests/compose-bridge-config.cjs index de14b553..8b3abd1d 100644 --- a/tests/compose-bridge-config.cjs +++ b/tests/compose-bridge-config.cjs @@ -15,6 +15,12 @@ function render(overrides) { CODEAPI_BRIDGE_WORKER_ID: '', CODEAPI_BRIDGE_TOKEN: '', CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS: '', + CODEAPI_BRIDGE_RECOVERY_SERVER_ID: '', + CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS: '', + CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS: '', + CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE: '', + CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE: '', + CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE: '', ...overrides, }, })); @@ -28,6 +34,12 @@ for (const overrides of [ CODEAPI_BRIDGE_DYNAMIC_WORKERS: 'false', CODEAPI_BRIDGE_WORKER_ID: 'test-worker', CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS: '4', + CODEAPI_BRIDGE_RECOVERY_SERVER_ID: 'https://code.example.test', + CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS: '86400', + CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS: '90', + CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE: '8', + CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE: '16', + CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE: '400', }, ]) { const config = render(overrides); @@ -39,10 +51,20 @@ for (const overrides of [ assert.equal(env.CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS, overrides.CODEAPI_BRIDGE_MAX_WORKSPACE_LEASE_SLOTS ?? '1'); assert.equal(env.CODEAPI_BRIDGE_DYNAMIC_WORKERS, overrides.CODEAPI_BRIDGE_DYNAMIC_WORKERS ?? 'true'); assert.equal(env.CODEAPI_BRIDGE_WORKER_ID, overrides.CODEAPI_BRIDGE_WORKER_ID ?? ''); + assert.equal(env.CODEAPI_BRIDGE_RECOVERY_SERVER_ID, overrides.CODEAPI_BRIDGE_RECOVERY_SERVER_ID ?? ''); + assert.equal(env.CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS, overrides.CODEAPI_BRIDGE_ENROLLMENT_TTL_SECONDS ?? '0'); + assert.equal(env.CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS, overrides.CODEAPI_BRIDGE_RECOVERY_CHALLENGE_TTL_SECONDS ?? '60'); + assert.equal(env.CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE, overrides.CODEAPI_BRIDGE_RECOVERY_MAX_CHALLENGES_PER_MINUTE ?? '12'); + assert.equal(env.CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE, overrides.CODEAPI_BRIDGE_RECOVERY_MAX_ATTEMPTS_PER_MINUTE ?? '30'); + assert.equal(env.CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE, overrides.CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE ?? '240'); } for (const name of ['egress_gateway', 'sandbox-runner']) { assert.equal(config.services[name].environment.CODEAPI_BRIDGE_TOKEN, undefined); + assert.equal(config.services[name].environment.CODEAPI_BRIDGE_RECOVERY_SERVER_ID, undefined); + assert.equal(config.services[name].environment.CODEAPI_BRIDGE_RECOVERY_MAX_UNTRUSTED_PER_MINUTE, undefined); } + assert.match(JSON.stringify(config.services.redis.command), /--appendonly.*yes/); + assert.ok(config.services.redis.volumes.some(volume => volume.target === '/data' && volume.type === 'volume')); } assert.equal(render({}).services.api.environment.CODEAPI_BRIDGE_TOKEN, ''); -console.log('Compose bridge configuration passed (dynamic/fixed pairing, no default secret).'); +console.log('Compose bridge configuration passed (pairing, opt-in recovery, durable Redis).');