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
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions docs/remote-bridge/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -384,3 +384,58 @@ 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. 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
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.

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. 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
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.
105 changes: 105 additions & 0 deletions service/openapi.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down
55 changes: 55 additions & 0 deletions service/src/bridge/admission.ts
Original file line number Diff line number Diff line change
@@ -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(
Expand Down Expand Up @@ -57,6 +74,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<number | undefined> {
const rank = await this.redis.zrank(this.keys(workerId)[0], id);
return rank == null ? undefined : rank + 1;
}

async isHead(workerId: string, id: string): Promise<boolean> {
const [order, deadlines, , workspaces] = this.keys(workerId);
return (
Expand Down
22 changes: 19 additions & 3 deletions service/src/bridge/slots.ts
Original file line number Diff line number Diff line change
@@ -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;
Expand Down Expand Up @@ -37,6 +38,8 @@ export class BridgeWorkspaceSlots {
workspaceId: string;
capacity: number;
expiresAtMs: number;
refreshOwned?: boolean;
guard?: { key: string; claimKey: string; token: string };
}): Promise<number | undefined> {
if (
!Number.isSafeInteger(args.capacity) ||
Expand All @@ -48,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.
Expand Down Expand Up @@ -83,7 +89,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',
Expand Down Expand Up @@ -125,14 +139,16 @@ 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,
args.capacity,
args.expiresAtMs,
Date.now(),
args.refreshOwned === true ? '1' : '0',
args.guard?.token ?? '',
),
);
if (result === -2)
Expand Down
Loading
Loading