From 3142a2f6541641488d572da453efba728a6c0539 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Fri, 11 Sep 2026 11:10:37 -0400 Subject: [PATCH 1/2] fix: fail denied input downloads once per batch (#177) --- api/src/download.test.ts | 70 +++++++++++++++++++++++++++++++++++++++- api/src/job.ts | 40 +++++++++++++++++++---- 2 files changed, 103 insertions(+), 7 deletions(-) diff --git a/api/src/download.test.ts b/api/src/download.test.ts index 5372a40c..2c7248e8 100644 --- a/api/src/download.test.ts +++ b/api/src/download.test.ts @@ -1,4 +1,4 @@ -import { describe, it, expect, beforeEach, afterEach, beforeAll, afterAll } from 'bun:test'; +import { describe, it, expect, beforeEach, afterEach, beforeAll, afterAll, spyOn } from 'bun:test'; import * as fsp from 'fs/promises'; import * as path from 'path'; import * as os from 'os'; @@ -439,6 +439,74 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => { expect(contents).toBe('hi'); }); + it.each([401, 403])('does not retry an HTTP %i authorization denial', async status => { + const file: TFile = { id: 'denied', storage_session_id: 'previous', name: 'denied.txt' }; + let requests = 0; + routes.set('/sessions/previous/objects/denied', { + status, + onRequest: () => { requests++; }, + }); + const job = makeJob([file]); + asInternals(job).submissionDir = tmpDir; + + await expect(job.downloadAndWriteFile(file, 5, 1)).rejects.toThrow(`HTTP error: ${status}`); + expect(requests).toBe(1); + expect(await fsp.readdir(tmpDir)).toEqual([]); + }); + + it.each([404, 408, 429, 503])('still retries transient HTTP %i responses', async status => { + const file: TFile = { id: 'transient', storage_session_id: 'previous', name: 'ready.txt' }; + let requests = 0; + const route: Route = { + status, + body: 'ready', + onRequest: () => { if (++requests === 2) route.status = 200; }, + }; + routes.set('/sessions/previous/objects/transient', route); + const job = makeJob([file]); + asInternals(job).submissionDir = tmpDir; + + await expect(job.downloadAndWriteFile(file, 5, 1)).resolves.toBe('ready.txt'); + expect(requests).toBe(2); + expect(await fsp.readFile(path.join(tmpDir, 'ready.txt'), 'utf8')).toBe('ready'); + }); + + it('accounts for a denied 240-file batch once and stops queued downloads', async () => { + const files: TFile[] = Array.from({ length: 240 }, (_, index) => ({ + id: `file-${index}`, storage_session_id: 'previous', name: `file-${index}.txt`, + })); + let requests = 0; + routes.set('/sessions/previous/objects', { status: 200, body: '[]' }); + for (const file of files) { + routes.set(`/sessions/previous/objects/${file.id}`, { + status: 403, + delayMs: file.id === 'file-0' ? 0 : 30, + onRequest: () => { requests++; }, + }); + } + let dirty = false; + const job = makeJob(files, sessionWorkspaceAt(tmpDir, 'batch-test', () => { dirty = true; })); + const log = (job as unknown as { log: import('pino').Logger }).log; + const errorLog = spyOn(log, 'error'); + const originalConcurrency = config.prime_concurrency; + config.prime_concurrency = 8; + try { + await expect(job.prime()).rejects.toBeInstanceOf(SessionWorkspaceDirtyError); + expect(dirty).toBe(true); + expect(requests).toBeGreaterThan(0); + expect(requests).toBeLessThanOrEqual(8); + expect(errorLog).toHaveBeenCalledTimes(1); + expect(errorLog).toHaveBeenCalledWith(expect.objectContaining({ + inputCount: 240, completed: 0, failed: 1, cancelled: 7, notStarted: 232, + }), 'Input preparation batch failed'); + expect(await fsp.readdir(tmpDir)).toEqual([]); + } finally { + errorLog.mockRestore(); + config.prime_concurrency = originalConcurrency; + await job.cleanup(); + } + }); + it('fails when the server keeps 404-ing past the retry cap (no phantom write)', async () => { const file: TFile = { id: 'missing-id', diff --git a/api/src/job.ts b/api/src/job.ts index 610b748b..8121f51d 100644 --- a/api/src/job.ts +++ b/api/src/job.ts @@ -60,6 +60,14 @@ export { const AUTO_LOAD_DIRKEEP_TIMEOUT_MS = 10000; const AUTO_LOAD_DIRKEEP_RETRIES = 2; +/** Replaying the same sealed grant cannot repair an authorization denial. */ +class InputAuthorizationError extends Error { + constructor(status: number) { + super(`HTTP error: ${status}`); + this.name = 'InputAuthorizationError'; + } +} + /** * Bridges a `fetch` response body to a Node-stream Readable. The types at the * module boundary (Node's `stream/web` vs. lib.dom) don't overlap cleanly, @@ -993,11 +1001,18 @@ export class Job { submissionDir: this.submissionDir, identity: this.jobIdentity, }; + const startedAt = performance.now(); + let started = 0; + let completed = 0; + let cancelled = 0; let firstFailure: { error: unknown } | undefined; const runFileOperation = async (operation: () => Promise): Promise => { + started++; try { await operation(); + completed++; } catch (error) { + if (firstFailure) cancelled++; if (!firstFailure) { firstFailure = { error }; controller.abort(error); @@ -1031,6 +1046,15 @@ export class Job { Array.from({ length: workerCount }, () => runPrimeWorker()), ); if (firstFailure) { + this.log.error({ + inputCount: fileOps.length, + completed, + failed: 1, + cancelled, + notStarted: fileOps.length - started, + durationMs: Math.round(performance.now() - startedAt), + err: firstFailure.error, + }, 'Input preparation batch failed'); if (this.session) { /* A sibling may already have atomically replaced its destination. The * workspace now matches neither the previous checkpoint nor the full @@ -1302,6 +1326,9 @@ export class Job { if (!response.ok) { await response.body?.cancel().catch(() => {}); + if (response.status === 401 || response.status === 403) { + throw new InputAuthorizationError(response.status); + } throw new Error(`HTTP error: ${response.status}`); } @@ -1366,11 +1393,10 @@ export class Job { try { await fsp.unlink(tempPath); } catch { /* may not exist */ } throw abortReason(operation.signal); } - /* ValidationError is deterministic — a bad Content-Disposition - * filename will fail identically on every retry. Abort fast - * (cleanup + rethrow) instead of burning ~7.5s on exponential - * backoff and surfacing the error as a generic download failure. */ - if (error instanceof ValidationError) { + /* Invalid filenames and authorization denials cannot recover by + * replaying the same request. Abort the batch before exponential + * backoff amplifies the failure across its remaining files. */ + if (error instanceof ValidationError || error instanceof InputAuthorizationError) { try { await fsp.unlink(tempPath); } catch { /* may not exist */ } throw error; } @@ -1383,7 +1409,9 @@ export class Job { } } - this.log.error({ fileId: file.id, maxRetries, err: lastError }, 'Failed to download file'); + if (!context?.signal) { + this.log.error({ fileId: file.id, maxRetries, err: lastError }, 'Failed to download file'); + } try { await fsp.unlink(tempPath); } catch { /* may not exist */ } throw lastError ?? new Error(`Failed to download input ${file.id}`); } From 794df9e444acc2a8a2c6a0fbb65ddc0976c8ef41 Mon Sep 17 00:00:00 2001 From: Danny Avila Date: Fri, 11 Sep 2026 15:47:10 -0400 Subject: [PATCH 2/2] fix: preserve retries for transient egress ledger conflicts (#179) * fix: distinguish retryable ledger contention from scope denials * fix: preserve error classification through marker discovery and relay * test: exercise classified denials through gateway configuration * fix: honor bounded gateway retry hints during object downloads --- api/src/download.test.ts | 79 +++++++++++++++++++++++++++++- api/src/egress.ts | 1 + api/src/job.ts | 23 +++++++-- packages/code/src/relay.test.ts | 35 +++++++++++++ packages/code/src/relay.ts | 6 +++ service/src/egress-gateway.test.ts | 36 +++++++++++++- service/src/egress-gateway.ts | 4 ++ service/src/egress-grant.ts | 4 +- service/src/egress-ledger.ts | 2 +- 9 files changed, 182 insertions(+), 8 deletions(-) diff --git a/api/src/download.test.ts b/api/src/download.test.ts index 2c7248e8..bb8ec7e9 100644 --- a/api/src/download.test.ts +++ b/api/src/download.test.ts @@ -55,6 +55,7 @@ function makeRuntime(): Runtime { function makeJob(files: TFile[] = [], session?: SessionWorkspace): Job { return new Job({ session_id: 'test-session', + egress_grant: 'test-grant', runtime: makeRuntime(), files, args: [], @@ -440,10 +441,12 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => { }); it.each([401, 403])('does not retry an HTTP %i authorization denial', async status => { + config.egress_gateway_url = `http://127.0.0.1:${serverPort}`; const file: TFile = { id: 'denied', storage_session_id: 'previous', name: 'denied.txt' }; let requests = 0; routes.set('/sessions/previous/objects/denied', { status, + headers: { 'X-CodeAPI-Error-Code': 'scope_mismatch' }, onRequest: () => { requests++; }, }); const job = makeJob([file]); @@ -454,7 +457,8 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => { expect(await fsp.readdir(tmpDir)).toEqual([]); }); - it.each([404, 408, 429, 503])('still retries transient HTTP %i responses', async status => { + it.each([403, 404, 408, 429, 503])('still retries transient HTTP %i responses', async status => { + config.egress_gateway_url = `http://127.0.0.1:${serverPort}`; const file: TFile = { id: 'transient', storage_session_id: 'previous', name: 'ready.txt' }; let requests = 0; const route: Route = { @@ -471,7 +475,79 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => { expect(await fsp.readFile(path.join(tmpDir, 'ready.txt'), 'utf8')).toBe('ready'); }); + it.each([false, true])('honors conflict retry hints with cancellation=%s', async cancel => { + config.egress_gateway_url = `http://127.0.0.1:${serverPort}`; + const controller = new AbortController(); + const timestamps: number[] = []; + let timer: ReturnType | undefined; + const route: Route = { + status: 503, body: 'ready', + headers: { 'X-CodeAPI-Error-Code': 'ledger_conflict', 'Retry-After': '1' }, + onRequest: () => { + timestamps.push(performance.now()); + if (timestamps.length === 2) route.status = 200; + else if (cancel) timer = setTimeout(() => controller.abort(new Error('cancelled retry')), 25); + }, + }; + routes.set('/sessions/previous/objects/retry-hint', route); + const file: TFile = { id: 'retry-hint', storage_session_id: 'previous', name: 'ready.txt' }; + const job = makeJob([file]); + asInternals(job).submissionDir = tmpDir; + try { + const result = job.downloadAndWriteFile(file, 5, 1, { + submissionDir: tmpDir, identity: fallbackSandboxIdentity(), signal: controller.signal, + }); + if (cancel) { + await expect(result).rejects.toThrow('cancelled retry'); + expect(timestamps).toHaveLength(1); + expect(await fsp.readdir(tmpDir)).toEqual([]); + } else { + await expect(result).resolves.toBe('ready.txt'); + expect(timestamps).toHaveLength(2); + expect(timestamps[1] - timestamps[0]).toBeGreaterThanOrEqual(900); + } + } finally { + clearTimeout(timer); + } + }); + + it('does not retry an unclassified direct file-server denial', async () => { + config.egress_gateway_url = ''; + let requests = 0; + const file: TFile = { id: 'denied', storage_session_id: 'previous', name: 'denied.txt' }; + routes.set('/sessions/previous/objects/denied', { + status: 403, onRequest: () => { requests++; }, + }); + const job = makeJob([file]); + asInternals(job).submissionDir = tmpDir; + await expect(job.downloadAndWriteFile(file, 5, 1)).rejects.toThrow('HTTP error: 403'); + expect(requests).toBe(1); + }); + + it.each(['legacy', 'classified', 'direct'])('handles %s marker denials before priming', async mode => { + config.egress_gateway_url = mode === 'direct' ? '' : `http://127.0.0.1:${serverPort}`; + let requests = 0; + const route: Route = { + status: 403, body: '[]', + headers: mode === 'classified' ? { 'X-CodeAPI-Error-Code': 'scope_mismatch' } : {}, + onRequest: () => { if (++requests === 2) route.status = 200; }, + }; + routes.set('/sessions/previous/objects', route); + const file: TFile = { id: 'ready', storage_session_id: 'previous', name: 'ready.txt' }; + routes.set('/sessions/previous/objects/ready', { status: 200, body: 'ready' }); + const job = makeJob([file], sessionWorkspaceAt(tmpDir, 'marker-retry')); + if (mode === 'legacy') { + await job.prime(); + expect(requests).toBe(2); + expect(await fsp.readFile(path.join(tmpDir, 'ready.txt'), 'utf8')).toBe('ready'); + } else { + await expect(job.prime()).rejects.toThrow('HTTP error loading .dirkeep markers: 403'); + expect(requests).toBe(1); + } + }); + it('accounts for a denied 240-file batch once and stops queued downloads', async () => { + config.egress_gateway_url = `http://127.0.0.1:${serverPort}`; const files: TFile[] = Array.from({ length: 240 }, (_, index) => ({ id: `file-${index}`, storage_session_id: 'previous', name: `file-${index}.txt`, })); @@ -480,6 +556,7 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => { for (const file of files) { routes.set(`/sessions/previous/objects/${file.id}`, { status: 403, + headers: { 'X-CodeAPI-Error-Code': 'scope_mismatch' }, delayMs: file.id === 'file-0' ? 0 : 30, onRequest: () => { requests++; }, }); diff --git a/api/src/egress.ts b/api/src/egress.ts index 442872ab..69ecf6db 100644 --- a/api/src/egress.ts +++ b/api/src/egress.ts @@ -1 +1,2 @@ export const EGRESS_GRANT_HEADER = 'X-CodeAPI-Egress-Grant'; +export const EGRESS_ERROR_CODE_HEADER = 'X-CodeAPI-Error-Code'; diff --git a/api/src/job.ts b/api/src/job.ts index 8121f51d..bdcd5061 100644 --- a/api/src/job.ts +++ b/api/src/job.ts @@ -15,7 +15,7 @@ import { getRuntimes } from './runtime'; import { execute } from './nsjail'; import { config } from './config'; import { internalServiceHeaders } from './internal-service-auth'; -import { EGRESS_GRANT_HEADER } from './egress'; +import { EGRESS_GRANT_HEADER, EGRESS_ERROR_CODE_HEADER } from './egress'; import { injectTraceHeaders } from './telemetry'; import { applyReadOnlyInputPermissions, @@ -1196,6 +1196,11 @@ export class Job { } } + private isLegacyGatewayDenial(response: Response): boolean { + return !!config.egress_gateway_url && response.status === 403 && + !response.headers.has(EGRESS_ERROR_CODE_HEADER); + } + /** * Fetches normalized objects for one inherited session and returns the * `.dirkeep` markers belonging to exactly that session. Guards against: @@ -1219,7 +1224,7 @@ export class Job { signal: controller.signal, }, ); - if (res.status === 503 && attempt < AUTO_LOAD_DIRKEEP_RETRIES) { + if ((res.status === 503 || this.isLegacyGatewayDenial(res)) && attempt < AUTO_LOAD_DIRKEEP_RETRIES) { await res.body?.cancel().catch(() => {}); const retryAfterSeconds = Number(res.headers.get('retry-after')); await sleep( @@ -1326,7 +1331,10 @@ export class Job { if (!response.ok) { await response.body?.cancel().catch(() => {}); - if (response.status === 401 || response.status === 403) { + /* Older gateways also used 403 for transient ledger contention. + * Only classify 403 as permanent when the gateway distinguishes it. */ + if (response.status === 401 || + (response.status === 403 && !this.isLegacyGatewayDenial(response))) { throw new InputAuthorizationError(response.status); } throw new Error(`HTTP error: ${response.status}`); @@ -1402,7 +1410,14 @@ export class Job { } lastError = error instanceof Error ? error : new Error(String(error)); if (attempt < maxRetries) { - const delay = retryDelay * Math.pow(2, attempt - 1); + const backoff = retryDelay * Math.pow(2, attempt - 1); + const retryAfterSeconds = response?.status === 503 + ? Number(response.headers.get('retry-after')) : NaN; + /* Use the same bounded retry hint as marker discovery, without + * shortening exponential backoff or bypassing batch cancellation. */ + const delay = Number.isFinite(retryAfterSeconds) + ? Math.max(backoff, Math.min(1000, Math.max(0, retryAfterSeconds * 1000))) + : backoff; this.log.warn({ fileId: file.id, attempt, maxRetries, delay, err: lastError }, 'Download failed, retrying'); await sleep(delay, operation.signal); } diff --git a/packages/code/src/relay.test.ts b/packages/code/src/relay.test.ts index c0814408..a9b81fd0 100644 --- a/packages/code/src/relay.test.ts +++ b/packages/code/src/relay.test.ts @@ -381,3 +381,38 @@ test('file relay rejects plaintext remote upstreams', async () => { /HTTPS unless it is a local development host/, ); }); + +for (const [status, reason] of [[403, 'scope_mismatch'], [503, 'ledger_conflict']] as const) { + test(`file relay preserves ${reason} classification`, async () => { + const upstream = createServer((_req, res) => { + res.writeHead(status, { + 'X-CodeAPI-Error-Code': reason, + 'Retry-After': '1', + 'X-Internal-Secret': 'must-not-forward', + }).end('rejected'); + }); + const upstreamUrl = await listen(upstream); + const relay = await startFileRelay({ + host: '127.0.0.1', port: 0, upstreamUrl, token: 'relay-secret', + maxBytes: 1024, timeoutMs: 1000, + }); + try { + const response = await fetch(`${relay.url}/sessions/storage-1/objects/file-1`, { + headers: { + 'X-LibreChat-Code-Relay-Token': 'relay-secret', + 'X-CodeAPI-Egress-Grant': 'grant-1', + }, + }); + assert.equal(response.status, status); + assert.equal(response.headers.get('x-codeapi-error-code'), reason); + assert.equal(response.headers.get('retry-after'), '1'); + assert.equal(response.headers.get('x-internal-secret'), null); + assert.equal(await response.text(), 'rejected'); + } finally { + await relay.close(); + await new Promise((resolve, reject) => + upstream.close(error => error ? reject(error) : resolve()), + ); + } + }); +} diff --git a/packages/code/src/relay.ts b/packages/code/src/relay.ts index 255845e9..14e91280 100644 --- a/packages/code/src/relay.ts +++ b/packages/code/src/relay.ts @@ -240,6 +240,12 @@ export async function startFileRelay( )!, } : {}), + ...(upstreamResponse.headers.has('x-codeapi-error-code') + ? { 'X-CodeAPI-Error-Code': upstreamResponse.headers.get('x-codeapi-error-code')! } + : {}), + ...(upstreamResponse.headers.has('retry-after') + ? { 'Retry-After': upstreamResponse.headers.get('retry-after')! } + : {}), 'Content-Length': String(body.length), }); response.end(body); diff --git a/service/src/egress-gateway.test.ts b/service/src/egress-gateway.test.ts index 99303850..717309d7 100644 --- a/service/src/egress-gateway.test.ts +++ b/service/src/egress-gateway.test.ts @@ -1,6 +1,6 @@ process.env.CODEAPI_EGRESS_GATEWAY_AUTOSTART = 'false'; -import { afterAll, beforeAll, beforeEach, describe, expect, test } from 'bun:test'; +import { afterAll, beforeAll, beforeEach, describe, expect, test, spyOn } from 'bun:test'; import crypto from 'crypto'; import RedisMock from 'ioredis-mock'; import type { Server } from 'http'; @@ -435,6 +435,39 @@ describe('egress gateway routes', () => { } }); + test('reports exhausted ledger conflicts as retryable without forwarding the read', async () => { + const redis = new RedisMock(); + env.EGRESS_LEDGER_REQUIRED = true; + setEgressLedgerRedisForTest(redis as unknown as Parameters[0]); + const duplicate = redis.duplicate.bind(redis); + const duplication = spyOn(redis, 'duplicate').mockImplementation(() => { + const connection = duplicate(); + const transaction = { + set: () => transaction, + exec: async () => null, + }; + spyOn(connection, 'multi').mockImplementation(() => transaction as never); + return connection; + }); + try { + await createEgressLedger(claims()); + const readSession = sessionHandle({ dir: 'read', sessionId: 'sess_input' }); + const response = await gatewayFetch(`/sessions/${readSession}/objects?detail=normalized`, { + headers: grantHeader(), + }); + expect(response.status).toBe(503); + expect(response.headers.get('X-CodeAPI-Error-Code')).toBe('ledger_conflict'); + expect(response.headers.get('Retry-After')).toBe('1'); + expect(upstreamCalls).toHaveLength(0); + expect((await assertEgressGrantActive(claims())).request_count).toBe(0); + } finally { + duplication.mockRestore(); + setEgressLedgerRedisForTest(null); + redis.disconnect(); + env.EGRESS_LEDGER_REQUIRED = false; + } + }); + test('lists only scoped objects and injects internal credentials', async () => { upstreamResponse = Response.json([ { id: 'file_123', name: 'inputs/data.csv', storage_session_id: 'sess_input' }, @@ -557,6 +590,7 @@ describe('egress gateway routes', () => { }); expect(response.status).toBe(403); + expect(response.headers.get('X-CodeAPI-Error-Code')).toBe('scope_mismatch'); expect(upstreamCalls).toHaveLength(0); } finally { await redis.disconnect(); diff --git a/service/src/egress-gateway.ts b/service/src/egress-gateway.ts index d4d3e713..13a6de95 100644 --- a/service/src/egress-gateway.ts +++ b/service/src/egress-gateway.ts @@ -6,6 +6,7 @@ import { Readable } from 'stream'; import { env } from './config'; import { EGRESS_GRANT_HEADER, + EGRESS_ERROR_CODE_HEADER, EgressGrantError, egressGrantFromExecutionClaims, openEgressGrant, @@ -157,6 +158,7 @@ app.use((req: Request, res: Response, next: NextFunction) => { function errorStatus(error: EgressGrantError): number { if (error.reason === 'missing_secret' || error.reason === 'weak_secret') return 500; + if (error.reason === 'ledger_conflict') return 503; if (error.reason === 'malformed') return 400; if (error.reason === 'expired') return 401; return 403; @@ -165,6 +167,8 @@ function errorStatus(error: EgressGrantError): number { function sendEgressError(req: Request, res: Response, error: unknown): Response { if (error instanceof EgressGrantError) { const statusCode = errorStatus(error); + res.setHeader(EGRESS_ERROR_CODE_HEADER, error.reason); + if (error.reason === 'ledger_conflict') res.setHeader('Retry-After', '1'); logger.warn('Rejected egress gateway request', { requestId: requestId(res), reason: error.reason, diff --git a/service/src/egress-grant.ts b/service/src/egress-grant.ts index 8d8b1352..8148bec7 100644 --- a/service/src/egress-grant.ts +++ b/service/src/egress-grant.ts @@ -3,6 +3,7 @@ import type { ExecutionManifestClaims, ExecutionManifestInputFile } from './exec import type * as t from './types'; export const EGRESS_GRANT_HEADER = 'X-CodeAPI-Egress-Grant'; +export const EGRESS_ERROR_CODE_HEADER = 'X-CodeAPI-Error-Code'; export const EGRESS_GRANT_VERSION = 1; const TOKEN_PREFIX = 'ceg1'; @@ -19,7 +20,8 @@ export type EgressGrantErrorReason = | 'malformed' | 'expired' | 'wrong_type' - | 'scope_mismatch'; + | 'scope_mismatch' + | 'ledger_conflict'; export class EgressGrantError extends Error { readonly reason: EgressGrantErrorReason; diff --git a/service/src/egress-ledger.ts b/service/src/egress-ledger.ts index 24dda87c..b6bbf4fb 100644 --- a/service/src/egress-ledger.ts +++ b/service/src/egress-ledger.ts @@ -264,7 +264,7 @@ async function mutateRecord( }); releaseMutationConnection(client); } - throw new EgressGrantError('scope_mismatch', 'Egress grant ledger update conflicted'); + throw new EgressGrantError('ledger_conflict', 'Egress grant ledger update conflicted'); } export async function assertEgressGrantActive(grant: EgressGrantClaims): Promise {