From ef02127c37a5b146598037dd1b1609bf0201e711 Mon Sep 17 00:00:00 2001 From: Lia Date: Sat, 3 Oct 2026 03:48:04 +0000 Subject: [PATCH 1/5] =?UTF-8?q?=F0=9F=A7=BE=20feat:=20Persist=20Workspace?= =?UTF-8?q?=20Admission=20Across=20API=20Restarts?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .github/workflows/ci.yml | 2 +- docs/remote-bridge/README.md | 45 +++ service/openapi.yml | 105 +++++++ service/src/bridge/admission.ts | 38 +++ service/src/bridge/slots.ts | 12 +- service/src/bridge/store.ts | 184 +++++++++++- service/src/workspace-tools/index.ts | 38 ++- service/src/workspace-tools/outcome.ts | 4 +- .../workspace-tools/requests-router.test.ts | 84 ++++++ service/src/workspace-tools/requests.test.ts | 277 ++++++++++++++++++ service/src/workspace-tools/requests.ts | 244 +++++++++++++++ service/src/workspace-tools/router.ts | 47 ++- 12 files changed, 1064 insertions(+), 16 deletions(-) create mode 100644 service/src/workspace-tools/requests-router.test.ts create mode 100644 service/src/workspace-tools/requests.test.ts create mode 100644 service/src/workspace-tools/requests.ts diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 87aaa066..5ee14372 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -170,7 +170,7 @@ jobs: redis_socket="$RUNNER_TEMP/byom-admission.sock" redis-server --port 0 --unixsocket "$redis_socket" --save '' --appendonly no --daemonize yes trap 'redis-cli -s "$redis_socket" shutdown nosave' EXIT - BRIDGE_TEST_REDIS_URL="$redis_socket" bun test src/bridge/fleet.test.ts + BRIDGE_TEST_REDIS_URL="$redis_socket" bun test src/bridge/fleet.test.ts src/workspace-tools/requests.test.ts - name: Build service run: bun run build diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index b2d33db7..706e94e8 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -384,3 +384,48 @@ capabilities. LibreChat's owner-scoped environment registry can issue these principal-bound pairings without changing the worker execution protocol or moving code tools into the Agents SDK. + +### Durable workspace requests + +New clients can probe authenticated `GET /v1/workspace-tools/capabilities` for +`durableWorkspaceRequests: 1`. Older servers return 404 or advertise 0. Keep the +legacy synchronous `/workspace-tools/execute` path for those servers. Do not +fall back to synchronous execution after submitting durable work. + +Submit the same workspace-tool body to `POST /v1/workspace-tools/requests`, with +`X-LibreChat-Workspace-Request-Id` set to a unique 16–128 character identifier +(letters, digits, `_`, `-`). The queue-wait header and command timeout retain their +existing limits. HTTP 202 acknowledges durable acceptance, not execution. +Disconnecting does not cancel accepted work. + +- Repeat an identical submission with the same ID to recover a lost response. + An ID reused for different work returns 409 `REQUEST_CONFLICT`. +- Poll `GET /v1/workspace-tools/requests/:requestId`. States are `queued`, + `admitted`, `completed`, `failed`, and `cancelled`. Lookup is tenant/user scoped; + unknown or another principal's IDs return 404. Results are non-consuming reads. +- Cancel with `DELETE /v1/workspace-tools/requests/:requestId`. Cancellation is a + request, not proof of termination. A committed result can win the race. + `ASSIGNMENT_EXPIRED` means execution may have occurred. Never replay it. +- Status includes worker/workspace/lane metadata, queue position when available, + measured queue wait, and the execution deadline after admission. Queue position + is advisory: negotiated independent roots may run concurrently. + +Requests and results are retained for 24 hours from acceptance. Idempotency is +bounded by that retention window. A client must not resubmit an expired ID after +404; its outcome is unknown. Automatic result claims and wake-ups remain the +client's responsibility. + +API replicas reconcile accepted work from Redis. FIFO admission position survives +API restarts. Assignment enqueue and the durable `admitted` transition are atomic. +Recovery resumes only queued work and observes admitted assignments without +creating another command. Worker identity/incarnation changes fail pending work +rather than transferring it to another machine. Execution gets its full budget +after admission. Existing lease fencing, cancellation and workspace quarantine +continue to apply. + +This capability covers workspace tools, including `execute_command`, not +programmatic `/execute` jobs. It requires Redis state to survive the API restart; +Redis data loss is not recoverable. All API replicas behind an endpoint must +support the capability before a client enables it. Roll back clients to the legacy +path only for new work; drain durable requests before rolling back servers. +Worker protocol, worker concurrency and deployed rate limits are unchanged. diff --git a/service/openapi.yml b/service/openapi.yml index a257fd18..07d2d7e3 100644 --- a/service/openapi.yml +++ b/service/openapi.yml @@ -64,6 +64,40 @@ components: - $ref: '#/components/schemas/ExecutionProfileError' schemas: + WorkspaceToolRequest: + type: object + description: Workspace tool protocol v1 body; operation-specific fields follow the worker protocol. + required: [protocolVersion, operation, workspaceId] + properties: + protocolVersion: { type: integer, enum: [1] } + operation: + type: string + enum: [list_files, search_text, read_file, write_file, preview_edit, edit_file, execute_command] + workspaceId: { type: string } + additionalProperties: true + + WorkspaceRequestStatus: + type: object + required: [requestId, state, workerId, workspaceId, queueWaitMs, cancelRequested] + properties: + requestId: { type: string } + state: { type: string, enum: [queued, admitted, completed, failed, cancelled] } + workerId: { type: string } + workspaceId: { type: string } + workspaceInstanceId: { type: string } + worktree: { type: string } + queuePosition: { type: integer, minimum: 1 } + queueWaitMs: { type: integer, minimum: 0 } + executionDeadlineAt: { type: string, format: date-time } + cancelRequested: { type: boolean } + result: { type: object, additionalProperties: true } + error: + type: object + required: [code, message] + properties: + code: { type: string } + message: { type: string } + FileRef: type: object properties: @@ -312,6 +346,77 @@ components: type: string paths: + /workspace-tools/capabilities: + get: + summary: Discover opt-in durable workspace request support + responses: + '200': + description: All replicas must support version 1 before a client enables durable submission. + content: + application/json: + schema: + type: object + properties: + durableWorkspaceRequests: { type: integer, enum: [0, 1] } + + /workspace-tools/requests: + post: + summary: Submit an idempotent durable workspace request + description: >- + Disconnect does not cancel accepted work. IDs and results are retained for 24 hours + from acceptance. Never resubmit after that retention window or fall back to synchronous + execution after durable acceptance. The execution budget begins after admission. + parameters: + - $ref: '#/components/parameters/ExpectedExecutionProfile' + - in: header + name: X-LibreChat-Workspace-Request-Id + required: true + schema: { type: string, pattern: '^[A-Za-z0-9_-]{16,128}$' } + - in: header + name: X-LibreChat-Workspace-Queue-Wait-Ms + schema: { type: integer, minimum: 1, maximum: 300000, default: 30000 } + requestBody: + required: true + content: + application/json: + schema: { $ref: '#/components/schemas/WorkspaceToolRequest' } + responses: + '202': + description: Accepted, or an identical request with this principal-scoped ID already exists. + content: + application/json: + schema: { $ref: '#/components/schemas/WorkspaceRequestStatus' } + '400': { $ref: '#/components/responses/BadRequest' } + '409': { $ref: '#/components/responses/Conflict' } + '429': { description: Admission capacity or execution rate limit reached } + '503': { description: Code environment unavailable } + + /workspace-tools/requests/{requestId}: + parameters: + - in: path + name: requestId + required: true + schema: { type: string, pattern: '^[A-Za-z0-9_-]{16,128}$' } + - $ref: '#/components/parameters/ExpectedExecutionProfile' + get: + summary: Read a retained workspace request and result without consuming it + responses: + '200': + description: Status scoped to the authenticated tenant and user. ASSIGNMENT_EXPIRED is not proof of non-execution. + content: + application/json: + schema: { $ref: '#/components/schemas/WorkspaceRequestStatus' } + '404': { description: Unknown, expired, or another principal's request } + delete: + summary: Request cancellation of the same logical call + responses: + '200': + description: Cancellation recorded, or terminal status returned. A committed result can win the race. + content: + application/json: + schema: { $ref: '#/components/schemas/WorkspaceRequestStatus' } + '404': { description: Unknown, expired, or another principal's request } + /hosted-apps: post: summary: Start or reassert a resident hosted app diff --git a/service/src/bridge/admission.ts b/service/src/bridge/admission.ts index 4002ed6d..db375507 100644 --- a/service/src/bridge/admission.ts +++ b/service/src/bridge/admission.ts @@ -57,6 +57,44 @@ export class BridgeAdmissionQueue { ); } + async submit(args: { + workerId: string; id: string; deadlineAtMs: number; workspaceId?: string; + key: string; activeKey: string; fingerprint: string; record: string; retentionMs: number; + }): Promise<'accepted' | 'existing' | 'conflict' | 'full'> { + const result = Number(await this.redis.eval([ + 'local fingerprint = redis.call(\'HGET\', KEYS[5], \'fingerprint\')', + 'if fingerprint then', + ' if fingerprint == ARGV[6] then return 2 end', + ' return -1', + 'end', + 'local expired = redis.call(\'ZRANGEBYSCORE\', KEYS[2], \'-inf\', ARGV[2])', + 'for _, id in ipairs(expired) do', + ' redis.call(\'ZREM\', KEYS[1], id); redis.call(\'ZREM\', KEYS[2], id); redis.call(\'HDEL\', KEYS[4], id)', + 'end', + 'if redis.call(\'ZCARD\', KEYS[1]) >= tonumber(ARGV[4]) then return 0 end', + 'local sequence = redis.call(\'INCR\', KEYS[3])', + 'redis.call(\'ZADD\', KEYS[1], sequence, ARGV[1])', + 'redis.call(\'ZADD\', KEYS[2], ARGV[3], ARGV[1])', + 'if ARGV[5] ~= \'\' then redis.call(\'HSET\', KEYS[4], ARGV[1], ARGV[5]) end', + 'local latest = redis.call(\'ZREVRANGE\', KEYS[2], 0, 0, \'WITHSCORES\')', + 'for i = 1, 4 do redis.call(\'PEXPIREAT\', KEYS[i], tonumber(latest[2]) + 30000) end', + 'redis.call(\'HSET\', KEYS[5], \'record\', ARGV[7], \'fingerprint\', ARGV[6], \'state\', \'queued\', \'queueDeadlineAtMs\', ARGV[3])', + 'redis.call(\'PEXPIRE\', KEYS[5], ARGV[8])', + 'redis.call(\'ZADD\', KEYS[6], ARGV[2], KEYS[5])', + 'return 1', + ].join('\n'), 6, ...this.keys(args.workerId), args.key, args.activeKey, + args.id, Date.now(), args.deadlineAtMs, this.capacity, args.workspaceId ?? '', + args.fingerprint, args.record, args.retentionMs)); + if (result === 2) return 'existing'; + if (result === -1) return 'conflict'; + return result === 1 ? 'accepted' : 'full'; + } + + async position(workerId: string, id: string): Promise { + const rank = await this.redis.zrank(this.keys(workerId)[0], id); + return rank == null ? undefined : rank + 1; + } + async isHead(workerId: string, id: string): Promise { const [order, deadlines, , workspaces] = this.keys(workerId); return ( diff --git a/service/src/bridge/slots.ts b/service/src/bridge/slots.ts index 81f60ebe..e88c2b22 100644 --- a/service/src/bridge/slots.ts +++ b/service/src/bridge/slots.ts @@ -37,6 +37,7 @@ export class BridgeWorkspaceSlots { workspaceId: string; capacity: number; expiresAtMs: number; + refreshOwned?: boolean; }): Promise { if ( !Number.isSafeInteger(args.capacity) || @@ -83,7 +84,15 @@ export class BridgeWorkspaceSlots { " redis.call('HDEL', KEYS[1], 'a:' .. slot, 'i:' .. slot, 'w:' .. slot, 'e:' .. slot)", ' occupied = false', ' else', - ' if entry[1] == ARGV[2] then return slot end', + ' if entry[1] == ARGV[2] then', + ' if ARGV[7] == \'1\' then', + ' redis.call(\'HSET\', KEYS[1], \'e:\' .. slot, ARGV[5])', + ' for _, key in ipairs({KEYS[1], KEYS[2], KEYS[3]}) do', + ' if redis.call(\'PTTL\', key) < tonumber(ARGV[5]) - tonumber(ARGV[6]) then redis.call(\'PEXPIREAT\', key, ARGV[5]) end', + ' end', + ' end', + ' return slot', + ' end', ' busy[entry[3]] = true', ' local busyParent = parentOf(entry[3])', ' if busyParent then busyChildren[busyParent] = true end', @@ -133,6 +142,7 @@ export class BridgeWorkspaceSlots { args.capacity, args.expiresAtMs, Date.now(), + args.refreshOwned === true ? '1' : '0', ), ); if (result === -2) diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index f64e748a..a2ba6115 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -23,6 +23,7 @@ import { import type { BridgeWorkerBinding } from './pairing'; import { BridgeAdmissionQueue } from './admission'; import { BridgeWorkspaceSlots } from './slots'; +import type { StoredWorkspaceRequest } from '../workspace-tools/requests'; const PREFIX = 'codeapi:bridge:v1'; const POLL_INTERVAL_MS = 100; @@ -90,6 +91,22 @@ interface StoredAssignment extends CodeBridgeAssignment { leaseTokenHash: string; workerIdentityId?: string; workspaceFence?: string; + durableRequestKey?: string; +} + +export type DurableWorkspaceAssignment = StoredAssignment; + +export interface DurableWorkspaceGuard { + key: string; + claimKey: string; + token: string; +} + +export interface DurableWorkspaceOutcome { + state: 'completed' | 'failed' | 'cancelled'; + result?: WorkspaceToolResult; + error?: { code: string; message: string }; + settlement?: CodeBridgeWorkspaceSettlement; } type AssignmentOwnership = Pick< @@ -775,6 +792,125 @@ export class RedisBridgeStore { } } + async prepareDurableWorkspaceTool(args: { + workerId: string; tenantId: string; requireTenantBinding: boolean; + request: WorkspaceToolRequest; executionTimeoutMs: number; queueWaitMs: number; + }): Promise { + if (!isWorkspaceToolRequest(args.request) || + !Number.isSafeInteger(args.executionTimeoutMs) || args.executionTimeoutMs < 1 || args.executionTimeoutMs > 305_000 || + !Number.isSafeInteger(args.queueWaitMs) || args.queueWaitMs < 1 || args.queueWaitMs > 300_000) { + throw new BridgeStoreError('ASSIGNMENT_INVALID', 'Invalid durable workspace request'); + } + const current = await boundedCommand(this.dispatchableRegistration(args.workerId), this.redisCommandTimeoutMs, 'Durable worker registration'); + if (current == null) throw new BridgeStoreError('WORKER_OFFLINE', 'Code environment is offline'); + const registration = current.registration; + if ((args.requireTenantBinding && registration.binding == null) || + (registration.binding != null && registration.binding.tenantId !== args.tenantId)) { + throw new BridgeStoreError('WORKER_UNAUTHORIZED', 'Code environment is not authorized for this tenant'); + } + if ((registration.capabilities.workspaceLeaseSlots ?? 1) > this.maxWorkspaceLeaseSlots || + !supportsWorkspaceTool(registration, args.request)) { + throw new BridgeStoreError('WORKER_MISMATCH', 'Code environment does not support the requested workspace operation'); + } + return registration; + } + + async advanceDurableWorkspaceTool( + record: StoredWorkspaceRequest, + guard: DurableWorkspaceGuard, + ): Promise { + if (record.state === 'queued') { + if (record.cancelRequested === true) return { state: 'cancelled' }; + if (Date.now() >= record.queueDeadlineAtMs) return { + state: 'failed', error: { code: 'WORKSPACE_QUEUE_TIMEOUT', message: 'Workspace admission deadline expired. The operation was not started.' }, + }; + const registration = await this.prepareDurableWorkspaceTool({ + ...record, queueWaitMs: Math.max(1, record.queueDeadlineAtMs - Date.now()), + }); + if (registration.incarnationId !== record.registration.incarnationId || + registration.identityId !== record.registration.identityId || + registration.binding?.tenantId !== record.registration.binding?.tenantId) { + throw new BridgeStoreError('WORKER_FENCED', 'Code environment changed while the request was waiting'); + } + const admission = new BridgeAdmissionQueue(this.redis); + const workspace = workspaceAdmissionId(record.request.workspaceId, record.request.workspaceInstanceId, record.request.worktree); + const capacity = registration.capabilities.workspaceLeaseSlots ?? 1; + if (capacity === 1 && !await admission.isHead(record.workerId, record.id)) return; + const expiresAtMs = Date.now() + record.executionTimeoutMs; + const ttlSeconds = assignmentTtlSeconds(expiresAtMs); + let slot: number | undefined; + if (capacity > 1) { + slot = await new BridgeWorkspaceSlots(this.redis).reserve({ + workerId: record.workerId, incarnationId: registration.incarnationId, + assignmentId: record.id, workspaceId: workspace, capacity, + expiresAtMs: Date.now() + ttlSeconds * 1000, refreshOwned: true, + }); + if (slot === undefined) return; + } else if (!await this.acquireLock(record.workerId, record.id, registration.incarnationId, ttlSeconds, true)) return; + if (Date.now() >= record.queueDeadlineAtMs) return { + state: 'failed', error: { code: 'WORKSPACE_QUEUE_TIMEOUT', message: 'Workspace admission deadline expired. The operation was not started.' }, + }; + const leaseToken = randomBytes(32).toString('base64url'); + const assignment: StoredAssignment = { + protocolVersion: BRIDGE_PROTOCOL_VERSION, + assignmentId: record.id, workerId: record.workerId, + incarnationId: registration.incarnationId, + generation: await this.redis.incr(generationKey(record.workerId)), + leaseToken, leaseTokenHash: tokenHash(leaseToken), + workerIdentityId: registration.identityId, + workspaceFence: `${NATIVE_WORKSPACE_FENCE_PREFIX}${workspace}`, + workspaceLeaseSlot: slot, executionKind: 'workspace_tool', durableRequestKey: record.key, + expiresAt: new Date(Date.now() + record.executionTimeoutMs).toISOString(), + request: record.request, + }; + const ready = await this.dispatchableRegistration(record.workerId); + if (ready == null || ready.registration.incarnationId !== registration.incarnationId || + ready.registration.identityId !== registration.identityId || + ready.registration.binding?.tenantId !== registration.binding?.tenantId || + !supportsWorkspaceTool(ready.registration, record.request)) { + throw new BridgeStoreError('WORKER_FENCED', 'Code environment changed before admission'); + } + await this.enqueueForActiveIncarnation(assignment, assignmentTtlSeconds(Date.parse(assignment.expiresAt)), ready.readyToken, guard); + return; + } + const assignment = record.assignment; + if (assignment == null) throw new Error('Durable admitted request has no assignment'); + let settlement: CodeBridgeWorkspaceSettlement; + try { + const raw = await boundedCommand(this.redis.get(settlementKey(record.id)), this.redisCommandTimeoutMs, 'Durable settlement lookup'); + if (raw == null && record.cancelRequested !== true && Date.now() < Date.parse(assignment.expiresAt)) return; + if (raw != null) settlement = JSON.parse(raw) as CodeBridgeWorkspaceSettlement; + else { + const controller = new AbortController(); + if (record.cancelRequested === true) controller.abort(); + settlement = await this.waitForSettlement(assignment, Date.parse(assignment.expiresAt), controller.signal) as unknown as CodeBridgeWorkspaceSettlement; + } + } catch (error) { + if (!(error instanceof BridgeStoreError) || error.code !== 'ASSIGNMENT_EXPIRED') throw error; + return { state: 'failed', error: { code: error.code, message: 'Assignment ended without a confirmed result. Do not replay this operation.' } }; + } + if (settlement.status === 'fulfilled') { + if (!isWorkspaceToolResult(record.request, settlement.result, record.registration.capabilities.workspaceTools)) { + return { state: 'failed', error: { code: 'RESULT_INVALID', message: 'Code environment returned an invalid workspace result. Do not replay this operation.' } }; + } + return { state: 'completed', result: settlement.result, settlement }; + } + return { + state: record.cancelRequested === true ? 'cancelled' : 'failed', + error: { code: settlement.errorCode ?? 'WORKSPACE_TOOL_REJECTED', message: settlement.error }, settlement, + }; + } + + async finishDurableWorkspaceTool(record: StoredWorkspaceRequest): Promise { + if (record.assignment != null) { + if (record.settlement != null) await this.commitPendingWorkspace(record.assignment, record.settlement, record.key); + await this.cleanupDispatch(record.workerId, record.id, record.assignment); + } else { + await this.cleanupUnassignedSlot(record.workerId, record.registration.incarnationId, record.id); + } + await new BridgeAdmissionQueue(this.redis).leave(record.workerId, record.id); + } + async dispatchWorkspaceTool(args: { workerId: string; tenantId?: string; @@ -1414,6 +1550,7 @@ export class RedisBridgeStore { leaseTokenHash: _leaseTokenHash, workerIdentityId: _workerIdentityId, workspaceFence: _workspaceFence, + durableRequestKey: _durableRequestKey, ...wireAssignment } = assignment; return { @@ -1706,7 +1843,9 @@ export class RedisBridgeStore { 'Bridge assignment has expired', ); } - const ttlSeconds = assignmentTtlSeconds(Date.parse(assignment.expiresAt)); + const ttlSeconds = assignment.durableRequestKey == null + ? assignmentTtlSeconds(Date.parse(assignment.expiresAt)) + : 86400; const settlementKeys = [ assignmentKey(assignmentId), settlementKey(assignmentId), @@ -2237,17 +2376,25 @@ export class RedisBridgeStore { assignment: StoredAssignment, ttlSeconds: number, readyToken?: string, + guard?: DurableWorkspaceGuard, ): Promise { const script = [ + 'local keyCount = tonumber(ARGV[10])', + ...(guard == null ? [] : [ + 'if redis.call(\'GET\', KEYS[keyCount + 2]) ~= ARGV[11] then return -2 end', + 'if redis.call(\'HGET\', KEYS[keyCount + 1], \'state\') ~= \'queued\' then return -2 end', + 'if redis.call(\'HGET\', KEYS[keyCount + 1], \'cancelRequested\') == \'1\' then return -2 end', + 'if tonumber(redis.call(\'HGET\', KEYS[keyCount + 1], \'queueDeadlineAtMs\')) <= tonumber(ARGV[12]) then return -2 end', + ]), "if redis.call('GET', KEYS[1]) ~= ARGV[1] then return 0 end", 'if ARGV[7] ~= "" and redis.call(\'GET\', KEYS[6]) ~= ARGV[7] then return 0 end', - "if #KEYS >= 7 and redis.call('EXISTS', KEYS[7]) == 1 then return -1 end", + 'if keyCount >= 7 and redis.call(\'EXISTS\', KEYS[7]) == 1 then return -1 end', // A lane refuses a quarantined checkout. A checkout reaches enqueue only once no // lane beneath it holds a slot, so any lane fence still indexed here is stuck. "if ARGV[9] == 'lane' and redis.call('EXISTS', KEYS[10]) == 1 then return -1 end", "if ARGV[9] == 'checkout' then", - " for i = 11, #KEYS do if redis.call('EXISTS', KEYS[i]) == 1 then return -1 end end", - " for i = 11, #KEYS do redis.call('SREM', KEYS[10], KEYS[i]) end", + ' for i = 11, keyCount do if redis.call(\'EXISTS\', KEYS[i]) == 1 then return -1 end end', + ' for i = 11, keyCount do redis.call(\'SREM\', KEYS[10], KEYS[i]) end', 'end', 'redis.call(\'SET\', KEYS[2], ARGV[2], \"EX\", ARGV[3])', "redis.call('RPUSH', KEYS[3], ARGV[4])", @@ -2256,15 +2403,18 @@ export class RedisBridgeStore { ? ['redis.call(\'SET\', KEYS[4], ARGV[1], \"PX\", ARGV[5])'] : []), 'redis.call(\'SET\', KEYS[5], "1", \"PXAT\", ARGV[6])', - "if #KEYS >= 7 then redis.call('SET', KEYS[7], ARGV[4]) end", + 'if keyCount >= 7 then redis.call(\'SET\', KEYS[7], ARGV[4]) end', "if ARGV[9] == 'lane' then redis.call('SADD', KEYS[11], KEYS[7]) end", - 'if #KEYS >= 9 then', + 'if keyCount >= 9 then', " local epoch = redis.call('GET', KEYS[9])", " if type(epoch) ~= 'string' then epoch = '0'; redis.call('SET', KEYS[9], epoch, 'EX', ARGV[3]) end", " if redis.call('PTTL', KEYS[9]) < tonumber(ARGV[3]) * 1000 then redis.call('EXPIRE', KEYS[9], ARGV[3]) end", " redis.call('HSET', KEYS[8], 'metadata', ARGV[8], 'epoch', epoch)", - " redis.call('EXPIRE', KEYS[8], ARGV[3])", + ' redis.call(\'EXPIRE\', KEYS[8], ARGV[13])', 'end', + ...(guard == null ? [] : [ + 'redis.call(\'HSET\', KEYS[keyCount + 1], \'state\', \'admitted\', \'assignment\', ARGV[2], \'admittedAtMs\', ARGV[12])', + ]), 'return 1', ].join('\n'); const keys = [ @@ -2322,6 +2472,8 @@ export class RedisBridgeStore { fenceScope = 'checkout'; } } + const keyCount = keys.length; + if (guard != null) keys.push(guard.key, guard.claimKey); const result = await this.redis.eval( script, keys.length, @@ -2335,6 +2487,10 @@ export class RedisBridgeStore { readyToken ?? '', JSON.stringify(receipt), fenceScope, + keyCount, + guard?.token ?? '', + Date.now(), + guard == null ? ttlSeconds : 86400, ); if (Number(result) === -1) { throw new BridgeStoreError( @@ -2350,8 +2506,10 @@ export class RedisBridgeStore { assignmentId: string, incarnationId: string, ttlSeconds: number, + reuseOwned = false, ): Promise { const script = [ + 'if ARGV[4] == \'1\' and redis.call(\'GET\', KEYS[1]) == ARGV[1] and redis.call(\'GET\', KEYS[2]) == ARGV[2] then redis.call(\'PEXPIRE\', KEYS[1], ARGV[3]); redis.call(\'PEXPIRE\', KEYS[2], ARGV[3]); return 1 end', 'if redis.call(\'EXISTS\', KEYS[1]) == 1 then return 0 end', 'redis.call(\'SET\', KEYS[1], ARGV[1], \"PX\", ARGV[3])', 'redis.call(\'SET\', KEYS[2], ARGV[2], \"PX\", ARGV[3])', @@ -2365,6 +2523,7 @@ export class RedisBridgeStore { assignmentId, incarnationId, String(ttlSeconds * 1000), + reuseOwned ? '1' : '0', ); return Number(result) === 1; } @@ -2412,6 +2571,7 @@ export class RedisBridgeStore { private async commitPendingWorkspace( assignment: StoredAssignment, settlement: AnyCodeBridgeSettlement, + durableKey?: string, ): Promise { if (assignment.workspaceLeaseSlot !== undefined) { const committed = Number( @@ -2450,8 +2610,13 @@ export class RedisBridgeStore { } const runtimeSessionId = assignmentWorkspace(assignment)!; const script = [ + ...(durableKey == null ? [] : [ + 'if redis.call(\'EXISTS\', KEYS[2]) == 1 then return 1 end', + ]), "if redis.call('GET', KEYS[1]) == ARGV[1] then", - " return redis.call('DEL', KEYS[1])", + ' redis.call(\'DEL\', KEYS[1])', + ...(durableKey == null ? [] : [' redis.call(\'SET\', KEYS[2], \'1\', \'PX\', 86400000)']), + ' return 1', 'end', 'return 0', ].join('\n'); @@ -2459,8 +2624,9 @@ export class RedisBridgeStore { await boundedCommand( this.redis.eval( script, - 1, + durableKey == null ? 1 : 2, workspaceQuarantineKey(assignment.workerId, runtimeSessionId), + ...(durableKey == null ? [] : [`${durableKey}:committed:${assignment.assignmentId}`]), assignment.assignmentId, ), // Once settlement wins, caller cancellation must not prevent its diff --git a/service/src/workspace-tools/index.ts b/service/src/workspace-tools/index.ts index 6f105f16..53f750c1 100644 --- a/service/src/workspace-tools/index.ts +++ b/service/src/workspace-tools/index.ts @@ -4,12 +4,48 @@ import { bridgeStore } from '../bridge'; import { env } from '../config'; import { executionLimiter } from '../middleware/limits'; import { createWorkspaceToolsRouter } from './router'; +import { RedisWorkspaceRequests } from './requests'; +import { connection } from '../queue'; +import logger from '../logger'; +import { checkServiceStartUp, checkServiceShutDown } from '../lifecycle'; +import { withSpan } from '../telemetry'; const router = Router(); -router.use('/workspace-tools/execute', executionLimiter); +const requests = new RedisWorkspaceRequests(connection, bridgeStore, record => { + const attributes = { + 'codeapi.request.id': record.id, + 'codeapi.admission.state': record.state, + 'codeapi.worker.id': record.workerId, + 'codeapi.workspace.id': record.request.workspaceId, + 'codeapi.workspace.instance_id': record.request.workspaceInstanceId ?? '', + 'codeapi.workspace.worktree': record.request.worktree ?? '', + 'codeapi.admission.wait_ms': Math.max(0, (record.admittedAtMs ?? Math.min(record.finishedAtMs ?? Date.now(), record.queueDeadlineAtMs)) - record.createdAtMs), + 'codeapi.admission.error_code': record.error?.code ?? '', + }; + logger.info('Durable workspace admission transition', attributes); + void withSpan('codeapi.workspace.admission.transition', attributes, () => undefined).catch(error => { + logger.warn('Admission transition telemetry failed', { error }); + }); +}); +let reconciling = false; +let unavailable = false; +const reconciliation = setInterval(() => { + if (reconciling || checkServiceStartUp() || checkServiceShutDown()) return; + reconciling = true; + void requests.reconcile().then(() => { unavailable = false; }).catch(error => { + if (!unavailable) logger.warn('Durable workspace reconciliation unavailable', { error }); + unavailable = true; + }).finally(() => { reconciling = false; }); +}, 250); +reconciliation.unref(); +router.use(['/workspace-tools/execute', '/workspace-tools/requests'], (req, res, next) => { + if (req.method === 'POST') executionLimiter(req, res, next); + else next(); +}); router.use( createWorkspaceToolsRouter({ store: bridgeStore, + requests, backend: env.SANDBOX_BACKEND, configuredWorkerId: env.BRIDGE_WORKER_ID, dynamicWorkers: env.BRIDGE_DYNAMIC_WORKERS, diff --git a/service/src/workspace-tools/outcome.ts b/service/src/workspace-tools/outcome.ts index d965d515..e8dc100f 100644 --- a/service/src/workspace-tools/outcome.ts +++ b/service/src/workspace-tools/outcome.ts @@ -25,7 +25,7 @@ const earlyErrorCodes: Record = { 500: 'INTERNAL_ERROR', }; -export function getWorkspaceToolOutcome(res: Response): WorkspaceToolOutcome { +export function getWorkspaceToolOutcome(res: Response, route = '/workspace-tools/execute'): WorkspaceToolOutcome { const existing = outcomes.get(res); if (existing != null) return existing; const startedAt = performance.now(); @@ -39,7 +39,7 @@ export function getWorkspaceToolOutcome(res: Response): WorkspaceToolOutcome { outcomes.delete(res); const finished = res.writableFinished; logger.log(finished && res.statusCode < 400 ? 'info' : 'warn', 'Workspace tool request completed', { - route: '/workspace-tools/execute', + route, operation: outcome.operation, workerId: outcome.workerId, status: finished ? res.statusCode : undefined, diff --git a/service/src/workspace-tools/requests-router.test.ts b/service/src/workspace-tools/requests-router.test.ts new file mode 100644 index 00000000..27ea4710 --- /dev/null +++ b/service/src/workspace-tools/requests-router.test.ts @@ -0,0 +1,84 @@ +import { createServer } from 'node:http'; +import type { Server } from 'node:http'; +import { afterEach, expect, test } from 'bun:test'; +import express, { json } from 'express'; +import RedisMock from 'ioredis-mock'; +import type Redis from 'ioredis'; +import { applyPrincipal } from '../auth/principal'; +import { RedisBridgeStore } from '../bridge/store'; +import { createWorkspaceToolsRouter } from './router'; +import { RedisWorkspaceRequests } from './requests'; + +let server: Server | undefined; +afterEach(() => { server?.close(); server = undefined; }); +const body = { protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }; +const requestId = 'http-request-000001'; + +async function setup(durable = true): Promise { + const redis = new RedisMock() as unknown as Redis; + const bridge = new RedisBridgeStore(redis); + await bridge.register({ protocolVersion: 1, workerId: 'worker', incarnationId: 'incarnation-00000001', capabilities: { + sandboxProfile: 'native-srt', statefulWorkspace: false, runtimes: [], + workspaceTools: { protocolVersion: 1, operations: ['read_file'], workspaces: [{ id: 'primary' }] }, + } }); + const app = express(); + app.use(json()); + app.use((req, _res, next) => { + if (req.header('X-Test-User') !== 'unauthenticated') applyPrincipal(req, { + tenantId: 'tenant', userId: req.header('X-Test-User') ?? 'user', principalSource: 'local', codeWorkerId: 'worker', + }); + next(); + }); + app.use(createWorkspaceToolsRouter({ + store: bridge, requests: durable ? new RedisWorkspaceRequests(redis, bridge) : undefined, + backend: 'remote-bridge', configuredWorkerId: 'worker', dynamicWorkers: false, + })); + 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'); + return `http://127.0.0.1:${address.port}`; +} + +function post(url: string, id = requestId, value = body): Promise { + return fetch(`${url}/workspace-tools/requests`, { method: 'POST', headers: { + 'Content-Type': 'application/json', 'X-LibreChat-Workspace-Request-Id': id, + }, body: JSON.stringify(value) }); +} + +test.each([false, true])('capability discovery explicitly distinguishes durable support (%s)', async durable => { + const url = await setup(durable); + expect(await (await fetch(`${url}/workspace-tools/capabilities`)).json()).toEqual({ durableWorkspaceRequests: durable ? 1 : 0 }); +}); + +test('submission, duplicate lookup and cancellation use the same principal-scoped handle without exposing leases', async () => { + const url = await setup(); + const accepted = await post(url); + expect(accepted.status).toBe(202); + expect(await accepted.json()).toMatchObject({ requestId, state: 'queued', queuePosition: 1 }); + const duplicate = await post(url); + expect(duplicate.status).toBe(202); + const status = await (await fetch(`${url}/workspace-tools/requests/${requestId}`)).json(); + expect(status).toMatchObject({ requestId, state: 'queued' }); + expect(status).not.toHaveProperty('assignment'); + expect(status).not.toHaveProperty('leaseToken'); + expect((await fetch(`${url}/workspace-tools/requests/${requestId}`, { headers: { 'X-Test-User': 'another' } })).status).toBe(404); + const cancelled = await fetch(`${url}/workspace-tools/requests/${requestId}`, { method: 'DELETE' }); + expect(await cancelled.json()).toMatchObject({ requestId, cancelRequested: true }); +}); + +test('missing IDs and conflicting requests fail without dispatching', async () => { + const url = await setup(); + expect((await post(url, '')).status).toBe(400); + expect((await post(url)).status).toBe(202); + const conflict = await post(url, requestId, { ...body, path: 'different' }); + expect(conflict.status).toBe(409); + expect(await conflict.json()).toMatchObject({ code: 'REQUEST_CONFLICT' }); +}); + +test('unauthenticated request status and capability discovery are rejected', async () => { + const url = await setup(); + for (const path of ['capabilities', `requests/${requestId}`]) { + expect((await fetch(`${url}/workspace-tools/${path}`, { headers: { 'X-Test-User': 'unauthenticated' } })).status).toBe(401); + } +}); diff --git a/service/src/workspace-tools/requests.test.ts b/service/src/workspace-tools/requests.test.ts new file mode 100644 index 00000000..992060cc --- /dev/null +++ b/service/src/workspace-tools/requests.test.ts @@ -0,0 +1,277 @@ +import { afterEach, expect, spyOn, test } from 'bun:test'; +import { createHash } from 'node:crypto'; +import RedisMock from 'ioredis-mock'; +import Redis from 'ioredis'; +import { RedisBridgeStore } from '../bridge/store'; +import type { CodeBridgeAssignment } from '../bridge/store'; +import { RedisWorkspaceRequests } from './requests'; + +const testRedisUrl = process.env.BRIDGE_TEST_REDIS_URL; +const redis = testRedisUrl !== undefined && testRedisUrl.length > 0 + ? new Redis(testRedisUrl) + : new RedisMock() as unknown as Redis; +const bridge = new RedisBridgeStore(redis, 600, 1000, 2); +const owner = { tenantId: 'tenant', userId: 'user' }; +const workerId = 'durable-worker'; +const incarnationId = 'incarnation-00000001'; +const requestId = 'request-00000000001'; +const activeKey = 'codeapi:workspace-requests:v1:active'; +const queueKey = `codeapi:bridge:v1:worker:${workerId}:admission`; + +// Only a disposable test Redis may be supplied. +afterEach(async () => { await redis.flushall(); }); + +async function setup(slots = 1): Promise { + const generation = await bridge.register({ protocolVersion: 1, workerId, incarnationId, + capabilities: { + sandboxProfile: 'native-srt', statefulWorkspace: false, runtimes: ['bash'], + ...(slots > 1 ? { workspaceLeaseSlots: slots, requiresReadyConfirmation: true } : {}), + workspaceTools: { protocolVersion: 1, operations: ['read_file', 'execute_command'], + workspaces: [{ id: 'primary' }, { id: 'independent' }] }, + }, + }); + if (slots > 1) await bridge.confirmReady(workerId, incarnationId, generation); + return new RedisWorkspaceRequests(redis, bridge); +} + +function submit(requests: RedisWorkspaceRequests, id = requestId, workspaceId = 'primary'): ReturnType { + return requests.submit({ owner, requestId: id, workerId, requireTenantBinding: false, + request: { protocolVersion: 1, operation: 'read_file', workspaceId, path: 'README.md' }, + queueWaitMs: 300_000, executionTimeoutMs: 30_000, + }); +} + +async function tick(requests: RedisWorkspaceRequests): Promise { + const keys = await redis.zrange(activeKey, 0, -1); + for (const key of keys) await redis.zadd(activeKey, 0, key); + await requests.reconcile(); +} + +async function settle(assignment: CodeBridgeAssignment, fulfilled = true): Promise { + await bridge.acknowledgeLease(workerId, incarnationId, assignment.assignmentId, + assignment.generation, assignment.leaseToken); + await bridge.settle(workerId, assignment.assignmentId, { + protocolVersion: 1, incarnationId, generation: assignment.generation, + leaseToken: assignment.leaseToken, + ...(fulfilled ? { status: 'fulfilled' as const, result: { protocolVersion: 1, operation: 'read_file', workspaceId: (assignment.request as { workspaceId: string }).workspaceId, path: 'README.md', content: 'hello', startLine: 1, endLine: 1, truncated: false } } + : { status: 'rejected' as const, error: 'stopped' }), + }); +} + +async function mutateRecord(update: (record: Record) => void): Promise { + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const record = JSON.parse((await redis.hget(key, 'record'))!); + update(record); + await redis.hset(key, 'record', JSON.stringify(record)); +} + +test('submission is idempotent across replicas, scoped by principal, and conflicts do not move FIFO position', async () => { + const requests = await setup(); + await submit(requests); + await submit(requests, 'request-00000000002'); + const replica = new RedisWorkspaceRequests(redis, bridge); + expect((await submit(replica)).queuePosition).toBe(1); + expect(await redis.zcard(queueKey)).toBe(2); + await expect(submit(replica, requestId, 'independent')).rejects.toThrow('different work'); + expect(await replica.get({ ...owner, userId: 'someone-else' }, requestId)).toBeUndefined(); + expect(await replica.cancel({ ...owner, tenantId: 'someone-else' }, requestId)).toBeUndefined(); +}); + +test('restart before dispatch preserves the queue and starts a full execution budget after sixty seconds waiting', async () => { + const requests = await setup(); + await submit(requests); + const clock = spyOn(Date, 'now').mockReturnValue(Date.now() + 60_000); + const restarted = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis)); + let assignment: CodeBridgeAssignment | undefined; + try { + await tick(restarted); + assignment = await bridge.lease(workerId, incarnationId, 0); + expect(assignment).toBeDefined(); + expect(assignment!.remainingMs).toBeGreaterThan(29_000); + expect((await restarted.get(owner, requestId))?.queueWaitMs).toBeGreaterThanOrEqual(60_000); + } finally { clock.mockRestore(); } + await settle(assignment!); + await tick(restarted); + expect(await restarted.get(owner, requestId)).toMatchObject({ state: 'completed', result: { content: 'hello' } }); +}); + +test.each([false, true])('restart around admission acknowledgement never enqueues another assignment (ack=%s)', async acknowledged => { + const requests = await setup(); + await submit(requests); + await tick(requests); + const assignment = await bridge.lease(workerId, incarnationId, 0); + if (acknowledged) await bridge.acknowledgeLease(workerId, incarnationId, assignment!.assignmentId, + assignment!.generation, assignment!.leaseToken); + const restarted = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis)); + await tick(restarted); + await submit(restarted); + expect((await bridge.lease(workerId, incarnationId, 0))?.assignmentId).toBe(assignment!.assignmentId); + expect(await redis.llen(`codeapi:bridge:v1:worker:${workerId}:incarnation:${incarnationId}:assignments`)).toBe(0); + await settle(assignment!); + await tick(restarted); + expect((await restarted.get(owner, requestId))?.state).toBe('completed'); +}); + +test('a lost enqueue acknowledgement is reconciled rather than replayed', async () => { + const requests = await setup(); + await submit(requests); + const original = redis.eval.bind(redis); + const failure = spyOn(redis, 'eval').mockImplementation((async (...args: Parameters) => { + const result = await original(...args); + if (String(args[0]).includes('\'state\', \'admitted\', \'assignment\'')) throw new Error('lost acknowledgement'); + return result; + }) as Redis['eval']); + try { await expect(tick(requests)).rejects.toThrow('lost acknowledgement'); } finally { failure.mockRestore(); } + const assignment = await bridge.lease(workerId, incarnationId, 0); + expect(assignment).toBeDefined(); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect((await requests.get(owner, requestId))?.state).toBe('admitted'); + await settle(assignment!); + await tick(requests); + expect((await requests.get(owner, requestId))?.state).toBe('completed'); +}); + +test('cancel before admission is durable and does not dispatch', async () => { + const requests = await setup(); + await submit(requests); + await requests.cancel(owner, requestId); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect((await requests.get(owner, requestId))?.state).toBe('cancelled'); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); + expect(await redis.zcard(queueKey)).toBe(0); +}); + +test('cancel racing a persisted fulfillment returns the winning result without consuming it', async () => { + const requests = await setup(); + await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await settle(assignment); + await requests.cancel(owner, requestId); + await tick(new RedisWorkspaceRequests(redis, bridge)); + const statuses = await Promise.all([requests.get(owner, requestId), requests.get(owner, requestId), requests.cancel(owner, requestId)]); + for (const status of statuses) expect(status).toMatchObject({ state: 'completed', result: { content: 'hello' } }); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); +}); + +test('restart after result persistence repeats only idempotent cleanup', async () => { + const requests = await setup(); + await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await settle(assignment); + const cleanup = spyOn(bridge, 'finishDurableWorkspaceTool').mockRejectedValue(new Error('process stopped')); + try { await expect(tick(requests)).rejects.toThrow('process stopped'); } finally { cleanup.mockRestore(); } + expect((await requests.get(owner, requestId))?.state).toBe('completed'); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect(await redis.zcard(activeKey)).toBe(0); + await submit(requests); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); +}); + +test('pending work rejects a replacement worker incarnation and never transfers to it', async () => { + const requests = await setup(); + await submit(requests); + await bridge.register({ protocolVersion: 1, workerId, incarnationId: 'replacement-00001', + capabilities: { sandboxProfile: 'native-srt', statefulWorkspace: false, runtimes: [], + workspaceTools: { protocolVersion: 1, operations: ['read_file'], workspaces: [{ id: 'primary' }] } }, + }); + await tick(requests); + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'failed', error: { code: 'WORKER_FENCED' } }); +}); + +test('durable requests preserve negotiated independent-root concurrency', async () => { + const requests = await setup(2); + await submit(requests); await submit(requests, 'request-00000000002', 'independent'); + await tick(requests); await tick(requests); + const first = await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0); + const second = await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 1); + expect(first).toBeDefined(); expect(second).toBeDefined(); + await settle(first!); await settle(second!); await tick(requests); + expect((await requests.get(owner, requestId))?.state).toBe('completed'); + expect((await requests.get(owner, 'request-00000000002'))?.state).toBe('completed'); +}); + +test('restart after capacity reservation but before enqueue refreshes its own reservation without duplicating work', async () => { + const requests = await setup(); + await submit(requests); + const generation = spyOn(redis, 'incr').mockRejectedValue(new Error('process stopped before enqueue')); + try { await expect(tick(requests)).rejects.toThrow('process stopped before enqueue'); } finally { generation.mockRestore(); } + expect((await requests.get(owner, requestId))?.state).toBe('queued'); + await redis.pexpire(`codeapi:bridge:v1:worker:${workerId}:lock`, 5000); + await tick(new RedisWorkspaceRequests(redis, bridge)); + const assignment = await bridge.lease(workerId, incarnationId, 0); + expect(assignment).toBeDefined(); + expect(await redis.pttl(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeGreaterThan(30_000); +}); + +test('a stale coordinator cannot dispatch or overwrite a newer claim', async () => { + const requests = await setup(); await submit(requests); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const record = JSON.parse((await redis.hget(key, 'record'))!); + await redis.set(`${key}:claim`, 'new-owner', 'PX', 10000); + await bridge.advanceDurableWorkspaceTool(record, { key, claimKey: `${key}:claim`, token: 'old-owner' }); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); + expect((await requests.get(owner, requestId))?.state).toBe('queued'); +}); + +test('a queued deadline expires as definitely unstarted work', async () => { + const requests = await setup(); await submit(requests); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + await mutateRecord(record => { record.queueDeadlineAtMs = Date.now() - 1; }); + await redis.hset(key, 'queueDeadlineAtMs', Date.now() - 1); + await tick(requests); + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'failed', error: { code: 'WORKSPACE_QUEUE_TIMEOUT' } }); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); +}); + +test('an expired acknowledged command is unknown, remains fenced, and is not re-enqueued on restart', async () => { + const requests = await setup(); await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await bridge.acknowledgeLease(workerId, incarnationId, assignment.assignmentId, assignment.generation, assignment.leaseToken); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const stored = JSON.parse((await redis.hget(key, 'assignment'))!); + stored.expiresAt = new Date(Date.now() - 1).toISOString(); + await redis.hset(key, 'assignment', JSON.stringify(stored)); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'failed', error: { code: 'ASSIGNMENT_EXPIRED' } }); + await submit(requests); + expect(await redis.llen(`codeapi:bridge:v1:worker:${workerId}:incarnation:${incarnationId}:assignments`)).toBe(0); + expect(await redis.get(`codeapi:bridge:v1:worker:${workerId}:workspace:${createHash('sha256').update('native-workspace:primary').digest('hex')}:quarantined`)).toBe(assignment.assignmentId); +}); + +test('cancel during command execution waits for clean rejection and does not replay', async () => { + const requests = await setup(); + await requests.submit({ owner, requestId, workerId, requireTenantBinding: false, + request: { protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'sleep 30' }, + queueWaitMs: 300_000, executionTimeoutMs: 35_000, + }); + await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await bridge.acknowledgeLease(workerId, incarnationId, assignment.assignmentId, assignment.generation, assignment.leaseToken); + await requests.cancel(owner, requestId); + const cancelling = tick(requests); + const deadline = Date.now() + 1000; + while (!await bridge.cancelled(workerId, incarnationId, assignment.assignmentId)) { + if (Date.now() >= deadline) throw new Error('Worker never received cancellation'); + await new Promise(resolve => setTimeout(resolve, 5)); + } + await settle(assignment, false); + await cancelling; + expect((await requests.get(owner, requestId))?.state).toBe('cancelled'); +}); + +test('restart between workspace commit and cleanup preserves the result and releases the same assignment', async () => { + const requests = await setup(); await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; await settle(assignment); + const original = redis.eval.bind(redis); + let failed = false; + const failure = spyOn(redis, 'eval').mockImplementation((async (...args: Parameters) => { + if (!failed && String(args[0]).includes('local queued = redis.call(\'LREM\'')) { + failed = true; throw new Error('process stopped after commit'); + } + return original(...args); + }) as Redis['eval']); + try { await expect(tick(requests)).rejects.toThrow('process stopped after commit'); } finally { failure.mockRestore(); } + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect((await requests.get(owner, requestId))?.state).toBe('completed'); + expect(await redis.zcard(activeKey)).toBe(0); +}); diff --git a/service/src/workspace-tools/requests.ts b/service/src/workspace-tools/requests.ts new file mode 100644 index 00000000..0cb301d2 --- /dev/null +++ b/service/src/workspace-tools/requests.ts @@ -0,0 +1,244 @@ +import { createHash, randomBytes } from 'crypto'; +import type Redis from 'ioredis'; +import type { WorkspaceToolRequest, WorkspaceToolResult } from '../../../packages/code/src/protocol'; +import type { RedisBridgeStore, RegisteredBridgeWorker, DurableWorkspaceAssignment, CodeBridgeWorkspaceSettlement } from '../bridge/store'; +import { BridgeAdmissionQueue } from '../bridge/admission'; +import { BridgeStoreError, workspaceAdmissionId } from '../bridge/store'; + +const PREFIX = 'codeapi:workspace-requests:v1'; +const ACTIVE = `${PREFIX}:active`; +const RETENTION_MS = 24 * 60 * 60_000; +const CLAIM_MS = 10_000; +export const WORKSPACE_REQUEST_ID_PATTERN = /^[A-Za-z0-9_-]{16,128}$/; + +export class WorkspaceRequestConflict extends Error {} + +export interface WorkspaceRequestOwner { + tenantId: string; + userId: string; +} + +export interface StoredWorkspaceRequest { + id: string; + key: string; + fingerprint: string; + workerId: string; + tenantId: string; + requireTenantBinding: boolean; + request: WorkspaceToolRequest; + registration: RegisteredBridgeWorker; + createdAtMs: number; + queueDeadlineAtMs: number; + executionTimeoutMs: number; + assignment?: DurableWorkspaceAssignment; + settlement?: CodeBridgeWorkspaceSettlement; + admittedAtMs?: number; + finishedAtMs?: number; + state: 'queued' | 'admitted' | 'completed' | 'failed' | 'cancelled'; + cancelRequested?: boolean; + result?: WorkspaceToolResult; + error?: { code: string; message: string }; +} + +export interface WorkspaceRequestStatus { + requestId: string; + state: StoredWorkspaceRequest['state']; + workerId: string; + workspaceId: string; + workspaceInstanceId?: string; + worktree?: string; + queuePosition?: number; + queueWaitMs: number; + executionDeadlineAt?: string; + cancelRequested: boolean; + result?: WorkspaceToolResult; + error?: StoredWorkspaceRequest['error']; +} + +function canonical(value: unknown): string { + if (Array.isArray(value)) return `[${value.map(canonical).join(',')}]`; + if (value !== null && typeof value === 'object') { + const record = value as Record; + return `{${Object.keys(record).filter(key => record[key] !== undefined).sort() + .map(key => `${JSON.stringify(key)}:${canonical(record[key])}`).join(',')}}`; + } + return JSON.stringify(value); +} + +async function deadline(operation: Promise): Promise { + let timer: ReturnType | undefined; + try { + return await Promise.race([operation, new Promise((_, reject) => { + timer = setTimeout(() => reject(new Error('Durable workspace transition timed out')), 8_000); + timer.unref(); + })]); + } finally { clearTimeout(timer); } +} + +function requestKey(owner: WorkspaceRequestOwner, requestId: string): string { + return `${PREFIX}:${createHash('sha256').update(canonical([owner.tenantId, owner.userId, requestId])).digest('hex')}`; +} + +/** Redis owns requests; API replicas only advance bounded, fenced transitions. */ +export class RedisWorkspaceRequests { + constructor( + private readonly redis: Redis, + private readonly bridge: RedisBridgeStore, + private readonly transition?: (record: StoredWorkspaceRequest) => void, + ) {} + + async submit(args: { + owner: WorkspaceRequestOwner; + requestId: string; + workerId: string; + requireTenantBinding: boolean; + request: WorkspaceToolRequest; + queueWaitMs: number; + executionTimeoutMs: number; + }): Promise { + if (!WORKSPACE_REQUEST_ID_PATTERN.test(args.requestId)) { + throw new BridgeStoreError('ASSIGNMENT_INVALID', 'Invalid durable workspace request ID'); + } + const key = requestKey(args.owner, args.requestId); + const fingerprint = createHash('sha256').update(canonical({ + workerId: args.workerId, request: args.request, + queueWaitMs: args.queueWaitMs, executionTimeoutMs: args.executionTimeoutMs, + requireTenantBinding: args.requireTenantBinding, + })).digest('hex'); + const existing = await this.read(key); + if (existing != null) { + if (existing.fingerprint !== fingerprint) throw new WorkspaceRequestConflict('Request ID was already used for different work'); + return this.status(existing, args.requestId); + } + const registration = await this.bridge.prepareDurableWorkspaceTool({ + ...args, tenantId: args.owner.tenantId, + }); + const now = Date.now(); + const record: StoredWorkspaceRequest = { + id: `durable-${randomBytes(18).toString('base64url')}`, key, fingerprint, + workerId: args.workerId, tenantId: args.owner.tenantId, + requireTenantBinding: args.requireTenantBinding, + request: args.request, registration, createdAtMs: now, + queueDeadlineAtMs: now + args.queueWaitMs, + executionTimeoutMs: args.executionTimeoutMs, state: 'queued', + }; + const queue = new BridgeAdmissionQueue(this.redis); + const accepted = await queue.submit({ + workerId: record.workerId, id: record.id, deadlineAtMs: record.queueDeadlineAtMs, + workspaceId: (registration.capabilities.workspaceLeaseSlots ?? 1) > 1 + ? workspaceAdmissionId(record.request.workspaceId, record.request.workspaceInstanceId, record.request.worktree) + : undefined, + key, activeKey: ACTIVE, fingerprint, record: JSON.stringify(record), retentionMs: RETENTION_MS, + }); + if (accepted === 'conflict') throw new WorkspaceRequestConflict('Request ID was already used for different work'); + if (accepted === 'full') throw new BridgeStoreError('WORKER_QUEUE_FULL', 'Bridge worker pending request limit reached'); + const submitted = (await this.read(key))!; + if (accepted === 'accepted') this.transition?.(submitted); + return this.status(submitted, args.requestId); + } + + async get(owner: WorkspaceRequestOwner, requestId: string): Promise { + if (!WORKSPACE_REQUEST_ID_PATTERN.test(requestId)) return undefined; + const record = await this.read(requestKey(owner, requestId)); + return record == null ? undefined : this.status(record, requestId); + } + + async cancel(owner: WorkspaceRequestOwner, requestId: string): Promise { + if (!WORKSPACE_REQUEST_ID_PATTERN.test(requestId)) return undefined; + const key = requestKey(owner, requestId); + await this.redis.eval([ + 'local state = redis.call(\'HGET\', KEYS[1], \'state\')', + 'if state ~= \'queued\' and state ~= \'admitted\' then return 0 end', + 'redis.call(\'HSET\', KEYS[1], \'cancelRequested\', \'1\')', + 'redis.call(\'ZADD\', KEYS[2], ARGV[1], KEYS[1])', + 'return 1', + ].join('\n'), 2, key, ACTIVE, Date.now()); + return this.get(owner, requestId); + } + + private async read(key: string): Promise { + const fields: Partial> = await this.redis.hgetall(key); + if (fields.record == null) return undefined; + return { + ...JSON.parse(fields.record) as StoredWorkspaceRequest, + ...(fields.outcome == null ? {} : JSON.parse(fields.outcome)), + state: fields.state as StoredWorkspaceRequest['state'], + assignment: fields.assignment == null ? undefined : JSON.parse(fields.assignment), + admittedAtMs: fields.admittedAtMs == null ? undefined : Number(fields.admittedAtMs), + cancelRequested: fields.cancelRequested === '1', + }; + } + + private async status(record: StoredWorkspaceRequest, requestId: string): Promise { + const rank = record.state === 'queued' + ? await new BridgeAdmissionQueue(this.redis).position(record.workerId, record.id) + : undefined; + return { + requestId, state: record.state, workerId: record.workerId, + workspaceId: record.request.workspaceId, + workspaceInstanceId: record.request.workspaceInstanceId, worktree: record.request.worktree, + queuePosition: rank, queueWaitMs: Math.max(0, (record.admittedAtMs ?? Math.min(record.finishedAtMs ?? Date.now(), record.queueDeadlineAtMs)) - record.createdAtMs), + executionDeadlineAt: record.assignment?.expiresAt, + cancelRequested: record.cancelRequested === true, + result: record.result, error: record.error, + }; + } + + async reconcile(): Promise { + const keys = await deadline(this.redis.zrangebyscore(ACTIVE, '-inf', Date.now(), 'LIMIT', 0, 32)); + const results = await Promise.allSettled(keys.map(key => deadline(this.advance(key)))); + const failure = results.find(result => result.status === 'rejected'); + if (failure?.status === 'rejected') throw failure.reason; + } + + private async advance(key: string): Promise { + const claimKey = `${key}:claim`; + const token = randomBytes(18).toString('hex'); + const claimed = await this.redis.eval([ + 'if not redis.call(\'SET\', KEYS[1], ARGV[1], \'PX\', ARGV[2], \'NX\') then return 0 end', + 'redis.call(\'ZADD\', KEYS[2], ARGV[3], KEYS[3])', + 'return 1', + ].join('\n'), 3, claimKey, ACTIVE, key, token, CLAIM_MS, Date.now() + CLAIM_MS); + if (Number(claimed) !== 1) return; + try { + let record = await this.read(key); + if (record == null) { await this.redis.zrem(ACTIVE, key); return; } + if (record.state === 'queued' || record.state === 'admitted') { + const previousState = record.state; + try { + const outcome = await deadline(this.bridge.advanceDurableWorkspaceTool(record, { key, claimKey, token })); + if (outcome != null) { + await this.redis.eval([ + 'if redis.call(\'GET\', KEYS[2]) ~= ARGV[1] then return 0 end', + 'local state = redis.call(\'HGET\', KEYS[1], \'state\')', + 'if state ~= \'queued\' and state ~= \'admitted\' then return 0 end', + 'redis.call(\'HSET\', KEYS[1], \'state\', ARGV[3], \'outcome\', ARGV[2])', + 'return 1', + ].join('\n'), 2, key, claimKey, token, JSON.stringify({ ...outcome, finishedAtMs: Date.now() }), outcome.state); + } + } catch (error) { + // Infrastructure uncertainty is reconciled, never retried as a new command. + if (!(error instanceof BridgeStoreError)) throw error; + const current = await this.read(key); + if (current?.state === 'admitted') throw error; + await this.redis.eval([ + 'if redis.call(\'GET\', KEYS[2]) ~= ARGV[1] then return 0 end', + 'if redis.call(\'HGET\', KEYS[1], \'state\') ~= \'queued\' then return 0 end', + 'redis.call(\'HSET\', KEYS[1], \'state\', \'failed\', \'outcome\', ARGV[2])', + 'return 1', + ].join('\n'), 2, key, claimKey, token, JSON.stringify({ state: 'failed', finishedAtMs: Date.now(), error: { code: error.code, message: error.message } })); + } + record = (await this.read(key))!; + if (record.state !== previousState) this.transition?.(record); + } + if (record.state !== 'queued' && record.state !== 'admitted') { + await deadline(this.bridge.finishDurableWorkspaceTool(record)); + await this.redis.zrem(ACTIVE, key); + } else { + await this.redis.zadd(ACTIVE, Date.now() + 250, key); + } + } finally { + await this.redis.eval('if redis.call(\'GET\', KEYS[1]) == ARGV[1] then return redis.call(\'DEL\', KEYS[1]) end return 0', 1, claimKey, token); + } + } +} diff --git a/service/src/workspace-tools/router.ts b/service/src/workspace-tools/router.ts index 78dbd1e7..83576e38 100644 --- a/service/src/workspace-tools/router.ts +++ b/service/src/workspace-tools/router.ts @@ -19,6 +19,8 @@ import { resolveBridgeWorkerSelection, } from '../bridge/selection'; import { principalWorkspaceInstanceId } from '../bridge/workspace-instance'; +import { WorkspaceRequestConflict } from './requests'; +import type { RedisWorkspaceRequests } from './requests'; const MAX_WORKSPACE_QUEUE_WAIT_MS = 5 * 60_000; const DEFAULT_WORKSPACE_QUEUE_WAIT_MS = 30_000; @@ -26,6 +28,7 @@ const WORKSPACE_QUEUE_WAIT_HEADER = 'X-LibreChat-Workspace-Queue-Wait-Ms'; interface WorkspaceToolsRouterOptions { store: Pick; + requests?: RedisWorkspaceRequests; backend: 'http' | 'lambda-microvm' | 'remote-bridge'; configuredWorkerId: string; dynamicWorkers: boolean; @@ -68,9 +71,9 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) const router = Router(); router.post( - '/workspace-tools/execute', + ['/workspace-tools/execute', '/workspace-tools/requests'], asyncRoute(async (req, res) => { - const outcome = getWorkspaceToolOutcome(res); + const outcome = getWorkspaceToolOutcome(res, req.path); const principal = getPrincipalOrReject(req, res); if (!principal) { outcome.errorCode = 'UNAUTHENTICATED'; @@ -151,6 +154,29 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) } outcome.workerId = selection.workerId; + if (req.path === '/workspace-tools/requests') { + if (options.requests == null) { + res.status(404).json({ error: 'Durable workspace requests are unavailable' }); + return; + } + try { + const status = await options.requests.submit({ + owner: principal, requestId: req.header('X-LibreChat-Workspace-Request-Id') ?? '', + workerId: selection.workerId, + requireTenantBinding: selection.explicit && (options.dynamicWorkers || selection.workerId !== options.configuredWorkerId), + request, queueWaitMs: queueBudgetMs, executionTimeoutMs: executionBudgetMs, + }); + res.status(202).json(status); + } catch (error) { + if (error instanceof WorkspaceRequestConflict) { + res.status(409).json({ error: error.message, code: 'REQUEST_CONFLICT' }); + } else if (error instanceof BridgeStoreError) { + res.status(bridgeStoreStatus(error)).json({ error: error.message, code: error.code }); + } else throw error; + } + return; + } + const controller = new AbortController(); const abort = (): void => controller.abort(); req.once('aborted', abort); @@ -223,5 +249,22 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) }), ); + router.get('/workspace-tools/capabilities', asyncRoute(async (req, res) => { + if (!getPrincipalOrReject(req, res)) return; + res.json({ durableWorkspaceRequests: options.requests != null ? 1 : 0 }); + })); + + const requestStatus = (cancel: boolean): RequestHandler => asyncRoute(async (req, res) => { + const principal = getPrincipalOrReject(req, res); + if (!principal) return; + const status = cancel + ? await options.requests?.cancel(principal, req.params.requestId) + : await options.requests?.get(principal, req.params.requestId); + if (status == null) { res.status(404).json({ error: 'Workspace request not found' }); return; } + res.json(status); + }); + router.get('/workspace-tools/requests/:requestId', requestStatus(false)); + router.delete('/workspace-tools/requests/:requestId', requestStatus(true)); + return router; } From 1bfe89c95ff346864dc520a4271e74bb6ae10974 Mon Sep 17 00:00:00 2001 From: Lia Date: Sat, 3 Oct 2026 04:00:15 +0000 Subject: [PATCH 2/5] =?UTF-8?q?=F0=9F=A7=BE=20fix:=20Fence=20Durable=20Adm?= =?UTF-8?q?ission=20Recovery=20Across=20Replicas?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/remote-bridge/README.md | 9 +- service/src/bridge/store.ts | 17 +- service/src/workspace-tools/coordination.ts | 21 ++ service/src/workspace-tools/index.ts | 37 +- .../workspace-tools/requests-router.test.ts | 16 +- service/src/workspace-tools/requests.test.ts | 114 ++++++- service/src/workspace-tools/requests.ts | 35 +- service/src/workspace-tools/router.ts | 323 +++++++++--------- 8 files changed, 371 insertions(+), 201 deletions(-) create mode 100644 service/src/workspace-tools/coordination.ts diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index 706e94e8..e5c4d177 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -415,8 +415,13 @@ bounded by that retention window. A client must not resubmit an expired ID after 404; its outcome is unknown. Automatic result claims and wake-ups remain the client's responsibility. -API replicas reconcile accepted work from Redis. FIFO admission position survives -API restarts. Assignment enqueue and the durable `admitted` transition are atomic. +Bridge-enabled API replicas reconcile accepted work from Redis within a shared +bridge-policy scope. Disabled bridges and different scheduling policies never +claim that work. Keep backend/profile, bridge auth mode, worker selection, slot +ceiling and command-timeout policy identical on replicas behind one endpoint. +Drain durable requests before changing those settings. FIFO admission position +survives API restarts and temporary worker registration/readiness loss until the +queue deadline. Assignment enqueue and the durable `admitted` transition are atomic. Recovery resumes only queued work and observes admitted assignments without creating another command. Worker identity/incarnation changes fail pending work rather than transferring it to another machine. Execution gets its full budget diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index a2ba6115..ea173e74 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -864,7 +864,8 @@ export class RedisBridgeStore { request: record.request, }; const ready = await this.dispatchableRegistration(record.workerId); - if (ready == null || ready.registration.incarnationId !== registration.incarnationId || + if (ready == null) throw new BridgeStoreError('WORKER_OFFLINE', 'Code environment is temporarily unavailable'); + if (ready.registration.incarnationId !== registration.incarnationId || ready.registration.identityId !== registration.identityId || ready.registration.binding?.tenantId !== registration.binding?.tenantId || !supportsWorkspaceTool(ready.registration, record.request)) { @@ -2407,8 +2408,8 @@ export class RedisBridgeStore { "if ARGV[9] == 'lane' then redis.call('SADD', KEYS[11], KEYS[7]) end", 'if keyCount >= 9 then', " local epoch = redis.call('GET', KEYS[9])", - " if type(epoch) ~= 'string' then epoch = '0'; redis.call('SET', KEYS[9], epoch, 'EX', ARGV[3]) end", - " if redis.call('PTTL', KEYS[9]) < tonumber(ARGV[3]) * 1000 then redis.call('EXPIRE', KEYS[9], ARGV[3]) end", + ' if type(epoch) ~= \'string\' then epoch = \'0\'; redis.call(\'SET\', KEYS[9], epoch, \'EX\', ARGV[13]) end', + ' if redis.call(\'PTTL\', KEYS[9]) < tonumber(ARGV[13]) * 1000 then redis.call(\'EXPIRE\', KEYS[9], ARGV[13]) end', " redis.call('HSET', KEYS[8], 'metadata', ARGV[8], 'epoch', epoch)", ' redis.call(\'EXPIRE\', KEYS[8], ARGV[13])', 'end', @@ -2614,9 +2615,13 @@ export class RedisBridgeStore { 'if redis.call(\'EXISTS\', KEYS[2]) == 1 then return 1 end', ]), "if redis.call('GET', KEYS[1]) == ARGV[1] then", - ' redis.call(\'DEL\', KEYS[1])', - ...(durableKey == null ? [] : [' redis.call(\'SET\', KEYS[2], \'1\', \'PX\', 86400000)']), - ' return 1', + ...(durableKey == null + ? [' return redis.call(\'DEL\', KEYS[1])'] + : [ + ' redis.call(\'DEL\', KEYS[1])', + ' redis.call(\'SET\', KEYS[2], \'1\', \'PX\', 86400000)', + ' return 1', + ]), 'end', 'return 0', ].join('\n'); diff --git a/service/src/workspace-tools/coordination.ts b/service/src/workspace-tools/coordination.ts new file mode 100644 index 00000000..0bcd5e91 --- /dev/null +++ b/service/src/workspace-tools/coordination.ts @@ -0,0 +1,21 @@ +import { createHash } from 'crypto'; + +export interface WorkspaceRequestCoordinationPolicy { + bridgeEnabled: boolean; + backend: string; + executionProfile: string; + authMode: string; + configuredWorkerId: string; + dynamicWorkers: boolean; + maxWorkspaceLeaseSlots: number; + maxCommandTimeoutMs: number; +} + +/** Only replicas with the same bridge policy may reconcile accepted work. */ +export function workspaceRequestCoordinationScope(policy: WorkspaceRequestCoordinationPolicy): string | undefined { + if (!policy.bridgeEnabled) return undefined; + return createHash('sha256').update(JSON.stringify([ + policy.backend, policy.executionProfile, policy.authMode, policy.configuredWorkerId, + policy.dynamicWorkers, policy.maxWorkspaceLeaseSlots, policy.maxCommandTimeoutMs, + ])).digest('hex'); +} diff --git a/service/src/workspace-tools/index.ts b/service/src/workspace-tools/index.ts index 53f750c1..3c7caee2 100644 --- a/service/src/workspace-tools/index.ts +++ b/service/src/workspace-tools/index.ts @@ -9,9 +9,17 @@ import { connection } from '../queue'; import logger from '../logger'; import { checkServiceStartUp, checkServiceShutDown } from '../lifecycle'; import { withSpan } from '../telemetry'; +import { isBridgeEnabled } from '../bridge/enabled'; +import { workspaceRequestCoordinationScope } from './coordination'; const router = Router(); -const requests = new RedisWorkspaceRequests(connection, bridgeStore, record => { +const coordinationScope = workspaceRequestCoordinationScope({ + bridgeEnabled: isBridgeEnabled(), backend: env.SANDBOX_BACKEND, + executionProfile: env.EXECUTION_PROFILE, authMode: env.BRIDGE_AUTH_MODE, + configuredWorkerId: env.BRIDGE_WORKER_ID, dynamicWorkers: env.BRIDGE_DYNAMIC_WORKERS, + maxWorkspaceLeaseSlots: env.BRIDGE_MAX_WORKSPACE_LEASE_SLOTS, maxCommandTimeoutMs: env.JOB_TIMEOUT, +}); +const requests = coordinationScope === undefined ? undefined : new RedisWorkspaceRequests(connection, bridgeStore, record => { const attributes = { 'codeapi.request.id': record.id, 'codeapi.admission.state': record.state, @@ -26,18 +34,21 @@ const requests = new RedisWorkspaceRequests(connection, bridgeStore, record => { void withSpan('codeapi.workspace.admission.transition', attributes, () => undefined).catch(error => { logger.warn('Admission transition telemetry failed', { error }); }); -}); -let reconciling = false; -let unavailable = false; -const reconciliation = setInterval(() => { - if (reconciling || checkServiceStartUp() || checkServiceShutDown()) return; - reconciling = true; - void requests.reconcile().then(() => { unavailable = false; }).catch(error => { - if (!unavailable) logger.warn('Durable workspace reconciliation unavailable', { error }); - unavailable = true; - }).finally(() => { reconciling = false; }); -}, 250); -reconciliation.unref(); +}, coordinationScope); +if (requests != null) { + const coordinator = requests; + let reconciling = false; + let unavailable = false; + const reconciliation = setInterval(() => { + if (reconciling || checkServiceStartUp() || checkServiceShutDown()) return; + reconciling = true; + void coordinator.reconcile().then(() => { unavailable = false; }).catch(error => { + if (!unavailable) logger.warn('Durable workspace reconciliation unavailable', { error }); + unavailable = true; + }).finally(() => { reconciling = false; }); + }, 250); + reconciliation.unref(); +} router.use(['/workspace-tools/execute', '/workspace-tools/requests'], (req, res, next) => { if (req.method === 'POST') executionLimiter(req, res, next); else next(); diff --git a/service/src/workspace-tools/requests-router.test.ts b/service/src/workspace-tools/requests-router.test.ts index 27ea4710..856a7a38 100644 --- a/service/src/workspace-tools/requests-router.test.ts +++ b/service/src/workspace-tools/requests-router.test.ts @@ -14,7 +14,7 @@ afterEach(() => { server?.close(); server = undefined; }); const body = { protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }; const requestId = 'http-request-000001'; -async function setup(durable = true): Promise { +async function setup(durable = true, rejectSynchronous = false): Promise { const redis = new RedisMock() as unknown as Redis; const bridge = new RedisBridgeStore(redis); await bridge.register({ protocolVersion: 1, workerId: 'worker', incarnationId: 'incarnation-00000001', capabilities: { @@ -30,7 +30,7 @@ async function setup(durable = true): Promise { next(); }); app.use(createWorkspaceToolsRouter({ - store: bridge, requests: durable ? new RedisWorkspaceRequests(redis, bridge) : undefined, + store: rejectSynchronous ? { async dispatchWorkspaceTool(): Promise { throw new Error('Durable URL dispatched synchronously'); } } : bridge, requests: durable ? new RedisWorkspaceRequests(redis, bridge) : undefined, backend: 'remote-bridge', configuredWorkerId: 'worker', dynamicWorkers: false, })); server = createServer(app); @@ -82,3 +82,15 @@ test('unauthenticated request status and capability discovery are rejected', asy expect((await fetch(`${url}/workspace-tools/${path}`, { headers: { 'X-Test-User': 'unauthenticated' } })).status).toBe(401); } }); + +test.each(['/workspace-tools/requests/', '/workspace-tools/REQUESTS'])('durable route aliases remain idempotent (%s)', async path => { + const url = await setup(true, true); + for (let attempt = 0; attempt < 2; attempt++) { + const response = await fetch(`${url}${path}`, { method: 'POST', headers: { + 'Content-Type': 'application/json', 'X-LibreChat-Workspace-Request-Id': requestId, + }, body: JSON.stringify(body) }); + expect(response.status).toBe(202); + expect(await response.json()).toMatchObject({ requestId, state: 'queued', queuePosition: 1 }); + } + expect(await (await fetch(`${url}/workspace-tools/requests/${requestId}`)).json()).toMatchObject({ state: 'queued' }); +}); diff --git a/service/src/workspace-tools/requests.test.ts b/service/src/workspace-tools/requests.test.ts index 992060cc..850a266d 100644 --- a/service/src/workspace-tools/requests.test.ts +++ b/service/src/workspace-tools/requests.test.ts @@ -5,6 +5,8 @@ import Redis from 'ioredis'; import { RedisBridgeStore } from '../bridge/store'; import type { CodeBridgeAssignment } from '../bridge/store'; import { RedisWorkspaceRequests } from './requests'; +import { workspaceRequestCoordinationScope } from './coordination'; +import type { WorkspaceRequestCoordinationPolicy } from './coordination'; const testRedisUrl = process.env.BRIDGE_TEST_REDIS_URL; const redis = testRedisUrl !== undefined && testRedisUrl.length > 0 @@ -190,15 +192,15 @@ test('durable requests preserve negotiated independent-root concurrency', async expect((await requests.get(owner, 'request-00000000002'))?.state).toBe('completed'); }); -test('restart after capacity reservation but before enqueue refreshes its own reservation without duplicating work', async () => { - const requests = await setup(); +test.each([1, 2])('restart after capacity reservation refreshes its own reservation without duplicating work (slots=%s)', async slots => { + const requests = await setup(slots); await submit(requests); const generation = spyOn(redis, 'incr').mockRejectedValue(new Error('process stopped before enqueue')); try { await expect(tick(requests)).rejects.toThrow('process stopped before enqueue'); } finally { generation.mockRestore(); } expect((await requests.get(owner, requestId))?.state).toBe('queued'); await redis.pexpire(`codeapi:bridge:v1:worker:${workerId}:lock`, 5000); await tick(new RedisWorkspaceRequests(redis, bridge)); - const assignment = await bridge.lease(workerId, incarnationId, 0); + const assignment = await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0); expect(assignment).toBeDefined(); expect(await redis.pttl(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeGreaterThan(30_000); }); @@ -275,3 +277,109 @@ test('restart between workspace commit and cleanup preserves the result and rele expect((await requests.get(owner, requestId))?.state).toBe('completed'); expect(await redis.zcard(activeKey)).toBe(0); }); + +test('simultaneous submissions from replicas keep one request identity and one FIFO entry', async () => { + const requests = await setup(); + const replica = new RedisWorkspaceRequests(redis, bridge); + const statuses = await Promise.all([submit(requests), submit(replica)]); + for (const status of statuses) expect(status).toMatchObject({ requestId, state: 'queued', queuePosition: 1 }); + expect(await redis.zcard(queueKey)).toBe(1); + expect(await redis.zcard(activeKey)).toBe(1); + await tick(requests); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeDefined(); +}); + +test.each([false, true])('temporary registration loss preserves FIFO until the original worker reconnects (replacement=%s)', async replacement => { + const requests = await setup(); + await submit(requests); await submit(requests, 'request-00000000002'); + // Registration/readiness TTL expiry, without losing accepted work. + const worker = `codeapi:bridge:v1:worker:${workerId}`; + await redis.del(worker, `${worker}:incarnation`, `${worker}:ready`); + const restarted = new RedisWorkspaceRequests(redis, bridge); + await tick(restarted); + expect(await restarted.get(owner, requestId)).toMatchObject({ state: 'queued', queuePosition: 1 }); + expect(await redis.zcard(activeKey)).toBe(2); + if (replacement) { + await bridge.register({ protocolVersion: 1, workerId, incarnationId: 'replacement-00001', + capabilities: { sandboxProfile: 'native-srt', statefulWorkspace: false, runtimes: [], + workspaceTools: { protocolVersion: 1, operations: ['read_file'], workspaces: [{ id: 'primary' }] } }, + }); + await tick(restarted); + expect(await restarted.get(owner, requestId)).toMatchObject({ state: 'failed', error: { code: 'WORKER_FENCED' } }); + expect(await bridge.lease(workerId, 'replacement-00001', 0)).toBeUndefined(); + } else { + await setup(); await tick(restarted); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + expect(assignment).toBeDefined(); + await settle(assignment); await tick(restarted); + expect((await restarted.get(owner, requestId))?.state).toBe('completed'); + } +}); + +test('an offline worker cannot extend the queue deadline', async () => { + const requests = await setup(); await submit(requests); + const worker = `codeapi:bridge:v1:worker:${workerId}`; + await redis.del(worker, `${worker}:incarnation`, `${worker}:ready`); + await tick(requests); + expect((await requests.get(owner, requestId))?.state).toBe('queued'); + await mutateRecord(record => { record.queueDeadlineAtMs = Date.now() - 1; }); + await tick(requests); + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'failed', error: { code: 'WORKSPACE_QUEUE_TIMEOUT' } }); + expect(await redis.zcard(activeKey)).toBe(0); +}); + +test('durable reset epochs outlive their receipts and reject late quarantine after reset', async () => { + const requests = await setup(2); await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0))!; + await settle(assignment); await tick(requests); + const receipt = `codeapi:bridge:v1:assignment:${assignment.assignmentId}:workspace-fence-owner`; + const fence = `codeapi:bridge:v1:worker:${workerId}:workspace:${createHash('sha256').update('native-workspace:primary').digest('hex')}:quarantined`; + expect(await redis.pttl(`${fence}:epoch`)).toBeGreaterThanOrEqual(await redis.pttl(receipt) - 50); + expect(await redis.pttl(`${fence}:epoch`)).toBeGreaterThan(86_000_000); + // Simulate execution-ownership TTL expiry. The durable receipt and epoch survive. + await redis.del(`codeapi:bridge:v1:worker:${workerId}:workspace-slots`, + `codeapi:bridge:v1:worker:${workerId}:lock`, `codeapi:bridge:v1:worker:${workerId}:lock:incarnation`); + await bridge.resetWorkspace(workerId, incarnationId, 'native-workspace:primary'); + await submit(requests, 'request-00000000002'); await tick(requests); + const next = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0))!; + await expect(bridge.settle(workerId, assignment.assignmentId, { + protocolVersion: 1, incarnationId, generation: assignment.generation, leaseToken: assignment.leaseToken, + status: 'rejected', error: 'delayed local cleanup failure', + }, undefined, undefined, true)).rejects.toMatchObject({ code: 'ASSIGNMENT_FENCED' }); + expect(await redis.get(fence)).toBe(next.assignmentId); + await settle(next); await tick(requests); + expect((await requests.get(owner, 'request-00000000002'))?.state).toBe('completed'); +}); + +const coordinationPolicy: WorkspaceRequestCoordinationPolicy = { + bridgeEnabled: true, backend: 'remote-bridge', executionProfile: 'default', authMode: 'paired', + configuredWorkerId: '', dynamicWorkers: true, maxWorkspaceLeaseSlots: 2, maxCommandTimeoutMs: 30_000, +}; + +test('disabled bridges do not participate and each scheduling policy has an isolated scope', () => { + expect(workspaceRequestCoordinationScope({ ...coordinationPolicy, bridgeEnabled: false })).toBeUndefined(); + const scope = workspaceRequestCoordinationScope(coordinationPolicy); + expect(workspaceRequestCoordinationScope({ ...coordinationPolicy })).toBe(scope); + const variants: Partial[] = [ + { maxWorkspaceLeaseSlots: 1 }, { backend: 'http' }, { executionProfile: 'stateful' }, + { authMode: 'static' }, { configuredWorkerId: 'another' }, { dynamicWorkers: false }, + { maxCommandTimeoutMs: 60_000 }, + ]; + for (const variant of variants) expect(workspaceRequestCoordinationScope({ ...coordinationPolicy, ...variant })).not.toBe(scope); +}); + +test('a different policy cannot claim or fail work accepted by a two-slot bridge API', async () => { + await setup(2); + const scope = workspaceRequestCoordinationScope(coordinationPolicy)!; + const requests = new RedisWorkspaceRequests(redis, bridge, undefined, scope); + await submit(requests); + const foreignScope = workspaceRequestCoordinationScope({ ...coordinationPolicy, maxWorkspaceLeaseSlots: 1 })!; + const foreign = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis), undefined, foreignScope); + await foreign.reconcile(); + expect((await requests.get(owner, requestId))?.state).toBe('queued'); + expect(await foreign.get(owner, requestId)).toBeUndefined(); + const restarted = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis, 600, 1000, 2), undefined, scope); + await restarted.reconcile(); + const assignment = await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0); + expect(assignment).toBeDefined(); +}); diff --git a/service/src/workspace-tools/requests.ts b/service/src/workspace-tools/requests.ts index 0cb301d2..a88b9acb 100644 --- a/service/src/workspace-tools/requests.ts +++ b/service/src/workspace-tools/requests.ts @@ -75,17 +75,21 @@ async function deadline(operation: Promise): Promise { } finally { clearTimeout(timer); } } -function requestKey(owner: WorkspaceRequestOwner, requestId: string): string { - return `${PREFIX}:${createHash('sha256').update(canonical([owner.tenantId, owner.userId, requestId])).digest('hex')}`; +function requestKey(owner: WorkspaceRequestOwner, requestId: string, scope: string): string { + return `${PREFIX}:${createHash('sha256').update(canonical([scope, owner.tenantId, owner.userId, requestId])).digest('hex')}`; } /** Redis owns requests; API replicas only advance bounded, fenced transitions. */ export class RedisWorkspaceRequests { + private readonly activeKey: string; constructor( private readonly redis: Redis, private readonly bridge: RedisBridgeStore, private readonly transition?: (record: StoredWorkspaceRequest) => void, - ) {} + private readonly scope = '', + ) { + this.activeKey = scope.length === 0 ? ACTIVE : `${PREFIX}:${scope}:active`; + } async submit(args: { owner: WorkspaceRequestOwner; @@ -99,7 +103,7 @@ export class RedisWorkspaceRequests { if (!WORKSPACE_REQUEST_ID_PATTERN.test(args.requestId)) { throw new BridgeStoreError('ASSIGNMENT_INVALID', 'Invalid durable workspace request ID'); } - const key = requestKey(args.owner, args.requestId); + const key = requestKey(args.owner, args.requestId, this.scope); const fingerprint = createHash('sha256').update(canonical({ workerId: args.workerId, request: args.request, queueWaitMs: args.queueWaitMs, executionTimeoutMs: args.executionTimeoutMs, @@ -128,7 +132,7 @@ export class RedisWorkspaceRequests { workspaceId: (registration.capabilities.workspaceLeaseSlots ?? 1) > 1 ? workspaceAdmissionId(record.request.workspaceId, record.request.workspaceInstanceId, record.request.worktree) : undefined, - key, activeKey: ACTIVE, fingerprint, record: JSON.stringify(record), retentionMs: RETENTION_MS, + key, activeKey: this.activeKey, fingerprint, record: JSON.stringify(record), retentionMs: RETENTION_MS, }); if (accepted === 'conflict') throw new WorkspaceRequestConflict('Request ID was already used for different work'); if (accepted === 'full') throw new BridgeStoreError('WORKER_QUEUE_FULL', 'Bridge worker pending request limit reached'); @@ -139,20 +143,20 @@ export class RedisWorkspaceRequests { async get(owner: WorkspaceRequestOwner, requestId: string): Promise { if (!WORKSPACE_REQUEST_ID_PATTERN.test(requestId)) return undefined; - const record = await this.read(requestKey(owner, requestId)); + const record = await this.read(requestKey(owner, requestId, this.scope)); return record == null ? undefined : this.status(record, requestId); } async cancel(owner: WorkspaceRequestOwner, requestId: string): Promise { if (!WORKSPACE_REQUEST_ID_PATTERN.test(requestId)) return undefined; - const key = requestKey(owner, requestId); + const key = requestKey(owner, requestId, this.scope); await this.redis.eval([ 'local state = redis.call(\'HGET\', KEYS[1], \'state\')', 'if state ~= \'queued\' and state ~= \'admitted\' then return 0 end', 'redis.call(\'HSET\', KEYS[1], \'cancelRequested\', \'1\')', 'redis.call(\'ZADD\', KEYS[2], ARGV[1], KEYS[1])', 'return 1', - ].join('\n'), 2, key, ACTIVE, Date.now()); + ].join('\n'), 2, key, this.activeKey, Date.now()); return this.get(owner, requestId); } @@ -185,7 +189,7 @@ export class RedisWorkspaceRequests { } async reconcile(): Promise { - const keys = await deadline(this.redis.zrangebyscore(ACTIVE, '-inf', Date.now(), 'LIMIT', 0, 32)); + const keys = await deadline(this.redis.zrangebyscore(this.activeKey, '-inf', Date.now(), 'LIMIT', 0, 32)); const results = await Promise.allSettled(keys.map(key => deadline(this.advance(key)))); const failure = results.find(result => result.status === 'rejected'); if (failure?.status === 'rejected') throw failure.reason; @@ -198,11 +202,11 @@ export class RedisWorkspaceRequests { 'if not redis.call(\'SET\', KEYS[1], ARGV[1], \'PX\', ARGV[2], \'NX\') then return 0 end', 'redis.call(\'ZADD\', KEYS[2], ARGV[3], KEYS[3])', 'return 1', - ].join('\n'), 3, claimKey, ACTIVE, key, token, CLAIM_MS, Date.now() + CLAIM_MS); + ].join('\n'), 3, claimKey, this.activeKey, key, token, CLAIM_MS, Date.now() + CLAIM_MS); if (Number(claimed) !== 1) return; try { let record = await this.read(key); - if (record == null) { await this.redis.zrem(ACTIVE, key); return; } + if (record == null) { await this.redis.zrem(this.activeKey, key); return; } if (record.state === 'queued' || record.state === 'admitted') { const previousState = record.state; try { @@ -219,6 +223,11 @@ export class RedisWorkspaceRequests { } catch (error) { // Infrastructure uncertainty is reconciled, never retried as a new command. if (!(error instanceof BridgeStoreError)) throw error; + if (error.code === 'WORKER_OFFLINE') { + // The original worker may reconnect within the remaining admission budget. + await this.redis.zadd(this.activeKey, Date.now() + 250, key); + return; + } const current = await this.read(key); if (current?.state === 'admitted') throw error; await this.redis.eval([ @@ -233,9 +242,9 @@ export class RedisWorkspaceRequests { } if (record.state !== 'queued' && record.state !== 'admitted') { await deadline(this.bridge.finishDurableWorkspaceTool(record)); - await this.redis.zrem(ACTIVE, key); + await this.redis.zrem(this.activeKey, key); } else { - await this.redis.zadd(ACTIVE, Date.now() + 250, key); + await this.redis.zadd(this.activeKey, Date.now() + 250, key); } } finally { await this.redis.eval('if redis.call(\'GET\', KEYS[1]) == ARGV[1] then return redis.call(\'DEL\', KEYS[1]) end return 0', 1, claimKey, token); diff --git a/service/src/workspace-tools/router.ts b/service/src/workspace-tools/router.ts index 83576e38..d4d65d14 100644 --- a/service/src/workspace-tools/router.ts +++ b/service/src/workspace-tools/router.ts @@ -70,184 +70,183 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) } const router = Router(); - router.post( - ['/workspace-tools/execute', '/workspace-tools/requests'], - asyncRoute(async (req, res) => { - const outcome = getWorkspaceToolOutcome(res, req.path); - const principal = getPrincipalOrReject(req, res); - if (!principal) { - outcome.errorCode = 'UNAUTHENTICATED'; - return; - } - if ((options.isShuttingDown ?? checkServiceShutDown)()) { - outcome.errorCode = 'SERVICE_SHUTTING_DOWN'; - res.status(503).json({ error: 'Service is shutting down' }); - return; - } - if (!isWorkspaceToolRequest(req.body)) { - outcome.errorCode = 'INVALID_WORKSPACE_TOOL_REQUEST'; - res.status(400).json({ - error: 'Invalid workspace tool request', - }); - 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' }); + const execute = (durable: boolean): RequestHandler => asyncRoute(async (req, res) => { + const outcome = getWorkspaceToolOutcome(res, req.path); + const principal = getPrincipalOrReject(req, res); + if (!principal) { + outcome.errorCode = 'UNAUTHENTICATED'; + return; + } + if ((options.isShuttingDown ?? checkServiceShutDown)()) { + outcome.errorCode = 'SERVICE_SHUTTING_DOWN'; + res.status(503).json({ error: 'Service is shutting down' }); + return; + } + if (!isWorkspaceToolRequest(req.body)) { + outcome.errorCode = 'INVALID_WORKSPACE_TOOL_REQUEST'; + res.status(400).json({ + error: 'Invalid workspace tool request', + }); + 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 + : { + ...req.body, + workspaceInstanceId: principalWorkspaceInstanceId({ + instanceId: req.body.workspaceInstanceId, + tenantId: principal.tenantId, + principalId: principal.userId, + }), + }; + const request: WorkspaceToolRequest = principalRequest.operation === 'execute_command' + ? { ...principalRequest, timeoutMs: Math.min( + principalRequest.timeoutMs ?? BRIDGE_WORKSPACE_COMMAND_DEFAULT_TIMEOUT_MS, + options.timeoutMs ?? Number.MAX_SAFE_INTEGER, + ) } + : principalRequest; + const executionBudgetMs = request.operation === 'execute_command' + ? request.timeoutMs! + 5_000 + : Math.min(options.timeoutMs ?? 30_000, 30_000); + outcome.deadlineBudgetMs = queueBudgetMs + executionBudgetMs; + + let selection: { workerId: string; explicit: boolean } | undefined; + try { + selection = resolveBridgeWorkerSelection({ + backend: options.backend, + configuredWorkerId: options.configuredWorkerId, + dynamicWorkers: options.dynamicWorkers, + requestedWorkerId: req.header(CODEAPI_BRIDGE_WORKER_HEADER), + trustedWorkerId: principal.codeWorkerId, + }); + } catch (error) { + if (error instanceof BridgeWorkerSelectionError) { + outcome.errorCode = 'WORKER_SELECTION_REJECTED'; + res.status(error.status).json({ error: error.message }); 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' }); + throw error; + } + if (selection == null) { + outcome.errorCode = 'WORKSPACE_BACKEND_UNAVAILABLE'; + res.status(503).json({ + error: 'Workspace tools require the remote-bridge backend', + }); + return; + } + outcome.workerId = selection.workerId; + + if (durable) { + if (options.requests == null) { + res.status(404).json({ error: 'Durable workspace requests are unavailable' }); return; } - const queueBudgetMs = Math.min(requestedQueueWaitMs, queueCeilingMs); - outcome.operation = req.body.operation; - const principalRequest: WorkspaceToolRequest = req.body.workspaceInstanceId == null - ? req.body - : { - ...req.body, - workspaceInstanceId: principalWorkspaceInstanceId({ - instanceId: req.body.workspaceInstanceId, - tenantId: principal.tenantId, - principalId: principal.userId, - }), - }; - const request: WorkspaceToolRequest = principalRequest.operation === 'execute_command' - ? { ...principalRequest, timeoutMs: Math.min( - principalRequest.timeoutMs ?? BRIDGE_WORKSPACE_COMMAND_DEFAULT_TIMEOUT_MS, - options.timeoutMs ?? Number.MAX_SAFE_INTEGER, - ) } - : principalRequest; - const executionBudgetMs = request.operation === 'execute_command' - ? request.timeoutMs! + 5_000 - : Math.min(options.timeoutMs ?? 30_000, 30_000); - outcome.deadlineBudgetMs = queueBudgetMs + executionBudgetMs; - - let selection: { workerId: string; explicit: boolean } | undefined; try { - selection = resolveBridgeWorkerSelection({ - backend: options.backend, - configuredWorkerId: options.configuredWorkerId, - dynamicWorkers: options.dynamicWorkers, - requestedWorkerId: req.header(CODEAPI_BRIDGE_WORKER_HEADER), - trustedWorkerId: principal.codeWorkerId, + const status = await options.requests.submit({ + owner: principal, requestId: req.header('X-LibreChat-Workspace-Request-Id') ?? '', + workerId: selection.workerId, + requireTenantBinding: selection.explicit && (options.dynamicWorkers || selection.workerId !== options.configuredWorkerId), + request, queueWaitMs: queueBudgetMs, executionTimeoutMs: executionBudgetMs, }); + res.status(202).json(status); } catch (error) { - if (error instanceof BridgeWorkerSelectionError) { - outcome.errorCode = 'WORKER_SELECTION_REJECTED'; - res.status(error.status).json({ error: error.message }); - return; - } - throw error; + if (error instanceof WorkspaceRequestConflict) { + res.status(409).json({ error: error.message, code: 'REQUEST_CONFLICT' }); + } else if (error instanceof BridgeStoreError) { + res.status(bridgeStoreStatus(error)).json({ error: error.message, code: error.code }); + } else throw error; } - if (selection == null) { - outcome.errorCode = 'WORKSPACE_BACKEND_UNAVAILABLE'; - res.status(503).json({ - error: 'Workspace tools require the remote-bridge backend', - }); - return; - } - outcome.workerId = selection.workerId; + return; + } - if (req.path === '/workspace-tools/requests') { - if (options.requests == null) { - res.status(404).json({ error: 'Durable workspace requests are unavailable' }); - return; - } - try { - const status = await options.requests.submit({ - owner: principal, requestId: req.header('X-LibreChat-Workspace-Request-Id') ?? '', - workerId: selection.workerId, - requireTenantBinding: selection.explicit && (options.dynamicWorkers || selection.workerId !== options.configuredWorkerId), - request, queueWaitMs: queueBudgetMs, executionTimeoutMs: executionBudgetMs, - }); - res.status(202).json(status); - } catch (error) { - if (error instanceof WorkspaceRequestConflict) { - res.status(409).json({ error: error.message, code: 'REQUEST_CONFLICT' }); - } else if (error instanceof BridgeStoreError) { - res.status(bridgeStoreStatus(error)).json({ error: error.message, code: error.code }); - } else throw error; - } - return; - } - - const controller = new AbortController(); - const abort = (): void => controller.abort(); - req.once('aborted', abort); - const abortClosedResponse = (): void => { - if (!res.writableEnded) abort(); - }; - res.once('close', abortClosedResponse); - try { - outcome.dispatchPending = true; - const dispatchStartedAt = performance.now(); - const settlement = await options.store.dispatchWorkspaceTool({ - workerId: selection.workerId, - tenantId: principal.tenantId, - requireTenantBinding: + const controller = new AbortController(); + const abort = (): void => controller.abort(); + req.once('aborted', abort); + const abortClosedResponse = (): void => { + if (!res.writableEnded) abort(); + }; + res.once('close', abortClosedResponse); + try { + outcome.dispatchPending = true; + const dispatchStartedAt = performance.now(); + const settlement = await options.store.dispatchWorkspaceTool({ + workerId: selection.workerId, + tenantId: principal.tenantId, + requireTenantBinding: selection.explicit && (options.dynamicWorkers || selection.workerId !== options.configuredWorkerId), - request, - deadlineAtMs: Date.now() + queueBudgetMs, - executionTimeoutMs: executionBudgetMs, - signal: controller.signal, - }).finally(() => { - outcome.dispatchDurationMs = Math.round(performance.now() - dispatchStartedAt); - }); - if (settlement.status === 'rejected') { - outcome.errorCode = settlement.errorCode ?? 'WORKSPACE_TOOL_REJECTED'; - let status = 422; - if ( - settlement.errorCode === 'SEARCH_TIMEOUT' || + request, + deadlineAtMs: Date.now() + queueBudgetMs, + executionTimeoutMs: executionBudgetMs, + signal: controller.signal, + }).finally(() => { + outcome.dispatchDurationMs = Math.round(performance.now() - dispatchStartedAt); + }); + if (settlement.status === 'rejected') { + outcome.errorCode = settlement.errorCode ?? 'WORKSPACE_TOOL_REJECTED'; + let status = 422; + if ( + settlement.errorCode === 'SEARCH_TIMEOUT' || settlement.errorCode === 'LIST_TIMEOUT' || settlement.errorCode === 'COMMAND_TIMEOUT' - ) { - status = 504; - } - if ( - settlement.errorCode === 'SEARCH_UNAVAILABLE' || + ) { + status = 504; + } + if ( + settlement.errorCode === 'SEARCH_UNAVAILABLE' || settlement.errorCode === 'LIST_UNAVAILABLE' || settlement.errorCode === 'COMMAND_UNAVAILABLE' - ) { - status = 503; - } - if (settlement.errorCode === 'WRITE_DISABLED') status = 403; - if (settlement.errorCode === 'COMMAND_DISABLED') status = 403; - if (settlement.errorCode === 'WRITE_LIMIT_EXCEEDED') status = 413; - if (settlement.errorCode === 'WRITE_UNAVAILABLE') status = 503; - if (settlement.errorCode === 'EDIT_CONFLICT') status = 409; - res.status(status).json({ - error: settlement.error, - code: settlement.errorCode ?? 'WORKSPACE_TOOL_REJECTED', - }); - return; - } - res.status(200).json(settlement.result); - } catch (error) { - if (error instanceof BridgeStoreError) { - outcome.errorCode = error.code; - if (error.code === 'WORKSPACE_QUEUE_TIMEOUT') res.setHeader('Retry-After', '1'); - res.status(bridgeStoreStatus(error)).json({ - error: error.message, - code: error.code, - }); - return; + ) { + status = 503; } - outcome.errorCode = 'INTERNAL_ERROR'; - throw error; - } finally { - outcome.dispatchPending = false; - outcome.flush(); - req.removeListener('aborted', abort); - res.removeListener('close', abortClosedResponse); + if (settlement.errorCode === 'WRITE_DISABLED') status = 403; + if (settlement.errorCode === 'COMMAND_DISABLED') status = 403; + if (settlement.errorCode === 'WRITE_LIMIT_EXCEEDED') status = 413; + if (settlement.errorCode === 'WRITE_UNAVAILABLE') status = 503; + if (settlement.errorCode === 'EDIT_CONFLICT') status = 409; + res.status(status).json({ + error: settlement.error, + code: settlement.errorCode ?? 'WORKSPACE_TOOL_REJECTED', + }); + return; } - }), - ); + res.status(200).json(settlement.result); + } catch (error) { + if (error instanceof BridgeStoreError) { + outcome.errorCode = error.code; + if (error.code === 'WORKSPACE_QUEUE_TIMEOUT') res.setHeader('Retry-After', '1'); + res.status(bridgeStoreStatus(error)).json({ + error: error.message, + code: error.code, + }); + return; + } + outcome.errorCode = 'INTERNAL_ERROR'; + throw error; + } finally { + outcome.dispatchPending = false; + outcome.flush(); + req.removeListener('aborted', abort); + res.removeListener('close', abortClosedResponse); + } + }); + router.post('/workspace-tools/execute', execute(false)); + router.post('/workspace-tools/requests', execute(true)); router.get('/workspace-tools/capabilities', asyncRoute(async (req, res) => { if (!getPrincipalOrReject(req, res)) return; From a94534617e5bac2e90a9090aa379134c05b2c6a8 Mon Sep 17 00:00:00 2001 From: Lia Date: Sat, 3 Oct 2026 04:12:55 +0000 Subject: [PATCH 3/5] =?UTF-8?q?=F0=9F=A7=BE=20fix:=20Fence=20Durable=20Cap?= =?UTF-8?q?acity=20Acquisition=20Before=20Dispatch?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/remote-bridge/README.md | 5 +- service/src/bridge/admission.ts | 17 +++ service/src/bridge/slots.ts | 10 +- service/src/bridge/store.ts | 43 +++--- service/src/workspace-tools/requests.test.ts | 137 +++++++++++++++++-- 5 files changed, 182 insertions(+), 30 deletions(-) diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index e5c4d177..19be5b4d 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -421,7 +421,10 @@ claim that work. Keep backend/profile, bridge auth mode, worker selection, slot ceiling and command-timeout policy identical on replicas behind one endpoint. Drain durable requests before changing those settings. FIFO admission position survives API restarts and temporary worker registration/readiness loss until the -queue deadline. Assignment enqueue and the durable `admitted` transition are atomic. +queue deadline. Capacity acquisition and assignment enqueue atomically check the coordinator +claim, queued state, cancellation and admission deadline. Delayed callbacks +cannot reserve capacity after their claim ends. Assignment enqueue also checks +capacity ownership and atomically commits the durable `admitted` transition. Recovery resumes only queued work and observes admitted assignments without creating another command. Worker identity/incarnation changes fail pending work rather than transferring it to another machine. Execution gets its full budget diff --git a/service/src/bridge/admission.ts b/service/src/bridge/admission.ts index db375507..18ab1f49 100644 --- a/service/src/bridge/admission.ts +++ b/service/src/bridge/admission.ts @@ -1,5 +1,22 @@ import type Redis from 'ioredis'; +/** Redis time and claim ownership gate every durable admission mutation. */ +export function durableAdmissionFence( + requestKey: string, + claimKey: string, + token: string, + rejected: number, +): string[] { + return [ + `if redis.call('GET', ${claimKey}) ~= ${token} then return ${rejected} end`, + `if redis.call('HGET', ${requestKey}, 'state') ~= 'queued' then return ${rejected} end`, + `if redis.call('HGET', ${requestKey}, 'cancelRequested') == '1' then return ${rejected} end`, + 'local admissionTime = redis.call(\'TIME\')', + 'local admissionNowMs = tonumber(admissionTime[1]) * 1000.0 + math.floor(tonumber(admissionTime[2]) / 1000)', + `if tonumber(redis.call('HGET', ${requestKey}, 'queueDeadlineAtMs')) <= admissionNowMs then return ${rejected} end`, + ]; +} + /** Bounded FIFO admission shared by API replicas. Entries expire after caller deadlines. */ export class BridgeAdmissionQueue { constructor( diff --git a/service/src/bridge/slots.ts b/service/src/bridge/slots.ts index e88c2b22..f748054b 100644 --- a/service/src/bridge/slots.ts +++ b/service/src/bridge/slots.ts @@ -1,4 +1,5 @@ import type Redis from 'ioredis'; +import { durableAdmissionFence } from './admission'; /** Hard bound keeps every atomic scheduling scan constant-sized. */ export const MAX_WORKSPACE_LEASE_SLOTS = 8; @@ -38,6 +39,7 @@ export class BridgeWorkspaceSlots { capacity: number; expiresAtMs: number; refreshOwned?: boolean; + guard?: { key: string; claimKey: string; token: string }; }): Promise { if ( !Number.isSafeInteger(args.capacity) || @@ -49,9 +51,12 @@ export class BridgeWorkspaceSlots { ) { throw new Error('Invalid workspace slot reservation'); } + const keys = this.keys(args.workerId); + if (args.guard != null) keys.push(args.guard.key, args.guard.claimKey); const result = Number( await this.redis.eval( [ + ...(args.guard == null ? [] : durableAdmissionFence('KEYS[9]', 'KEYS[10]', 'ARGV[8]', -3)), "if redis.call('GET', KEYS[4]) ~= ARGV[1] then return -2 end", "if (redis.call('GET', KEYS[8]) or '1') ~= ARGV[4] then return -2 end", // Parent of a linked-worktree lane key; nil for a checkout key. @@ -134,8 +139,8 @@ export class BridgeWorkspaceSlots { "redis.call('SET', KEYS[3], ARGV[1], 'PXAT', latest)", 'return free', ].join('\n'), - 8, - ...this.keys(args.workerId), + keys.length, + ...keys, args.incarnationId, args.assignmentId, args.workspaceId, @@ -143,6 +148,7 @@ export class BridgeWorkspaceSlots { args.expiresAtMs, Date.now(), args.refreshOwned === true ? '1' : '0', + args.guard?.token ?? '', ), ); if (result === -2) diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index ea173e74..07f1ecb7 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -21,7 +21,7 @@ import { workspaceIsolationParent, } from '../../../packages/code/src/protocol'; import type { BridgeWorkerBinding } from './pairing'; -import { BridgeAdmissionQueue } from './admission'; +import { BridgeAdmissionQueue, durableAdmissionFence } from './admission'; import { BridgeWorkspaceSlots } from './slots'; import type { StoredWorkspaceRequest } from '../workspace-tools/requests'; @@ -843,10 +843,10 @@ export class RedisBridgeStore { slot = await new BridgeWorkspaceSlots(this.redis).reserve({ workerId: record.workerId, incarnationId: registration.incarnationId, assignmentId: record.id, workspaceId: workspace, capacity, - expiresAtMs: Date.now() + ttlSeconds * 1000, refreshOwned: true, + expiresAtMs: Date.now() + ttlSeconds * 1000, refreshOwned: true, guard, }); if (slot === undefined) return; - } else if (!await this.acquireLock(record.workerId, record.id, registration.incarnationId, ttlSeconds, true)) return; + } else if (!await this.acquireLock(record.workerId, record.id, registration.incarnationId, ttlSeconds, true, guard)) return; if (Date.now() >= record.queueDeadlineAtMs) return { state: 'failed', error: { code: 'WORKSPACE_QUEUE_TIMEOUT', message: 'Workspace admission deadline expired. The operation was not started.' }, }; @@ -2382,10 +2382,13 @@ export class RedisBridgeStore { const script = [ 'local keyCount = tonumber(ARGV[10])', ...(guard == null ? [] : [ - 'if redis.call(\'GET\', KEYS[keyCount + 2]) ~= ARGV[11] then return -2 end', - 'if redis.call(\'HGET\', KEYS[keyCount + 1], \'state\') ~= \'queued\' then return -2 end', - 'if redis.call(\'HGET\', KEYS[keyCount + 1], \'cancelRequested\') == \'1\' then return -2 end', - 'if tonumber(redis.call(\'HGET\', KEYS[keyCount + 1], \'queueDeadlineAtMs\')) <= tonumber(ARGV[12]) then return -2 end', + ...durableAdmissionFence('KEYS[keyCount + 1]', 'KEYS[keyCount + 2]', 'ARGV[11]', -2), + ...(assignment.workspaceLeaseSlot === undefined + ? ['if redis.call(\'GET\', KEYS[keyCount + 3]) ~= ARGV[4] or redis.call(\'GET\', KEYS[4]) ~= ARGV[1] then return -2 end'] + : [ + `if redis.call('HGET', KEYS[keyCount + 3], 'a:${assignment.workspaceLeaseSlot}') ~= ARGV[4] then return -2 end`, + `if redis.call('HGET', KEYS[keyCount + 3], 'i:${assignment.workspaceLeaseSlot}') ~= ARGV[1] then return -2 end`, + ]), ]), "if redis.call('GET', KEYS[1]) ~= ARGV[1] then return 0 end", 'if ARGV[7] ~= "" and redis.call(\'GET\', KEYS[6]) ~= ARGV[7] then return 0 end', @@ -2408,13 +2411,13 @@ export class RedisBridgeStore { "if ARGV[9] == 'lane' then redis.call('SADD', KEYS[11], KEYS[7]) end", 'if keyCount >= 9 then', " local epoch = redis.call('GET', KEYS[9])", - ' if type(epoch) ~= \'string\' then epoch = \'0\'; redis.call(\'SET\', KEYS[9], epoch, \'EX\', ARGV[13]) end', - ' if redis.call(\'PTTL\', KEYS[9]) < tonumber(ARGV[13]) * 1000 then redis.call(\'EXPIRE\', KEYS[9], ARGV[13]) end', + ' if type(epoch) ~= \'string\' then epoch = \'0\'; redis.call(\'SET\', KEYS[9], epoch, \'EX\', ARGV[12]) end', + ' if redis.call(\'PTTL\', KEYS[9]) < tonumber(ARGV[12]) * 1000 then redis.call(\'EXPIRE\', KEYS[9], ARGV[12]) end', " redis.call('HSET', KEYS[8], 'metadata', ARGV[8], 'epoch', epoch)", - ' redis.call(\'EXPIRE\', KEYS[8], ARGV[13])', + ' redis.call(\'EXPIRE\', KEYS[8], ARGV[12])', 'end', ...(guard == null ? [] : [ - 'redis.call(\'HSET\', KEYS[keyCount + 1], \'state\', \'admitted\', \'assignment\', ARGV[2], \'admittedAtMs\', ARGV[12])', + 'redis.call(\'HSET\', KEYS[keyCount + 1], \'state\', \'admitted\', \'assignment\', ARGV[2], \'admittedAtMs\', admissionNowMs)', ]), 'return 1', ].join('\n'); @@ -2474,7 +2477,9 @@ export class RedisBridgeStore { } } const keyCount = keys.length; - if (guard != null) keys.push(guard.key, guard.claimKey); + if (guard != null) keys.push(guard.key, guard.claimKey, + assignment.workspaceLeaseSlot === undefined ? lockKey(assignment.workerId) + : `${workerKey(assignment.workerId)}:workspace-slots`); const result = await this.redis.eval( script, keys.length, @@ -2490,7 +2495,6 @@ export class RedisBridgeStore { fenceScope, keyCount, guard?.token ?? '', - Date.now(), guard == null ? ttlSeconds : 86400, ); if (Number(result) === -1) { @@ -2508,23 +2512,30 @@ export class RedisBridgeStore { incarnationId: string, ttlSeconds: number, reuseOwned = false, + guard?: DurableWorkspaceGuard, ): Promise { const script = [ + ...(guard == null ? [] : [ + ...durableAdmissionFence('KEYS[3]', 'KEYS[4]', 'ARGV[5]', 0), + 'if redis.call(\'GET\', KEYS[5]) ~= ARGV[2] then return 0 end', + ]), 'if ARGV[4] == \'1\' and redis.call(\'GET\', KEYS[1]) == ARGV[1] and redis.call(\'GET\', KEYS[2]) == ARGV[2] then redis.call(\'PEXPIRE\', KEYS[1], ARGV[3]); redis.call(\'PEXPIRE\', KEYS[2], ARGV[3]); return 1 end', 'if redis.call(\'EXISTS\', KEYS[1]) == 1 then return 0 end', 'redis.call(\'SET\', KEYS[1], ARGV[1], \"PX\", ARGV[3])', 'redis.call(\'SET\', KEYS[2], ARGV[2], \"PX\", ARGV[3])', 'return 1', ].join('\n'); + const keys = [lockKey(workerId), lockIncarnationKey(workerId)]; + if (guard != null) keys.push(guard.key, guard.claimKey, workerIncarnationKey(workerId)); const result = await this.redis.eval( script, - 2, - lockKey(workerId), - lockIncarnationKey(workerId), + keys.length, + ...keys, assignmentId, incarnationId, String(ttlSeconds * 1000), reuseOwned ? '1' : '0', + guard?.token ?? '', ); return Number(result) === 1; } diff --git a/service/src/workspace-tools/requests.test.ts b/service/src/workspace-tools/requests.test.ts index 850a266d..b4329764 100644 --- a/service/src/workspace-tools/requests.test.ts +++ b/service/src/workspace-tools/requests.test.ts @@ -3,6 +3,7 @@ import { createHash } from 'node:crypto'; import RedisMock from 'ioredis-mock'; import Redis from 'ioredis'; import { RedisBridgeStore } from '../bridge/store'; +import { BridgeAdmissionQueue } from '../bridge/admission'; import type { CodeBridgeAssignment } from '../bridge/store'; import { RedisWorkspaceRequests } from './requests'; import { workspaceRequestCoordinationScope } from './coordination'; @@ -82,16 +83,20 @@ test('submission is idempotent across replicas, scoped by principal, and conflic test('restart before dispatch preserves the queue and starts a full execution budget after sixty seconds waiting', async () => { const requests = await setup(); await submit(requests); - const clock = spyOn(Date, 'now').mockReturnValue(Date.now() + 60_000); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + await mutateRecord(record => { + record.createdAtMs = Number(record.createdAtMs) - 60_000; + record.queueDeadlineAtMs = Number(record.queueDeadlineAtMs) - 60_000; + }); + const record = JSON.parse((await redis.hget(key, 'record'))!); + await redis.hset(key, 'queueDeadlineAtMs', record.queueDeadlineAtMs); const restarted = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis)); - let assignment: CodeBridgeAssignment | undefined; - try { - await tick(restarted); - assignment = await bridge.lease(workerId, incarnationId, 0); - expect(assignment).toBeDefined(); - expect(assignment!.remainingMs).toBeGreaterThan(29_000); - expect((await restarted.get(owner, requestId))?.queueWaitMs).toBeGreaterThanOrEqual(60_000); - } finally { clock.mockRestore(); } + await tick(restarted); + const assignment = await bridge.lease(workerId, incarnationId, 0); + expect(assignment).toBeDefined(); + expect(assignment!.remainingMs).toBeGreaterThan(29_000); + // Mock Redis TIME has subsecond clock skew. Execution still receives its full budget. + expect((await restarted.get(owner, requestId))?.queueWaitMs).toBeGreaterThan(59_000); await settle(assignment!); await tick(restarted); expect(await restarted.get(owner, requestId)).toMatchObject({ state: 'completed', result: { content: 'hello' } }); @@ -205,14 +210,16 @@ test.each([1, 2])('restart after capacity reservation refreshes its own reservat expect(await redis.pttl(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeGreaterThan(30_000); }); -test('a stale coordinator cannot dispatch or overwrite a newer claim', async () => { - const requests = await setup(); await submit(requests); +test.each([1, 2])('a stale coordinator cannot acquire capacity or dispatch (slots=%s)', async slots => { + const requests = await setup(slots); await submit(requests); const key = (await redis.zrange(activeKey, 0, -1))[0]; const record = JSON.parse((await redis.hget(key, 'record'))!); await redis.set(`${key}:claim`, 'new-owner', 'PX', 10000); await bridge.advanceDurableWorkspaceTool(record, { key, claimKey: `${key}:claim`, token: 'old-owner' }); expect(await bridge.lease(workerId, incarnationId, 0)).toBeUndefined(); expect((await requests.get(owner, requestId))?.state).toBe('queued'); + expect(await redis.get(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeNull(); + expect(await redis.hlen(`codeapi:bridge:v1:worker:${workerId}:workspace-slots`)).toBe(0); }); test('a queued deadline expires as definitely unstarted work', async () => { @@ -383,3 +390,111 @@ test('a different policy cannot claim or fail work accepted by a two-slot bridge const assignment = await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0); expect(assignment).toBeDefined(); }); + +test('a late FIFO response cannot relock a cancelled request after the coordinator timed out', async () => { + const requests = await setup(); await submit(requests); + let release!: () => void; + let entered!: () => void; + let finished!: () => void; + const gate = new Promise(resolve => { release = resolve; }); + const started = new Promise(resolve => { entered = resolve; }); + const drained = new Promise(resolve => { finished = resolve; }); + const originalHead = BridgeAdmissionQueue.prototype.isHead; + const originalAdvance = bridge.advanceDurableWorkspaceTool.bind(bridge); + const head = spyOn(BridgeAdmissionQueue.prototype, 'isHead').mockImplementation(async (worker, id) => { + const result = await originalHead.call(new BridgeAdmissionQueue(redis), worker, id); + entered(); await gate; return result; + }); + const advance = spyOn(bridge, 'advanceDurableWorkspaceTool').mockImplementation(async (...args) => { + try { return await originalAdvance(...args); } finally { finished(); } + }); + const pending = tick(requests); + void pending.catch(() => undefined); + try { + await started; + await expect(pending).rejects.toThrow('Durable workspace transition timed out'); + head.mockRestore(); advance.mockRestore(); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const claimDeadline = Date.now() + 1000; + while (await redis.get(`${key}:claim`) != null) { + if (Date.now() >= claimDeadline) throw new Error('Timed-out coordinator did not release its claim'); + await new Promise(resolve => setTimeout(resolve, 5)); + } + await requests.cancel(owner, requestId); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect((await requests.get(owner, requestId))?.state).toBe('cancelled'); + expect(await redis.zcard(activeKey)).toBe(0); + release(); await drained; + expect(await redis.get(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeNull(); + await submit(requests, 'request-00000000002'); await tick(requests); + expect(await bridge.lease(workerId, incarnationId, 0)).toBeDefined(); + } finally { release(); head.mockRestore(); advance.mockRestore(); await drained; } +}, 15_000); + +test.each([1, 2])('Redis capacity writes delayed past cancellation cannot strand a reservation (slots=%s)', async slots => { + const requests = await setup(slots); await submit(requests); + let release!: () => void; + let entered!: () => void; + const gate = new Promise(resolve => { release = resolve; }); + const started = new Promise(resolve => { entered = resolve; }); + const original = redis.eval.bind(redis); + let held = false; + const delay = spyOn(redis, 'eval').mockImplementation((async (...args: Parameters) => { + const script = String(args[0]); + if (!held && script.includes('local admissionTime') && !script.includes('local keyCount')) { + held = true; entered(); await gate; + } + return original(...args); + }) as Redis['eval']); + const pending = tick(requests); + void pending.catch(() => undefined); + try { + await started; + const key = (await redis.zrange(activeKey, 0, -1))[0]; + // Expire this claim while its Redis command is still undelivered. + await redis.del(`${key}:claim`); + await requests.cancel(owner, requestId); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect((await requests.get(owner, requestId))?.state).toBe('cancelled'); + release(); await pending; + expect(await redis.get(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeNull(); + expect(await redis.hlen(`codeapi:bridge:v1:worker:${workerId}:workspace-slots`)).toBe(0); + expect(await redis.zcard(activeKey)).toBe(0); + } finally { release(); delay.mockRestore(); await pending.catch(() => undefined); } +}); + +test.each([1, 2])('a durable enqueue refuses lost capacity ownership (slots=%s)', async slots => { + const requests = await setup(slots); await submit(requests); + const original = redis.eval.bind(redis); + const loss = spyOn(redis, 'eval').mockImplementation((async (...args: Parameters) => { + if (String(args[0]).includes('local keyCount')) { + await redis.del(`codeapi:bridge:v1:worker:${workerId}:lock`, + `codeapi:bridge:v1:worker:${workerId}:workspace-slots`); + } + return original(...args); + }) as Redis['eval']); + try { await tick(requests); } finally { loss.mockRestore(); } + expect((await requests.get(owner, requestId))?.state).toBe('queued'); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeUndefined(); + await tick(requests); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeDefined(); +}); + +test.each([1, 2])('capacity acquisition rejects terminal, cancelled, expired, and unclaimed requests (slots=%s)', async slots => { + for (const denial of ['terminal', 'cancelled', 'expired', 'unclaimed']) { + await redis.flushall(); + const requests = await setup(slots); await submit(requests); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const record = JSON.parse((await redis.hget(key, 'record'))!); + const token = 'current-coordinator'; + await redis.set(`${key}:claim`, token, 'PX', 10000); + if (denial === 'terminal') await redis.hset(key, 'state', 'cancelled'); + if (denial === 'cancelled') await redis.hset(key, 'cancelRequested', '1'); + if (denial === 'expired') await redis.hset(key, 'queueDeadlineAtMs', Date.now() - 2000); + if (denial === 'unclaimed') await redis.del(`${key}:claim`); + await bridge.advanceDurableWorkspaceTool(record, { key, claimKey: `${key}:claim`, token }); + expect(await redis.get(`codeapi:bridge:v1:worker:${workerId}:lock`)).toBeNull(); + expect(await redis.hlen(`codeapi:bridge:v1:worker:${workerId}:workspace-slots`)).toBe(0); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeUndefined(); + } +}); From 0a0fe93c7d4ae1ae174b56863418577035ad7fb2 Mon Sep 17 00:00:00 2001 From: Lia Date: Sat, 3 Oct 2026 04:21:27 +0000 Subject: [PATCH 4/5] =?UTF-8?q?=F0=9F=A7=BE=20fix:=20Retain=20Durable=20Qu?= =?UTF-8?q?arantine=20Outcomes=20Across=20Restarts?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- service/src/bridge/store.ts | 16 +++++--- service/src/workspace-tools/requests.test.ts | 43 ++++++++++++++++++++ 2 files changed, 54 insertions(+), 5 deletions(-) diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index 07f1ecb7..f93a06dc 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -121,6 +121,7 @@ type AssignmentOwnership = Pick< | 'workerIdentityId' | 'expiresAt' | 'runtimeSessionId' + | 'durableRequestKey' >; const NATIVE_WORKSPACE_FENCE_PREFIX = 'native-workspace:'; @@ -398,6 +399,12 @@ function assignmentTtlSeconds(deadlineAtMs: number): number { return Math.max(1, Math.ceil((deadlineAtMs - Date.now()) / 1000) + 30); } +function assignmentOutcomeTtlSeconds(assignment: AssignmentOwnership): number { + return assignment.durableRequestKey == null + ? assignmentTtlSeconds(Date.parse(assignment.expiresAt)) + : 24 * 60 * 60; +} + async function delay(ms: number, signal?: AbortSignal): Promise { if (signal?.aborted === true) return; await new Promise((resolve) => { @@ -1844,9 +1851,7 @@ export class RedisBridgeStore { 'Bridge assignment has expired', ); } - const ttlSeconds = assignment.durableRequestKey == null - ? assignmentTtlSeconds(Date.parse(assignment.expiresAt)) - : 86400; + const ttlSeconds = assignmentOutcomeTtlSeconds(assignment); const settlementKeys = [ assignmentKey(assignmentId), settlementKey(assignmentId), @@ -2125,7 +2130,7 @@ export class RedisBridgeStore { raw.epoch, assignmentId, JSON.stringify(settlement), - assignmentTtlSeconds(Date.parse(receipt.expiresAt)), + assignmentOutcomeTtlSeconds(receipt), quarantine ? '1' : '0', ), signal, @@ -2451,6 +2456,7 @@ export class RedisBridgeStore { leaseTokenHash: assignment.leaseTokenHash, workerIdentityId: assignment.workerIdentityId, expiresAt: assignment.expiresAt, + durableRequestKey: assignment.durableRequestKey, }; let fenceScope: '' | 'lane' | 'checkout' = ''; if (assignment.workspaceLeaseSlot !== undefined) { @@ -2495,7 +2501,7 @@ export class RedisBridgeStore { fenceScope, keyCount, guard?.token ?? '', - guard == null ? ttlSeconds : 86400, + guard == null ? ttlSeconds : assignmentOutcomeTtlSeconds(assignment), ); if (Number(result) === -1) { throw new BridgeStoreError( diff --git a/service/src/workspace-tools/requests.test.ts b/service/src/workspace-tools/requests.test.ts index b4329764..ff13cd4b 100644 --- a/service/src/workspace-tools/requests.test.ts +++ b/service/src/workspace-tools/requests.test.ts @@ -498,3 +498,46 @@ test.each([1, 2])('capacity acquisition rejects terminal, cancelled, expired, an expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeUndefined(); } }); + +test.each([false, true])('quarantine retains the winning durable outcome across two-minute recovery (settled=%s)', async settled => { + const requests = await setup(2); + await requests.submit({ owner, requestId, workerId, requireTenantBinding: false, + request: { protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, + queueWaitMs: 300_000, executionTimeoutMs: 30_000, + }); + await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0))!; + await bridge.acknowledgeLease(workerId, incarnationId, assignment.assignmentId, assignment.generation, assignment.leaseToken); + const firstError = 'Command stopped before local cleanup'; + const quarantineError = 'Workspace cleanup failed before settlement'; + const envelope = { protocolVersion: 1 as const, incarnationId, + generation: assignment.generation, leaseToken: assignment.leaseToken, + status: 'rejected' as const, + }; + if (settled) await bridge.settle(workerId, assignment.assignmentId, { ...envelope, error: firstError }); + await bridge.settle(workerId, assignment.assignmentId, { ...envelope, error: quarantineError }, undefined, undefined, true); + const resultKey = `codeapi:bridge:v1:assignment:${assignment.assignmentId}:settlement`; + const receiptKey = `codeapi:bridge:v1:assignment:${assignment.assignmentId}:workspace-fence-owner`; + const receipt = JSON.parse((await redis.hget(receiptKey, 'metadata'))!); + expect(receipt.durableRequestKey).toBe((await redis.zrange(activeKey, 0, -1))[0]); + const remaining = await redis.pttl(resultKey); + expect(remaining).toBeGreaterThan(86_000_000); + expect(await redis.get(`codeapi:bridge:v1:assignment:${assignment.assignmentId}`)).toBeNull(); + // Advance retained storage and the coordinator clock without a two-minute sleep. + await redis.pexpire(resultKey, Math.max(0, remaining - 120_000)); + const clock = spyOn(Date, 'now').mockReturnValue(Date.now() + 120_000); + const restarted = new RedisWorkspaceRequests(redis, new RedisBridgeStore(redis, 600, 1000, 2)); + try { + await tick(restarted); + expect(await restarted.get(owner, requestId)).toMatchObject({ + state: 'failed', error: { code: 'WORKSPACE_TOOL_REJECTED', message: settled ? firstError : quarantineError }, + }); + expect(await redis.pttl(resultKey)).toBeGreaterThan(85_000_000); + expect(await redis.zcard(activeKey)).toBe(0); + await submit(restarted, 'request-00000000002'); await tick(restarted); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0)).toBeUndefined(); + expect(await restarted.get(owner, 'request-00000000002')).toMatchObject({ + state: 'failed', error: { code: 'WORKSPACE_QUARANTINED' }, + }); + } finally { clock.mockRestore(); } +}); From 0b5ea20aa4bfa3e081de25a86a3984ff4f7048af Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 4 Oct 2026 01:41:40 +0000 Subject: [PATCH 5/5] =?UTF-8?q?=F0=9F=A7=BE=20fix:=20Preserve=20Durable=20?= =?UTF-8?q?Admission=20and=20Read=20Cancellation?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/remote-bridge/README.md | 4 +- service/src/bridge/store.ts | 15 ++- service/src/workspace-tools/requests.test.ts | 130 ++++++++++++++++++- service/src/workspace-tools/requests.ts | 4 +- 4 files changed, 142 insertions(+), 11 deletions(-) diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index 19be5b4d..dc7cb98b 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -404,7 +404,9 @@ Disconnecting does not cancel accepted work. `admitted`, `completed`, `failed`, and `cancelled`. Lookup is tenant/user scoped; unknown or another principal's IDs return 404. Results are non-consuming reads. - Cancel with `DELETE /v1/workspace-tools/requests/:requestId`. Cancellation is a - request, not proof of termination. A committed result can win the race. + request, not proof of termination. All operations notify the worker. A committed + result can win the race. Without a confirmed settlement, read-only calls retain + `cancelled`; mutations retain an unknown-outcome failure and their safety fence. `ASSIGNMENT_EXPIRED` means execution may have occurred. Never replay it. - Status includes worker/workspace/lane metadata, queue position when available, measured queue wait, and the execution deadline after admission. Queue position diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index f93a06dc..542c5903 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -166,6 +166,10 @@ export interface BridgeWorkerStatus { capabilities?: BridgeWorkerRegistration['capabilities']; } +function isWorkspaceMutation(request: WorkspaceToolRequest): boolean { + return request.operation === 'write_file' || request.operation === 'edit_file' || request.operation === 'execute_command'; +} + function supportsWorkspaceTool( registration: RegisteredBridgeWorker, request: WorkspaceToolRequest, @@ -895,6 +899,7 @@ export class RedisBridgeStore { } } catch (error) { if (!(error instanceof BridgeStoreError) || error.code !== 'ASSIGNMENT_EXPIRED') throw error; + if (record.cancelRequested === true && !isWorkspaceMutation(record.request)) return { state: 'cancelled' }; return { state: 'failed', error: { code: error.code, message: 'Assignment ended without a confirmed result. Do not replay this operation.' } }; } if (settlement.status === 'fulfilled') { @@ -2275,14 +2280,12 @@ export class RedisBridgeStore { const cancelledMutation = signal.aborted && (assignment.executionKind === 'workspace_programmatic' || - (workspaceRequest != null && - (workspaceRequest.operation === 'write_file' || - workspaceRequest.operation === 'edit_file' || - workspaceRequest.operation === 'execute_command'))); - if (cancelledMutation) { + (workspaceRequest != null && isWorkspaceMutation(workspaceRequest))); + const cancelledDurableWorkspace = signal.aborted && assignment.durableRequestKey != null && workspaceRequest != null; + if (cancelledMutation || cancelledDurableWorkspace) { try { // Keep the acknowledged assignment available long enough for the - // worker to terminate its process tree and commit a clean rejection. + // worker to stop its operation and commit a clean rejection. // Closing it first makes that rejection impossible to acknowledge and // leaves the worker's durable mutation guard armed. await this.cancel(assignment.assignmentId, assignment); diff --git a/service/src/workspace-tools/requests.test.ts b/service/src/workspace-tools/requests.test.ts index ff13cd4b..a0f7d99f 100644 --- a/service/src/workspace-tools/requests.test.ts +++ b/service/src/workspace-tools/requests.test.ts @@ -5,6 +5,7 @@ import Redis from 'ioredis'; import { RedisBridgeStore } from '../bridge/store'; import { BridgeAdmissionQueue } from '../bridge/admission'; import type { CodeBridgeAssignment } from '../bridge/store'; +import type { WorkspaceToolRequest } from '../../../packages/code/src/protocol'; import { RedisWorkspaceRequests } from './requests'; import { workspaceRequestCoordinationScope } from './coordination'; import type { WorkspaceRequestCoordinationPolicy } from './coordination'; @@ -29,7 +30,7 @@ async function setup(slots = 1): Promise { capabilities: { sandboxProfile: 'native-srt', statefulWorkspace: false, runtimes: ['bash'], ...(slots > 1 ? { workspaceLeaseSlots: slots, requiresReadyConfirmation: true } : {}), - workspaceTools: { protocolVersion: 1, operations: ['read_file', 'execute_command'], + workspaceTools: { protocolVersion: 1, operations: ['read_file', 'execute_command', 'list_files', 'search_text', 'preview_edit'], workspaces: [{ id: 'primary' }, { id: 'independent' }] }, }, }); @@ -541,3 +542,130 @@ test.each([false, true])('quarantine retains the winning durable outcome across }); } finally { clock.mockRestore(); } }); + +test('queued durable requests retain FIFO workspace metadata across same-incarnation slot changes', async () => { + const requests = await setup(); + await submit(requests); + await submit(requests, 'request-00000000002'); + await submit(requests, 'request-00000000003', 'independent'); + await submit(requests); + await setup(2); + await tick(requests); await tick(requests); + const first = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0))!; + const independent = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 1))!; + expect(first).toBeDefined(); expect(independent).toBeDefined(); + expect(first.request).toMatchObject({ workspaceId: 'primary' }); + expect(independent.request).toMatchObject({ workspaceId: 'independent' }); + expect(await requests.get(owner, 'request-00000000002')).toMatchObject({ state: 'queued', queuePosition: 2 }); + await settle(first); await settle(independent); await tick(requests); + await bridge.confirmWorkspaceCleanup(workerId, first.assignmentId, { + protocolVersion: 1, incarnationId, generation: first.generation, leaseToken: first.leaseToken, + status: 'rejected', error: 'local cleanup confirmed', + }); + await tick(requests); + const next = (await bridge.lease(workerId, incarnationId, 0, undefined, undefined, 0))!; + expect(next).toBeDefined(); + await settle(next); await tick(requests); + expect((await requests.get(owner, 'request-00000000002'))?.state).toBe('completed'); +}); + +const readOnlyRequests: WorkspaceToolRequest[] = [ + { protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, + { protocolVersion: 1, operation: 'list_files', workspaceId: 'primary' }, + { protocolVersion: 1, operation: 'search_text', workspaceId: 'primary', query: 'needle' }, + { protocolVersion: 1, operation: 'preview_edit', workspaceId: 'primary', path: 'README.md', oldText: 'old', newText: 'new' }, +]; +const readCancellationCases = readOnlyRequests.flatMap(request => [1, 2].flatMap(slots => + ['unleased', 'claimed', 'acknowledged'].map(phase => ({ request, slots, phase })), +)); + +test.each(readCancellationCases)('durable read cancellation notifies the worker before closing (%j)', async ({ request, slots, phase }) => { + const requests = await setup(slots); + await requests.submit({ owner, requestId, workerId, requireTenantBinding: false, + request, queueWaitMs: 300_000, executionTimeoutMs: 30_000, + }); + await tick(requests); + const slot = slots === 1 ? undefined : 0; + let assignment = phase === 'unleased' ? undefined : await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slot); + if (phase === 'acknowledged') await bridge.acknowledgeLease(workerId, incarnationId, assignment!.assignmentId, + assignment!.generation, assignment!.leaseToken); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const stored = JSON.parse((await redis.hget(key, 'assignment'))!); + await requests.cancel(owner, requestId); + const cancelling = tick(new RedisWorkspaceRequests(redis, bridge)); + const marker = `codeapi:bridge:v1:assignment:${stored.assignmentId}:cancelled`; + const limit = Date.now() + 1000; + try { + while (await redis.get(marker) !== '1') { + if (Date.now() >= limit) throw new Error('Read-only worker was never notified'); + await new Promise(resolve => setTimeout(resolve, 5)); + } + expect((await requests.get(owner, requestId))?.state).toBe('admitted'); + assignment ??= await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slot); + expect(assignment).toBeDefined(); + expect(await bridge.cancelled(workerId, incarnationId, assignment!.assignmentId)).toBe(true); + await settle(assignment!, false); + await cancelling; + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'cancelled', cancelRequested: true }); + expect(await redis.zcard(activeKey)).toBe(0); + } finally { await cancelling; } +}); + +test.each([1, 2])('a durable read cancellation remains cancelled without worker settlement (slots=%s)', async slots => { + const requests = await setup(slots); await submit(requests); await tick(requests); + const key = (await redis.zrange(activeKey, 0, -1))[0]; + const stored = JSON.parse((await redis.hget(key, 'assignment'))!); + await requests.cancel(owner, requestId); + await tick(new RedisWorkspaceRequests(redis, bridge)); + expect(await redis.get(`codeapi:bridge:v1:assignment:${stored.assignmentId}:cancelled`)).toBe('1'); + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'cancelled', cancelRequested: true }); + expect((await requests.get(owner, requestId))?.error).toBeUndefined(); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeUndefined(); + expect(await redis.zcard(activeKey)).toBe(0); + await submit(requests, 'request-00000000002'); await tick(requests); + expect(await bridge.lease(workerId, incarnationId, 0, undefined, undefined, slots === 1 ? undefined : 0)).toBeDefined(); +}, 10_000); + +test('fulfillment racing durable read cancellation keeps the winning result', async () => { + const requests = await setup(); await submit(requests); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await requests.cancel(owner, requestId); + const cancelling = tick(requests); + const marker = `codeapi:bridge:v1:assignment:${assignment.assignmentId}:cancelled`; + const limit = Date.now() + 1000; + try { + while (await redis.get(marker) !== '1') { + if (Date.now() >= limit) throw new Error('Read-only worker was never notified'); + await new Promise(resolve => setTimeout(resolve, 5)); + } + await settle(assignment); + await cancelling; + expect(await requests.get(owner, requestId)).toMatchObject({ state: 'completed', result: { content: 'hello' } }); + } finally { await cancelling; } +}); + +test('queued workspace metadata also permits a same-incarnation return to serial admission', async () => { + const requests = await setup(2); await submit(requests); await setup(); await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + expect(assignment).toBeDefined(); + expect(assignment.workspaceLeaseSlot).toBeUndefined(); + await settle(assignment); await tick(requests); + expect((await requests.get(owner, requestId))?.state).toBe('completed'); +}); + +test('unconfirmed durable mutation cancellation still reports an unknown outcome and retains its fence', async () => { + const requests = await setup(); + await requests.submit({ owner, requestId, workerId, requireTenantBinding: false, + request: { protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'sleep 30' }, + queueWaitMs: 300_000, executionTimeoutMs: 35_000, + }); + await tick(requests); + const assignment = (await bridge.lease(workerId, incarnationId, 0))!; + await bridge.acknowledgeLease(workerId, incarnationId, assignment.assignmentId, assignment.generation, assignment.leaseToken); + await requests.cancel(owner, requestId); await tick(requests); + expect(await requests.get(owner, requestId)).toMatchObject({ + state: 'failed', cancelRequested: true, error: { code: 'ASSIGNMENT_EXPIRED' }, + }); + const fence = `codeapi:bridge:v1:worker:${workerId}:workspace:${createHash('sha256').update('native-workspace:primary').digest('hex')}:quarantined`; + expect(await redis.get(fence)).toBe(assignment.assignmentId); +}, 10_000); diff --git a/service/src/workspace-tools/requests.ts b/service/src/workspace-tools/requests.ts index a88b9acb..2988c7cc 100644 --- a/service/src/workspace-tools/requests.ts +++ b/service/src/workspace-tools/requests.ts @@ -129,9 +129,7 @@ export class RedisWorkspaceRequests { const queue = new BridgeAdmissionQueue(this.redis); const accepted = await queue.submit({ workerId: record.workerId, id: record.id, deadlineAtMs: record.queueDeadlineAtMs, - workspaceId: (registration.capabilities.workspaceLeaseSlots ?? 1) > 1 - ? workspaceAdmissionId(record.request.workspaceId, record.request.workspaceInstanceId, record.request.worktree) - : undefined, + workspaceId: workspaceAdmissionId(record.request.workspaceId, record.request.workspaceInstanceId, record.request.worktree), key, activeKey: this.activeKey, fingerprint, record: JSON.stringify(record), retentionMs: RETENTION_MS, }); if (accepted === 'conflict') throw new WorkspaceRequestConflict('Request ID was already used for different work');