diff --git a/packages/api/src/code/bridge.spec.ts b/packages/api/src/code/bridge.spec.ts index faf724de7f3..bd9dae4bf0f 100644 --- a/packages/api/src/code/bridge.spec.ts +++ b/packages/api/src/code/bridge.spec.ts @@ -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({ @@ -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('Bad gateway', { 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>({ 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>({ reason: 'rejected' }), + ); + await expect(poll(params)).rejects.toEqual( + expect.objectContaining>({ 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>({ 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>({ 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>({ reason: 'rejected' }), + ); + await expect(poll(params)).rejects.toEqual( + expect.objectContaining>({ 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>({ reason: 'failed' }), + ); + }); +}); diff --git a/packages/api/src/code/bridge.ts b/packages/api/src/code/bridge.ts index c3cb7af0e53..0820901eca4 100644 --- a/packages/api/src/code/bridge.ts +++ b/packages/api/src/code/bridge.ts @@ -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, @@ -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, @@ -98,19 +116,78 @@ export function createCodeBridgeStatusPoller({ string, { expiresAt: number; request: Promise } >(); + const lastReady = new Map(); 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 => { + 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 => { + 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; @@ -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 =