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
169 changes: 168 additions & 1 deletion packages/api/src/code/bridge.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -482,7 +482,7 @@ describe('getCodeBridgeWorkerStatus', () => {
});
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(new Response(rejectedBody, { status: 503 }))
.mockResolvedValueOnce(new Response(rejectedBody, { status: 403 }))
.mockResolvedValue(
new Response(
JSON.stringify({
Expand Down Expand Up @@ -514,3 +514,170 @@ describe('getCodeBridgeWorkerStatus', () => {
});
});
});

describe('createCodeBridgeStatusPoller transient failures', () => {
const params = {
baseURL: 'https://code.example.com/v1',
token: 'administrator-token',
workerId: 'personal-vm',
};

const statusResponse = (online: boolean, ready: boolean, leaseExpiresInMs = 60_000) =>
new Response(
JSON.stringify({
protocolVersion: 1,
workerId: 'personal-vm',
online,
ready,
...(online
? { leaseExpiresInMs, capabilities: { sandboxProfile: 'native-srt', runtimes: ['bash'] } }
: {}),
}),
);
const networkFailure = () => Promise.reject(new TypeError('fetch failed'));

let now = 10_000;
beforeEach(() => {
now = 10_000;
jest.spyOn(Date, 'now').mockImplementation(() => now);
});
afterEach(() => {
jest.restoreAllMocks();
});

test('answers a failed poll with the last ready status while its lease runs', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true, 60_000))
.mockImplementationOnce(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await expect(poll(params)).resolves.toMatchObject({ status: 'ready' });
now += 15_000;
await expect(poll(params)).resolves.toMatchObject({
status: 'ready',
leaseExpiresInMs: 45_000,
});

expect(fetchImpl).toHaveBeenCalledTimes(2);
});

test('treats a 5xx from the proxy in front of the Code API as transient', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockResolvedValueOnce(new Response('<html>Bad gateway</html>', { status: 502 }));
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
now += 1_000;
await expect(poll(params)).resolves.toMatchObject({ status: 'ready' });
});

test('retries once when no ready status is remembered', async () => {
const fetchImpl = jest
.fn()
.mockImplementationOnce(networkFailure)
.mockResolvedValueOnce(statusResponse(true, true));
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await expect(poll(params)).resolves.toMatchObject({ status: 'ready' });
expect(fetchImpl).toHaveBeenCalledTimes(2);
});

test('fails after one retry once the remembered lease has run out', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true, 5_000))
.mockImplementation(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
now += 5_000;
await expect(poll(params)).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
expect(fetchImpl).toHaveBeenCalledTimes(3);
});

test('never covers a rejection, and forgets the remembered status after one', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockResolvedValueOnce(new Response('{}', { status: 401 }))
.mockImplementation(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
await expect(poll(params)).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'rejected' }),
);
await expect(poll(params)).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
expect(fetchImpl).toHaveBeenCalledTimes(4);
});

test('forgets a ready status once the worker reports it is no longer ready', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockResolvedValueOnce(statusResponse(false, false))
.mockImplementation(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
await expect(poll(params)).resolves.toMatchObject({ status: 'offline' });
await expect(poll(params)).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
});

test('gives bypassCache callers one fresh observation, never the remembered status', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockImplementationOnce(networkFailure)
.mockResolvedValueOnce(statusResponse(true, false));
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
await expect(poll({ ...params, bypassCache: true })).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
await expect(poll({ ...params, bypassCache: true })).resolves.toMatchObject({
status: 'starting',
});
expect(fetchImpl).toHaveBeenCalledTimes(3);
});

test('forgets the remembered status when a selection check is rejected', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockResolvedValueOnce(new Response('{}', { status: 404 }))
.mockImplementation(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
await expect(poll({ ...params, bypassCache: true })).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'rejected' }),
);
await expect(poll(params)).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
});

test('keeps remembered statuses separate per credential', async () => {
const fetchImpl = jest
.fn()
.mockResolvedValueOnce(statusResponse(true, true))
.mockImplementation(networkFailure);
const poll = createCodeBridgeStatusPoller({ fetchImpl, cacheTtlMs: 0 });

await poll(params);
await expect(poll({ ...params, token: 'rotated-administrator-token' })).rejects.toEqual(
expect.objectContaining<Partial<CodeBridgeStatusError>>({ reason: 'failed' }),
);
});
});
87 changes: 82 additions & 5 deletions packages/api/src/code/bridge.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { createHash } from 'node:crypto';
import { logger } from '@librechat/data-schemas';
import {
CODE_ENVIRONMENT_COMMAND_TIMEOUT_HARD_MAX_MS,
CODE_WORKSPACE_ID_PATTERN,
Expand Down Expand Up @@ -77,6 +78,23 @@ export class CodeBridgeStatusError extends Error {
}
}

