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
147 changes: 146 additions & 1 deletion api/src/download.test.ts
Original file line number Diff line number Diff line change
@@ -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';
Expand Down Expand Up @@ -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: [],
Expand Down Expand Up @@ -439,6 +440,150 @@ describe('downloadAndWriteFile / RFC 5987 round-trip', () => {
expect(contents).toBe('hi');
});

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]);
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([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 = {
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.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<typeof setTimeout> | 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`,
}));
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,
headers: { 'X-CodeAPI-Error-Code': 'scope_mismatch' },
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',
Expand Down
1 change: 1 addition & 0 deletions api/src/egress.ts
Original file line number Diff line number Diff line change
@@ -1 +1,2 @@
export const EGRESS_GRANT_HEADER = 'X-CodeAPI-Egress-Grant';
export const EGRESS_ERROR_CODE_HEADER = 'X-CodeAPI-Error-Code';
61 changes: 52 additions & 9 deletions api/src/job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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<void>): Promise<void> => {
started++;
try {
await operation();
completed++;
} catch (error) {
if (firstFailure) cancelled++;
if (!firstFailure) {
firstFailure = { error };
controller.abort(error);
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -1172,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:
Expand All @@ -1195,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(
Expand Down Expand Up @@ -1302,6 +1331,12 @@ export class Job {

if (!response.ok) {
await response.body?.cancel().catch(() => {});
/* 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}`);
}

Expand Down Expand Up @@ -1366,24 +1401,32 @@ 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;
}
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);
}
}
}

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}`);
}
Expand Down
35 changes: 35 additions & 0 deletions packages/code/src/relay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>((resolve, reject) =>
upstream.close(error => error ? reject(error) : resolve()),
);
}
});
}
6 changes: 6 additions & 0 deletions packages/code/src/relay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading