diff --git a/docs/byom-worker-admission.md b/docs/byom-worker-admission.md index 2fcd9f9f..0df8f17c 100644 --- a/docs/byom-worker-admission.md +++ b/docs/byom-worker-admission.md @@ -5,8 +5,13 @@ The limit is 32 admitted requests per worker, including the active request. When the limit is reached, the workspace endpoint returns HTTP 429 with `WORKER_QUEUE_FULL`. A different worker has an independent admission queue. -The Code API workspace HTTP endpoint allows up to the smaller of `JOB_TIMEOUT` -and five minutes for admission while the caller stays connected. After admission +The Code API workspace HTTP endpoint allows 30 seconds for admission when the +request has no `X-LibreChat-Workspace-Queue-Wait-Ms` header. A caller can +advertise a positive integer allowance in milliseconds through that header, +up to five minutes and any configured server queue ceiling. An invalid or +out-of-range value is rejected before dispatch. This queue budget is +independent of `JOB_TIMEOUT`; a shorter client or proxy deadline still ends +the wait. After admission and worker validation, a separate execution deadline starts. Commands receive their requested timeout (30 seconds by default, up to five minutes), capped by the operator's `JOB_TIMEOUT`, plus five seconds to settle the result. Other @@ -29,22 +34,27 @@ Existing workers still execute one assignment at a time. Parallel execution acro workspaces requires separate lease claims and isolated native sandbox contexts; this admission change does not advertise that capability. -Clients and reverse proxies must allow queue time plus execution/settlement time -and five seconds for HTTP delivery. With the default five-minute `JOB_TIMEOUT`, -that is at least 335 seconds for non-command tools, 340 seconds for default -commands, and 610 seconds for five-minute commands. With a smaller `JOB_TIMEOUT`, -use `min(JOB_TIMEOUT, 300s)` for the queue, plus `min(JOB_TIMEOUT, 30s)` for other +Clients and reverse proxies must allow the admitted queue budget plus +execution/settlement time and five seconds for HTTP delivery. With the default +five-minute `JOB_TIMEOUT` and no queue header, that is at least 65 seconds for +non-command tools, 70 seconds for default commands, and 340 seconds for +five-minute commands. At the maximum advertised five-minute queue allowance, +those totals become 335, 340, and 610 seconds respectively. With a smaller +`JOB_TIMEOUT`, use the advertised allowance (or 30 seconds without a header), +bounded by the server queue ceiling, plus `min(JOB_TIMEOUT, 30s)` for other operations or `min(JOB_TIMEOUT, requested command timeout) + 5s` for commands, -plus five seconds for delivery. +plus five seconds for delivery. The caller should advertise only the queue time +left after reserving execution, settlement, and delivery under its own HTTP +deadline; Code API does not receive that absolute deadline. -At the time of this change, LibreChat's `getWorkspaceToolTimeoutMs` still budgets -only 30 seconds for a single admission attempt (65/70/340 seconds in total). -Its `maxQueueWaitMs` is a retry horizon after a typed capacity rejection, **not** -a per-attempt HTTP timeout. Updating Code API alone therefore does not guarantee -the full wait. An earlier client, tool, or proxy timeout disconnects the request; -if work was already admitted, a mutation may have run and must not be blindly -retried. Match LibreChat's per-attempt timeout and each intermediary to the new -budget before relying on it. Existing workers do not need an update. +LibreChat's `maxQueueWaitMs` is a retry horizon after a typed capacity +rejection, **not** a per-attempt HTTP timeout. Without its opt-in +`maxRequestTimeoutMs`, LibreChat keeps a 30-second admission allowance per +attempt. Enabling a longer client budget requires LibreChat's header support on +every API replica and a timed canary through each intermediary; changing Code +API alone does not guarantee the full wait. An earlier client, tool, or proxy +timeout disconnects the request; if work was already admitted, a mutation may +have run and must not be blindly retried. Existing workers do not need an update. Focused regression coverage lives in `service/src/bridge/admission.test.ts`, `service/src/bridge/worker-admission.test.ts`, diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index c5bc573c..3b29122b 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -239,12 +239,14 @@ execution. worker. The lower API or worker slot ceiling wins, and assignments sharing the same workspace isolation key remain serialized while independent conversation worktrees may run concurrently. -- Workspace tool admission waits for capacity up to the smaller of `JOB_TIMEOUT` - and five minutes while the HTTP caller remains connected. Disconnects cancel - waiting, and admitted work receives a separate execution budget. A shorter - client or proxy timeout can end the wait sooner; Code API does not receive an - absolute caller deadline. See [BYOM worker admission](../byom-worker-admission.md) - for the caller and proxy timeout requirements. +- Workspace tool admission waits for capacity up to 30 seconds without a + `X-LibreChat-Workspace-Queue-Wait-Ms` header. A caller can advertise a longer + per-request allowance, bounded by five minutes and any server queue ceiling. + Disconnects cancel waiting, and admitted work receives a separate execution + budget capped by `JOB_TIMEOUT`. A shorter client or proxy timeout can end the + wait sooner; Code API does not receive an absolute caller deadline. See + [BYOM worker admission](../byom-worker-admission.md) for the total-request + and proxy timeout requirements. - Dynamic workers are fenced to their server-issued tenant before assignment. - Each assignment has an absolute deadline, generation, and random lease token. - Settlements with the wrong worker, generation, token, or expired deadline are diff --git a/packages/code/README.md b/packages/code/README.md index 87a6fe61..95d86485 100644 --- a/packages/code/README.md +++ b/packages/code/README.md @@ -791,17 +791,22 @@ Legacy requests without a conversation identity continue to use the selected source root. Older Code API deployments do not negotiate the capability, so the worker omits it until every request path understands the isolation boundary. -On an updated Code API, admission waits up to the smaller of `JOB_TIMEOUT` and -five minutes while the HTTP caller remains connected; older Code API versions -waited at most 30 seconds. A `WORKSPACE_QUEUE_TIMEOUT` response (HTTP 503, +On an updated Code API, admission waits up to 30 seconds without the +`X-LibreChat-Workspace-Queue-Wait-Ms` request header. A caller may advertise a +positive integer millisecond allowance up to five minutes, capped by any server +queue ceiling. This allowance is separate from the `JOB_TIMEOUT` execution +budget and cannot outlast a shorter client or proxy timeout. A +`WORKSPACE_QUEUE_TIMEOUT` response (HTTP 503, `Retry-After: 1`) means the operation was not assigned or started; wait for capacity before submitting it again. This is distinct from `ASSIGNMENT_EXPIRED` or a transport timeout after dispatch, where execution may have occurred and mutations must not be blindly retried. No automatic retry is added by this policy. -Align the client's per-attempt timeout and any proxy with the queue **plus** -execution budget before relying on the longer wait. See the +Align the client's per-attempt timeout and every proxy with the queue **plus** +execution, settlement, and delivery budget before relying on a longer wait. +LibreChat keeps the 30-second admission allowance unless its longer total HTTP +budget is explicitly enabled and the live ingress path is verified. See the [BYOM worker admission guide](../../docs/byom-worker-admission.md) for the -current client limitation and the timeout calculations. +timeout calculations. Keep the existing URL, pairing/identity, and network policy configuration. The primary root keeps its configured workspace ID (default `primary`). Repeat diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index 3d9d7e1c..0ca75566 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -65,6 +65,26 @@ export class BridgeStoreError extends Error { } } +function classifyPreEnqueueExpiry( + error: unknown, + args: { workspaceRequest?: WorkspaceToolRequest; workspaceId?: string; signal: AbortSignal }, +): unknown { + // A queue deadline reached before the assignment enqueue attempt is a + // definite non-execution, including expiry during the initial Redis reads. + if ( + (args.workspaceRequest != null || args.workspaceId != null) && + !args.signal.aborted && + error instanceof BridgeStoreError && + error.code === 'ASSIGNMENT_EXPIRED' + ) { + return new BridgeStoreError( + 'WORKSPACE_QUEUE_TIMEOUT', + 'Workspace capacity was unavailable before the queue deadline. The operation was not started. Wait for active work to finish or select an independent workspace on a machine with available capacity.', + ); + } + return error; +} + interface StoredAssignment extends CodeBridgeAssignment { leaseTokenHash: string; workerIdentityId?: string; @@ -778,12 +798,17 @@ export class RedisBridgeStore { 'Invalid workspace execution budget', ); } - this.assertDispatchActive(args.signal, args.deadlineAtMs); - const dispatchable = await this.dispatchCommand( - () => this.dispatchableRegistration(args.workerId), - args, - 'Bridge worker registration read', - ); + let dispatchable: Awaited>; + try { + this.assertDispatchActive(args.signal, args.deadlineAtMs); + dispatchable = await this.dispatchCommand( + () => this.dispatchableRegistration(args.workerId), + args, + 'Bridge worker registration read', + ); + } catch (error) { + throw classifyPreEnqueueExpiry(error, args); + } if (dispatchable == null) { throw new BridgeStoreError( 'WORKER_OFFLINE', @@ -844,17 +869,20 @@ export class RedisBridgeStore { `Bridge worker ${args.workerId} does not advertise programmatic execution for the selected workspace`, ); } - if ( - args.runtimeSessionId !== undefined && - (await this.dispatchCommand( - () => - this.redis.exists( - workspaceQuarantineKey(args.workerId, args.runtimeSessionId ?? ''), - ), - args, - 'Bridge workspace fence read', - )) === 1 - ) { + let quarantined = false; + const runtimeSessionId = args.runtimeSessionId; + if (runtimeSessionId !== undefined) { + try { + quarantined = (await this.dispatchCommand( + () => this.redis.exists(workspaceQuarantineKey(args.workerId, runtimeSessionId)), + args, + 'Bridge workspace fence read', + )) === 1; + } catch (error) { + throw classifyPreEnqueueExpiry(error, args); + } + } + if (quarantined) { throw new BridgeStoreError( 'WORKSPACE_QUARANTINED', 'Bridge workspace is quarantined after an incomplete result commit', @@ -1181,15 +1209,7 @@ export class RedisBridgeStore { } } catch (error) { // Once enqueue starts, even a lost Redis response may hide execution. - if ( - admission != null && !enqueueAttempted && !args.signal.aborted && - error instanceof BridgeStoreError && error.code === 'ASSIGNMENT_EXPIRED' - ) { - throw new BridgeStoreError( - 'WORKSPACE_QUEUE_TIMEOUT', - 'Workspace capacity was unavailable before the queue deadline. The operation was not started. Wait for active work to finish or select an independent workspace on a machine with available capacity.', - ); - } + if (!enqueueAttempted) throw classifyPreEnqueueExpiry(error, args); throw error; } finally { if (admission != null) { diff --git a/service/src/bridge/worker-admission.test.ts b/service/src/bridge/worker-admission.test.ts index 98363c34..c77ac19e 100644 --- a/service/src/bridge/worker-admission.test.ts +++ b/service/src/bridge/worker-admission.test.ts @@ -124,6 +124,32 @@ test('an expired queued call never reaches the worker and does not strand later await third; }); +test('an already-expired workspace deadline is a definite queue timeout before registration', async () => { + await register(); + await expect(dispatch('expired-before-read', new AbortController(), -1)).rejects.toMatchObject({ + code: 'WORKSPACE_QUEUE_TIMEOUT', + }); + expect(await redis.zcard(`codeapi:bridge:v1:worker:${workerId}:admission`)).toBe(0); + expect(await store.lease(workerId, incarnationId, 20)).toBeUndefined(); +}); + +test('expiry during the registration read never becomes an ambiguous assignment error', async () => { + await register(); + const registrationRead = spyOn(redis, 'mget').mockImplementation(async () => { + await new Promise(resolve => setTimeout(resolve, 25)); + throw new Error('registration read outlived its queue budget'); + }); + try { + await expect(dispatch('expired-during-read', new AbortController(), 1)).rejects.toMatchObject({ + code: 'WORKSPACE_QUEUE_TIMEOUT', + }); + } finally { + registrationRead.mockRestore(); + } + expect(await redis.zcard(`codeapi:bridge:v1:worker:${workerId}:admission`)).toBe(0); + expect(await store.lease(workerId, incarnationId, 20)).toBeUndefined(); +}); + test('a queued request is rejected if the worker withdraws its capability', async () => { await register(); const first = dispatch('first'); diff --git a/service/src/workspace-tools/router.test.ts b/service/src/workspace-tools/router.test.ts index 18d24a76..d3cb032e 100644 --- a/service/src/workspace-tools/router.test.ts +++ b/service/src/workspace-tools/router.test.ts @@ -67,16 +67,17 @@ test('binds instance admission to the authenticated tenant and user while preser expect(dispatched[1]?.workspaceInstanceId).toBeUndefined(); }); -test.each<[WorkspaceToolRequest, number, number?, number?]>([ - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 90_000 }, 95_000, 125_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 300_000 }, 305_000, undefined, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 6000, 1000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, 600_000, undefined], - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, 5000], -])('separates the admission deadline from execution budget for %j', async (request, expectedExecution, ceiling, queueTimeoutMs) => { +test.each<[WorkspaceToolRequest, number, number?, number?, number?]>([ + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 90_000 }, 95_000, 125_000, undefined, 90_000], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 300_000 }, 305_000, undefined, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 6000, 1000, undefined, undefined], + [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, 600_000, undefined, undefined], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, 125_000, 5000, 90_000], + [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined, undefined, 300_000], +])('separates the admission deadline from execution budget for %j', async (request, expectedExecution, ceiling, queueTimeoutMs, advertisedQueueWaitMs) => { const app = express(); app.use(json()); app.use((req, _res, next) => { @@ -102,12 +103,15 @@ test.each<[WorkspaceToolRequest, number, number?, number?]>([ const address = server.address(); if (address == null || typeof address === 'string') throw new Error('Missing listener'); const response = await fetch(`http://127.0.0.1:${address.port}/workspace-tools/execute`, { - method: 'POST', headers: { 'Content-Type': 'application/json' }, body: JSON.stringify(request), + method: 'POST', headers: { + 'Content-Type': 'application/json', + ...(advertisedQueueWaitMs === undefined ? {} : { 'X-LibreChat-Workspace-Queue-Wait-Ms': String(advertisedQueueWaitMs) }), + }, body: JSON.stringify(request), }); await response.json(); expect(executionBudget).toBe(expectedExecution); if (request.operation === 'execute_command') expect(commandTimeout).toBe(expectedExecution - 5000); - const expectedQueueBudget = queueTimeoutMs ?? Math.min(ceiling ?? 30_000, 300_000); + const expectedQueueBudget = Math.min(advertisedQueueWaitMs ?? 30_000, queueTimeoutMs ?? 300_000); expect(queueRemaining).toBeGreaterThan(expectedQueueBudget - 1000); expect(queueRemaining).toBeLessThanOrEqual(expectedQueueBudget); }); @@ -122,6 +126,32 @@ test.each([0, 300_001, Number.POSITIVE_INFINITY])('rejects an unbounded queue ov })).toThrow('Workspace queue timeout must be between 1 and 300000 milliseconds'); }); +test.each(['0', '-1', '300001', '1.5', '01', '1, 2', '999999999999999999999'])('rejects invalid per-request queue allowance %j before dispatch', async (queueWait) => { + const app = express(); + app.use(json()); + app.use((req, _res, next) => { + applyPrincipal(req, { userId: 'user-1', tenantId: 'tenant-1', principalSource: 'librechat_jwt', codeWorkerId: 'user-worker' }); + next(); + }); + let dispatched = false; + app.use(createWorkspaceToolsRouter({ + backend: 'remote-bridge', configuredWorkerId: 'user-worker', dynamicWorkers: false, + store: { async dispatchWorkspaceTool() { dispatched = true; throw new Error('must not dispatch'); } }, + })); + server = createServer(app); + await new Promise(resolve => server!.listen(0, '127.0.0.1', resolve)); + const address = server.address(); + if (address == null || typeof address === 'string') throw new Error('Missing listener'); + const response = await fetch(`http://127.0.0.1:${address.port}/workspace-tools/execute`, { + method: 'POST', + headers: { 'Content-Type': 'application/json', 'X-LibreChat-Workspace-Queue-Wait-Ms': queueWait }, + body: JSON.stringify({ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }), + }); + expect(response.status).toBe(400); + expect(await response.json()).toMatchObject({ code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + expect(dispatched).toBe(false); +}); + test('rejects new workspace dispatches while the service is shutting down', async () => { let dispatched = false; const app = express(); @@ -419,7 +449,7 @@ test.each([ status: expectedStatus, errorCode, outcome: 'completed', - deadlineBudgetMs: 330_000, + deadlineBudgetMs: 60_000, dispatchDurationMs: expect.any(Number), }), ); @@ -549,7 +579,8 @@ test('logs a disconnected dispatch once without inventing HTTP 200', async () => closeConnection(); await expect(response).rejects.toThrow(); await closed.promise; - expect(queueRemaining).toBeGreaterThan(120_000); + expect(queueRemaining).toBeGreaterThan(29_000); + expect(queueRemaining).toBeLessThanOrEqual(30_000); expect(dispatchAborted).toBe(true); expect(logSpy).not.toHaveBeenCalled(); settlementGate.resolve(); diff --git a/service/src/workspace-tools/router.ts b/service/src/workspace-tools/router.ts index 51e47f26..78dbd1e7 100644 --- a/service/src/workspace-tools/router.ts +++ b/service/src/workspace-tools/router.ts @@ -21,6 +21,8 @@ import { import { principalWorkspaceInstanceId } from '../bridge/workspace-instance'; const MAX_WORKSPACE_QUEUE_WAIT_MS = 5 * 60_000; +const DEFAULT_WORKSPACE_QUEUE_WAIT_MS = 30_000; +const WORKSPACE_QUEUE_WAIT_HEADER = 'X-LibreChat-Workspace-Queue-Wait-Ms'; interface WorkspaceToolsRouterOptions { store: Pick; @@ -57,12 +59,10 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) )) { throw new RangeError('Workspace execution timeout must be a positive safe integer'); } - // The HTTP disconnect cancels waiting; this bounds admission while the caller remains connected. - const queueBudgetMs = options.queueTimeoutMs ?? Math.min( - options.timeoutMs ?? 30_000, - MAX_WORKSPACE_QUEUE_WAIT_MS, - ); - if (!Number.isSafeInteger(queueBudgetMs) || queueBudgetMs < 1 || queueBudgetMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { + // A configured queue timeout is a ceiling. Legacy callers get 30 seconds even + // when their execution timeout is longer; only the per-request header opts in. + const queueCeilingMs = options.queueTimeoutMs ?? MAX_WORKSPACE_QUEUE_WAIT_MS; + if (!Number.isSafeInteger(queueCeilingMs) || queueCeilingMs < 1 || queueCeilingMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { throw new RangeError('Workspace queue timeout must be between 1 and 300000 milliseconds'); } const router = Router(); @@ -88,6 +88,21 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) }); return; } + const advertisedQueueWait = req.header(WORKSPACE_QUEUE_WAIT_HEADER); + if (advertisedQueueWait !== undefined && !/^[1-9]\d*$/.test(advertisedQueueWait)) { + outcome.errorCode = 'INVALID_WORKSPACE_QUEUE_WAIT'; + res.status(400).json({ error: 'Invalid workspace queue wait', code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + return; + } + const requestedQueueWaitMs = advertisedQueueWait === undefined + ? DEFAULT_WORKSPACE_QUEUE_WAIT_MS + : Number(advertisedQueueWait); + if (!Number.isSafeInteger(requestedQueueWaitMs) || requestedQueueWaitMs > MAX_WORKSPACE_QUEUE_WAIT_MS) { + outcome.errorCode = 'INVALID_WORKSPACE_QUEUE_WAIT'; + res.status(400).json({ error: 'Invalid workspace queue wait', code: 'INVALID_WORKSPACE_QUEUE_WAIT' }); + return; + } + const queueBudgetMs = Math.min(requestedQueueWaitMs, queueCeilingMs); outcome.operation = req.body.operation; const principalRequest: WorkspaceToolRequest = req.body.workspaceInstanceId == null ? req.body