Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 26 additions & 16 deletions docs/byom-worker-admission.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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`,
Expand Down
14 changes: 8 additions & 6 deletions docs/remote-bridge/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
17 changes: 11 additions & 6 deletions packages/code/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 46 additions & 26 deletions service/src/bridge/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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<ReturnType<typeof this.dispatchableRegistration>>;
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',
Expand Down Expand Up @@ -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',
Expand Down Expand Up @@ -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) {
Expand Down
26 changes: 26 additions & 0 deletions service/src/bridge/worker-admission.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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');
Expand Down
59 changes: 45 additions & 14 deletions service/src/workspace-tools/router.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) => {
Expand All @@ -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);
});
Expand All @@ -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<void>(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();
Expand Down Expand Up @@ -419,7 +449,7 @@ test.each([
status: expectedStatus,
errorCode,
outcome: 'completed',
deadlineBudgetMs: 330_000,
deadlineBudgetMs: 60_000,
dispatchDurationMs: expect.any(Number),
}),
);
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading