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
48 changes: 31 additions & 17 deletions docs/byom-worker-admission.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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`.
6 changes: 6 additions & 0 deletions docs/remote-bridge/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -239,6 +239,12 @@ 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.
- 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
10 changes: 8 additions & 2 deletions packages/code/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
45 changes: 45 additions & 0 deletions service/src/bridge/concurrent-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
14 changes: 9 additions & 5 deletions service/src/bridge/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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(
() =>
Expand Down Expand Up @@ -1070,7 +1074,7 @@ export class RedisBridgeStore {
enqueueAttempted = true;
return this.enqueueForActiveIncarnation(
assignment!,
ttlSeconds,
assignmentTtlSeconds(Date.parse(assignment!.expiresAt)),
readyToken,
);
},
Expand Down
66 changes: 66 additions & 0 deletions service/src/bridge/worker-admission.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
41 changes: 30 additions & 11 deletions service/src/workspace-tools/router.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) => {
Expand All @@ -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;
Expand All @@ -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 () => {
Expand Down Expand Up @@ -404,7 +419,7 @@ test.each([
status: expectedStatus,
errorCode,
outcome: 'completed',
deadlineBudgetMs: 60_000,
deadlineBudgetMs: 330_000,
dispatchDurationMs: expect.any(Number),
}),
);
Expand Down Expand Up @@ -486,6 +501,7 @@ test('logs a disconnected dispatch once without inventing HTTP 200', async () =>
const closed = Promise.withResolvers<void>();
const settlementGate = Promise.withResolvers<void>();
let dispatchAborted = false;
let queueRemaining: number | undefined;
let closeConnection = (): void => { throw new Error('connection not ready'); };
app.use(json());
app.use((req, res, next) => {
Expand All @@ -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(
Expand Down Expand Up @@ -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();
Expand Down
Loading
Loading