/** Transport failures that say nothing about the worker itself; its own heartbeat lease still does. */
function isTransientStatusFailure(error: unknown): boolean {
if (!(error instanceof CodeBridgeStatusError)) return false;
if (error.reason === 'timeout' || error.reason === 'failed') return true;
return error.reason === 'rejected' && error.upstreamStatus != null && error.upstreamStatus >= 500;
}

type RememberedReadyStatus = { status: CodeBridgeWorkerStatus; leaseExpiresAt: number };

/**
* Polls worker status with a short shared cache. A transient transport failure (timeout, network
* error, 5xx from the Code API or the proxy in front of it) is not evidence that the worker left:
* while the last `ready` observation's heartbeat lease is still running, that observation answers
* for it, and without one the poll is retried once. Rejections and malformed responses fail
* immediately and forget the remembered status. `bypassCache` callers validate a selection before
* it is persisted, so they get exactly one fresh observation: no retry, no remembered status.
*/
export function createCodeBridgeStatusPoller({
fetchImpl,
maxConcurrent = 32,
Expand All @@ -98,19 +116,78 @@ export function createCodeBridgeStatusPoller({
string,
{ expiresAt: number; request: Promise<CodeBridgeWorkerStatus> }
>();
const lastReady = new Map<string, RememberedReadyStatus>();
let active = 0;

const remember = (key: string, status: CodeBridgeWorkerStatus, startedAt: number): void => {
lastReady.delete(key);
if (status.status !== 'ready' || status.leaseExpiresInMs == null) return;
if (lastReady.size >= maxEntries) {
const oldest = lastReady.keys().next().value;
if (oldest != null) lastReady.delete(oldest);
}
/** The lease was read after the request started, so this bound never outlives it. */
lastReady.set(key, { status, leaseExpiresAt: startedAt + status.leaseExpiresInMs });
};

const recall = (key: string): CodeBridgeWorkerStatus | undefined => {
const entry = lastReady.get(key);
if (entry == null) return undefined;
const remainingMs = entry.leaseExpiresAt - Date.now();
if (remainingMs <= 0) {
lastReady.delete(key);
return undefined;
}
return { ...entry.status, leaseExpiresInMs: remainingMs };
};

const fetchStatus = async (
params: { baseURL: string; token: string; workerId: string },
key: string,
): Promise<CodeBridgeWorkerStatus> => {
const startedAt = Date.now();
try {
const status = await getCodeBridgeWorkerStatus({ ...params, fetchImpl });
remember(key, status, startedAt);
return status;
} catch (error) {
if (!isTransientStatusFailure(error)) lastReady.delete(key);
throw error;
}
};

const observe = async (
params: { baseURL: string; token: string; workerId: string },
key: string,
): Promise<CodeBridgeWorkerStatus> => {
try {
return await fetchStatus(params, key);
} catch (error) {
if (!isTransientStatusFailure(error)) throw error;
const remembered = recall(key);
if (remembered != null) {
logger.warn(
`[codeBridge] Worker status poll failed (${(error as CodeBridgeStatusError).reason}); ` +
`using the last ready status for ${remembered.leaseExpiresInMs}ms more of its lease`,
);
return remembered;
}
return await fetchStatus(params, key);
}
};

return (params) => {
const credentialId = createHash('sha256').update(params.token).digest('base64url');
const normalizedBaseURL = params.baseURL.trim().replace(/\/+$/, '');
const key = `${normalizedBaseURL}\u0000${params.workerId}\u0000${credentialId}`;
if (params.bypassCache) {
if (active >= maxConcurrent) return Promise.reject(new CodeBridgeStatusError('busy'));
active += 1;
// Mutation validation needs its own observation, but still shares the polling capacity.
return getCodeBridgeWorkerStatus({ ...params, fetchImpl }).finally(() => {
return fetchStatus(params, key).finally(() => {
active -= 1;
});
}
const credentialId = createHash('sha256').update(params.token).digest('base64url');
const normalizedBaseURL = params.baseURL.trim().replace(/\/+$/, '');
const key = `${normalizedBaseURL}\u0000${params.workerId}\u0000${credentialId}`;
const now = Date.now();
const cached = requests.get(key);
if (cached != null && cached.expiresAt > now) return cached.request;
Expand All @@ -126,7 +203,7 @@ export function createCodeBridgeStatusPoller({
}
active += 1;
const startedAt = Date.now();
const request = getCodeBridgeWorkerStatus({ ...params, fetchImpl })
const request = observe(params, key)
.then((status) => {
const completedAt = Date.now();
const ttl =
Expand Down
Loading