From 52d9a9d4ffdb0e5abfbe11c6f662c5b6ee213c93 Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 27 Sep 2026 13:22:45 +0000 Subject: [PATCH 1/2] =?UTF-8?q?=E2=8F=B3=20fix:=20Let=20BYOM=20Workspace?= =?UTF-8?q?=20Calls=20Queue=20Past=20Thirty=20Seconds?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/remote-bridge/README.md | 5 +++ service/src/workspace-tools/router.test.ts | 41 ++++++++++++++++------ service/src/workspace-tools/router.ts | 28 +++++++++------ 3 files changed, 52 insertions(+), 22 deletions(-) diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index efe84b84..9d41846a 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -239,6 +239,11 @@ 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. - 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/service/src/workspace-tools/router.test.ts b/service/src/workspace-tools/router.test.ts index a52b3f0b..18d24a76 100644 --- a/service/src/workspace-tools/router.test.ts +++ b/service/src/workspace-tools/router.test.ts @@ -67,13 +67,16 @@ test('binds instance admission to the authenticated tenant and user while preser expect(dispatched[1]?.workspaceInstanceId).toBeUndefined(); }); -test.each<[WorkspaceToolRequest, number, number?]>([ - [{ protocolVersion: 1, operation: 'read_file', workspaceId: 'primary', path: 'README.md' }, 30_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready', timeoutMs: 300_000 }, 305_000, undefined], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 6000, 1000], - [{ protocolVersion: 1, operation: 'execute_command', workspaceId: 'primary', command: 'echo ready' }, 35_000, 600_000], -])('separates the admission deadline from execution budget for %j', async (request, expectedExecution, ceiling) => { +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) => { const app = express(); app.use(json()); app.use((req, _res, next) => { @@ -86,6 +89,7 @@ test.each<[WorkspaceToolRequest, number, number?]>([ app.use(createWorkspaceToolsRouter({ backend: 'remote-bridge', configuredWorkerId: 'user-worker', dynamicWorkers: false, timeoutMs: ceiling, + queueTimeoutMs, store: { async dispatchWorkspaceTool(args) { executionBudget = args.executionTimeoutMs; if (args.request.operation === 'execute_command') commandTimeout = args.request.timeoutMs; @@ -103,8 +107,19 @@ test.each<[WorkspaceToolRequest, number, number?]>([ await response.json(); expect(executionBudget).toBe(expectedExecution); if (request.operation === 'execute_command') expect(commandTimeout).toBe(expectedExecution - 5000); - expect(queueRemaining).toBeGreaterThan(29_000); - expect(queueRemaining).toBeLessThanOrEqual(30_000); + const expectedQueueBudget = queueTimeoutMs ?? Math.min(ceiling ?? 30_000, 300_000); + expect(queueRemaining).toBeGreaterThan(expectedQueueBudget - 1000); + expect(queueRemaining).toBeLessThanOrEqual(expectedQueueBudget); +}); + +test.each([0, 300_001, Number.POSITIVE_INFINITY])('rejects an unbounded queue override (%s)', (queueTimeoutMs) => { + expect(() => createWorkspaceToolsRouter({ + backend: 'remote-bridge', + configuredWorkerId: 'user-worker', + dynamicWorkers: false, + queueTimeoutMs, + store: { async dispatchWorkspaceTool() { throw new Error('must not dispatch'); } }, + })).toThrow('Workspace queue timeout must be between 1 and 300000 milliseconds'); }); test('rejects new workspace dispatches while the service is shutting down', async () => { @@ -404,7 +419,7 @@ test.each([ status: expectedStatus, errorCode, outcome: 'completed', - deadlineBudgetMs: 60_000, + deadlineBudgetMs: 330_000, dispatchDurationMs: expect.any(Number), }), ); @@ -486,6 +501,7 @@ test('logs a disconnected dispatch once without inventing HTTP 200', async () => const closed = Promise.withResolvers(); const settlementGate = Promise.withResolvers(); let dispatchAborted = false; + let queueRemaining: number | undefined; let closeConnection = (): void => { throw new Error('connection not ready'); }; app.use(json()); app.use((req, res, next) => { @@ -499,8 +515,10 @@ test('logs a disconnected dispatch once without inventing HTTP 200', async () => backend: 'remote-bridge', configuredWorkerId: 'user-worker', dynamicWorkers: true, + timeoutMs: 125_000, store: { - async dispatchWorkspaceTool({ signal }) { + async dispatchWorkspaceTool({ deadlineAtMs, signal }) { + queueRemaining = deadlineAtMs - Date.now(); started.resolve(); return await new Promise((_resolve, reject) => { signal.addEventListener( @@ -531,6 +549,7 @@ 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(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 3e1f9fee..51e47f26 100644 --- a/service/src/workspace-tools/router.ts +++ b/service/src/workspace-tools/router.ts @@ -20,6 +20,8 @@ import { } from '../bridge/selection'; import { principalWorkspaceInstanceId } from '../bridge/workspace-instance'; +const MAX_WORKSPACE_QUEUE_WAIT_MS = 5 * 60_000; + interface WorkspaceToolsRouterOptions { store: Pick; backend: 'http' | 'lambda-microvm' | 'remote-bridge'; @@ -50,15 +52,19 @@ export function bridgeStoreStatus(error: BridgeStoreError): number { } export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions): Router { - const queueBudgetMs = options.queueTimeoutMs ?? 30_000; - if (!Number.isSafeInteger(queueBudgetMs) || queueBudgetMs < 1 || queueBudgetMs > 30_000) { - throw new RangeError('Workspace queue timeout must be between 1 and 30000 milliseconds'); - } if (options.timeoutMs !== undefined && ( !Number.isSafeInteger(options.timeoutMs) || options.timeoutMs < 1 )) { 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) { + throw new RangeError('Workspace queue timeout must be between 1 and 300000 milliseconds'); + } const router = Router(); router.post( @@ -86,13 +92,13 @@ export function createWorkspaceToolsRouter(options: WorkspaceToolsRouterOptions) const principalRequest: WorkspaceToolRequest = req.body.workspaceInstanceId == null ? req.body : { - ...req.body, - workspaceInstanceId: principalWorkspaceInstanceId({ - instanceId: req.body.workspaceInstanceId, - tenantId: principal.tenantId, - principalId: principal.userId, - }), - }; + ...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, From bbd246a2971abfe7ab44ef52974575ed671ca4c4 Mon Sep 17 00:00:00 2001 From: Lia Date: Sun, 27 Sep 2026 15:57:30 +0000 Subject: [PATCH 2/2] fix: Recompute BYOM Reservation TTLs After Admission --- docs/byom-worker-admission.md | 48 +++++++++------ docs/remote-bridge/README.md | 3 +- packages/code/README.md | 10 +++- service/src/bridge/concurrent-store.test.ts | 45 ++++++++++++++ service/src/bridge/store.ts | 14 +++-- service/src/bridge/worker-admission.test.ts | 66 +++++++++++++++++++++ 6 files changed, 161 insertions(+), 25 deletions(-) diff --git a/docs/byom-worker-admission.md b/docs/byom-worker-admission.md index 2c2b354b..2fcd9f9f 100644 --- a/docs/byom-worker-admission.md +++ b/docs/byom-worker-admission.md @@ -5,14 +5,17 @@ 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 workspace HTTP endpoint allows at most 30 seconds for admission. 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. -Read/search/list operations receive up to 30 seconds, also capped by `JOB_TIMEOUT`. -Disconnecting or cancelling removes the waiting -request without cancelling the active assignment. Expired entries are pruned; -Redis key expiry also bounds state left by a crashed API process. +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 +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 +operations receive up to 30 seconds, also capped by `JOB_TIMEOUT`. +Disconnecting or cancelling removes a waiting request without cancelling the +active assignment. Expired entries are pruned; Redis key expiry also bounds +state left by a crashed API process. Reservations derive their TTL at acquisition +from the remaining absolute deadline or a fresh execution budget; enqueued +assignment records use the final execution deadline, not the elapsed queue budget. After admission, the API revalidates the worker incarnation, identity, tenant binding and workspace operation. A waiting request cannot migrate to a replacement @@ -26,13 +29,24 @@ 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. -LibreChat must allow queue time plus execution/settlement time and five seconds -for HTTP delivery: 65 seconds for reads, 70 seconds for default commands, and -340 seconds for five-minute commands. Either side can be upgraded first. Older -clients still cancel at their earlier deadline; newer clients preserve errors from -older servers without retrying mutations. Both updates are needed for the full -waiting budget. Any reverse proxy request timeout must accommodate these totals. -The worker package does not need an update for the deadline change. +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 +operations or `min(JOB_TIMEOUT, requested command timeout) + 5s` for commands, +plus five seconds for delivery. -Focused regression coverage lives in `service/src/bridge/admission.test.ts` and -`service/src/bridge/worker-admission.test.ts`. +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. + +Focused regression coverage lives in `service/src/bridge/admission.test.ts`, +`service/src/bridge/worker-admission.test.ts`, +`service/src/bridge/concurrent-store.test.ts`, and +`service/src/workspace-tools/router.test.ts`. diff --git a/docs/remote-bridge/README.md b/docs/remote-bridge/README.md index 9d41846a..c5bc573c 100644 --- a/docs/remote-bridge/README.md +++ b/docs/remote-bridge/README.md @@ -243,7 +243,8 @@ execution. 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. + absolute caller deadline. See [BYOM worker admission](../byom-worker-admission.md) + for the caller 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 5a06b92a..9dc322e3 100644 --- a/packages/code/README.md +++ b/packages/code/README.md @@ -780,11 +780,17 @@ 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. -Admission waits at most 30 seconds. A `WORKSPACE_QUEUE_TIMEOUT` response (HTTP -503, `Retry-After: 1`) means the operation was not assigned or started; wait for +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, +`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 +[BYOM worker admission guide](../../docs/byom-worker-admission.md) for the +current client limitation and the 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/concurrent-store.test.ts b/service/src/bridge/concurrent-store.test.ts index 33f3d082..6bdfde88 100644 --- a/service/src/bridge/concurrent-store.test.ts +++ b/service/src/bridge/concurrent-store.test.ts @@ -357,6 +357,51 @@ test('same-root work waits while another root progresses', async () => { await Promise.all([nextA, b]); }); +test('a long same-root queue allowance does not extend slot or assignment TTLs', async () => { + await register(); + const active = dispatch('a'); + const activeAssignment = await store.lease( + workerId, incarnationId, 1000, undefined, undefined, 0, + ); + const controller = new AbortController(); + const queued = store.dispatchWorkspaceTool({ + workerId, + signal: controller.signal, + deadlineAtMs: Date.now() + 300_000, + executionTimeoutMs: 305_000, + request: { protocolVersion: 1, operation: 'read_file', workspaceId: 'a', path: 'second.txt' }, + }); + void queued.catch(() => undefined); + await new Promise((resolve) => setTimeout(resolve, 150)); + await settle(activeAssignment!); + await active; + + const assignment = await store.lease( + workerId, incarnationId, 1000, undefined, undefined, 0, + ); + if (assignment == null) { + controller.abort(); + await queued.catch(() => undefined); + throw new Error('Queued request was not leased'); + } + expect(assignment.request).toMatchObject({ path: 'second.txt' }); + try { + const expiresAtMs = Date.parse(assignment.expiresAt); + const slotExpiresAtMs = Number(await redis.hget( + `codeapi:bridge:v1:worker:${workerId}:workspace-slots`, 'e:0', + )); + expect(slotExpiresAtMs).toBeGreaterThan(expiresAtMs + 28_000); + expect(slotExpiresAtMs).toBeLessThan(expiresAtMs + 32_000); + const assignmentTtlMs = await redis.pttl(`codeapi:bridge:v1:assignment:${assignment.assignmentId}`); + const remainingMs = expiresAtMs - Date.now(); + expect(assignmentTtlMs).toBeGreaterThan(remainingMs + 28_000); + expect(assignmentTtlMs).toBeLessThan(remainingMs + 32_000); + } finally { + await settle(assignment); + await queued; + } +}); + test('conversation worktrees on one source use independent capacity lanes', async () => { await register(); const firstId = 'a'.repeat(64); diff --git a/service/src/bridge/store.ts b/service/src/bridge/store.ts index 2ffeec66..3d9d7e1c 100644 --- a/service/src/bridge/store.ts +++ b/service/src/bridge/store.ts @@ -863,10 +863,6 @@ export class RedisBridgeStore { const assignmentId = randomBytes(18).toString('base64url'); const leaseToken = randomBytes(32).toString('base64url'); - // The lock is acquired before admission finishes; it must outlive the later execution deadline. - const ttlSeconds = assignmentTtlSeconds( - args.deadlineAtMs + (args.executionTimeoutMs ?? 0), - ); const lockIncarnationId = registration.incarnationId; let assignment: StoredAssignment | undefined; let enqueueAttempted = false; @@ -931,6 +927,14 @@ export class RedisBridgeStore { ); continue; } + // A queued request may have waited nearly its full admission budget. + // Reserve for the *remaining* absolute deadline or a fresh execution + // budget, not the original queue window plus execution again. + const ttlSeconds = assignmentTtlSeconds( + args.executionTimeoutMs === undefined + ? args.deadlineAtMs + : Date.now() + args.executionTimeoutMs, + ); if (workspaceSlots != null) { workspaceLeaseSlot = await this.dispatchCommand( () => @@ -1070,7 +1074,7 @@ export class RedisBridgeStore { enqueueAttempted = true; return this.enqueueForActiveIncarnation( assignment!, - ttlSeconds, + assignmentTtlSeconds(Date.parse(assignment!.expiresAt)), readyToken, ); }, diff --git a/service/src/bridge/worker-admission.test.ts b/service/src/bridge/worker-admission.test.ts index 655748b1..98363c34 100644 --- a/service/src/bridge/worker-admission.test.ts +++ b/service/src/bridge/worker-admission.test.ts @@ -161,6 +161,72 @@ test('execution receives a fresh budget after waiting and the lock covers long c await second; }); +test('a long queue allowance does not extend serial lock or assignment TTLs after admission', async () => { + await register(); + const active = dispatch('first'); + const activeAssignment = await store.lease(workerId, incarnationId, 1000); + const controller = new AbortController(); + const queued = dispatch('second', controller, 300_000, 305_000); + void queued.catch(() => undefined); + await new Promise((resolve) => setTimeout(resolve, 150)); + await settle(activeAssignment); + await active; + + const assignment = await store.lease(workerId, incarnationId, 1000); + if (assignment == null) { + controller.abort(); + await queued.catch(() => undefined); + throw new Error('Queued request was not leased'); + } + expect(assignment.request).toMatchObject({ path: 'second' }); + try { + const expiresAtMs = Date.parse(assignment.expiresAt); + for (const key of [ + `codeapi:bridge:v1:worker:${workerId}:lock`, + `codeapi:bridge:v1:worker:${workerId}:lock:incarnation`, + `codeapi:bridge:v1:assignment:${assignment.assignmentId}`, + ]) { + const ttlMs = await redis.pttl(key); + const remainingMs = expiresAtMs - Date.now(); + expect(ttlMs).toBeGreaterThan(remainingMs + 28_000); + expect(ttlMs).toBeLessThan(remainingMs + 32_000); + } + } finally { + await settle(assignment); + await queued; + } +}); + +test('absolute-deadline callers keep only their remaining deadline in the serial lock TTL', async () => { + await register(); + const active = dispatch('first'); + const activeAssignment = await store.lease(workerId, incarnationId, 1000); + const controller = new AbortController(); + const queued = dispatch('absolute', controller, 2_500); + void queued.catch(() => undefined); + await new Promise((resolve) => setTimeout(resolve, 1450)); + await settle(activeAssignment); + await active; + + const assignment = await store.lease(workerId, incarnationId, 1000); + if (assignment == null) { + controller.abort(); + await queued.catch(() => undefined); + throw new Error('Queued request was not leased'); + } + expect(assignment.request).toMatchObject({ path: 'absolute' }); + try { + const remainingMs = Date.parse(assignment.expiresAt) - Date.now(); + expect(remainingMs).toBeGreaterThan(0); + const ttlMs = await redis.pttl(`codeapi:bridge:v1:worker:${workerId}:lock`); + expect(ttlMs).toBeGreaterThan(remainingMs + 28_000); + expect(ttlMs).toBeLessThan(remainingMs + 31_000); + } finally { + await settle(assignment); + await queued; + } +}); + test('execution expires independently of an unused queue allowance', async () => { await register(); const completion = dispatch('short', new AbortController(), 5000, 150);