Skip to content
Closed
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
77 changes: 77 additions & 0 deletions apps/web/src/lib/integrations/core/repository-read-limits.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
import { boundRepositoryResponse, withRepositoryReadDeadline } from './repository-read-limits';

const jsonHeaders = { 'content-type': 'application/json' };

afterEach(() => jest.useRealTimers());

describe('repository response bounds', () => {
it('accepts exactly 1 MiB without changing the decoded data', async () => {
const body = JSON.stringify('x'.repeat(1024 * 1024 - 2));
const response = await boundRepositoryResponse(new Response(body, { headers: jsonHeaders }));
expect(await response.json()).toHaveLength(1024 * 1024 - 2);
});

it.each(['stream', 'advertised', 'invalid length', 'content type', 'invalid bytes'])(
'rejects and cancels %s before parsing',
async failure => {
let cancelled = false;
const stream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(
failure === 'invalid bytes' ? new Uint8Array([255]) : new Uint8Array(1024 * 1024 + 1)
);
if (failure === 'invalid bytes') controller.close();
},
cancel() {
cancelled = true;
},
});
const headers = {
...jsonHeaders,
...(failure === 'advertised'
? { 'content-length': '1048577' }
: failure === 'invalid length'
? { 'content-length': 'invalid' }
: {}),
};
if (failure === 'content type') headers['content-type'] = 'text/html';
await expect(boundRepositoryResponse(new Response(stream, { headers }))).rejects.toThrow();
if (failure !== 'invalid bytes') expect(cancelled).toBe(true);
}
);

it('cancels a stalled body at the operation deadline', async () => {
jest.useFakeTimers();
let cancelled = false;
const response = new Response(
new ReadableStream({
cancel() {
cancelled = true;
},
}),
{ headers: jsonHeaders }
);
const result = withRepositoryReadDeadline({ bounded: true }, signal =>
boundRepositoryResponse(response, signal)
);
const rejection = expect(result).rejects.toThrow('Repository fetch timed out');
await jest.advanceTimersByTimeAsync(30_000);
await rejection;
expect(cancelled).toBe(true);
});

it('propagates caller cancellation without waiting for the deadline', async () => {
const controller = new AbortController();
const result = withRepositoryReadDeadline(
{ bounded: true, signal: controller.signal },
async signal => {
await new Promise<void>(resolve =>
signal?.addEventListener('abort', () => resolve(), { once: true })
);
signal?.throwIfAborted();
}
);
controller.abort(new Error('Read cancelled'));
await expect(result).rejects.toThrow('Read cancelled');
});
});
102 changes: 102 additions & 0 deletions apps/web/src/lib/integrations/core/repository-read-limits.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
import { z } from 'zod';

export const REPOSITORY_READ_LIMITS = {
pages: 2,
repositories: 50,
responseBytes: 1024 * 1024,
timeoutMs: 30_000,
} as const;

// Omitted options preserve complete legacy reads until those callers retire.
// Bounded results can be incomplete and must not replace a complete shared cache.
export type RepositoryReadOptions = { bounded?: boolean; signal?: AbortSignal };

export const repositoryPageSchema = z.custom<unknown[]>(
value => Array.isArray(value) && value.length <= REPOSITORY_READ_LIMITS.repositories,
'Invalid repository page'
);

export async function withRepositoryReadDeadline<T>(
options: RepositoryReadOptions | undefined,
read: (signal?: AbortSignal) => Promise<T>
): Promise<T> {
if (!options?.bounded) return read();
const controller = new AbortController();
const signal = options.signal
? AbortSignal.any([controller.signal, options.signal])
: controller.signal;
const timer = setTimeout(
() => controller.abort(new Error('Repository fetch timed out')),
REPOSITORY_READ_LIMITS.timeoutMs
);
const aborted = Promise.withResolvers<never>();
const onAbort = () => aborted.reject(signal.reason);
signal.addEventListener('abort', onAbort, { once: true });
try {
signal.throwIfAborted();
return await Promise.race([read(signal), aborted.promise]);
} finally {
clearTimeout(timer);
signal.removeEventListener('abort', onAbort);
}
}

export async function boundRepositoryResponse(
response: Response,
signal?: AbortSignal
): Promise<Response> {
if (!response.body) return response;
const reader = response.body.getReader();
const cancel = () => {
void reader.cancel(signal?.reason).catch(() => {});
};
signal?.addEventListener('abort', cancel, { once: true });
try {
signal?.throwIfAborted();
const length = response.headers.get('content-length');
if (
length &&
(!/^\d+$/.test(length) || Number(length) > REPOSITORY_READ_LIMITS.responseBytes)
) {
throw new Error('Repository response exceeded size limit');
}
if (
response.ok &&
response.headers.get('content-type')?.split(';')[0].trim().toLowerCase() !==
'application/json'
) {
throw new Error('Invalid repository response content type');
}
const bytes = new Uint8Array(REPOSITORY_READ_LIMITS.responseBytes);
let size = 0;
while (true) {
const chunk = await reader.read();
signal?.throwIfAborted();
if (chunk.done) break;
if (!(chunk.value instanceof Uint8Array)) throw new Error('Invalid repository response');
if (size + chunk.value.byteLength > bytes.byteLength) {
throw new Error('Repository response exceeded size limit');
}
bytes.set(chunk.value, size);
size += chunk.value.byteLength;
}
const text = new TextDecoder('utf-8', { fatal: true }).decode(bytes.subarray(0, size));
return new Response(text, response);
} catch (error) {
cancel();
throw error;
} finally {
signal?.removeEventListener('abort', cancel);
reader.releaseLock();
}
}

export function boundedRepositoryFetch(signal: AbortSignal): typeof fetch {
return async (input, init) => {
signal.throwIfAborted();
return boundRepositoryResponse(
await fetch(input, { ...init, signal, redirect: 'error' }),
signal
);
};
}
Original file line number Diff line number Diff line change
@@ -1,5 +1,125 @@
import { describe, expect, it } from '@jest/globals';
import { BitbucketRepositoryListResultSchema } from './token-service-client';
jest.mock('@/lib/config.server', () => ({
GIT_TOKEN_SERVICE_API_URL: 'https://token-service.example',
}));
jest.mock('@/lib/tokens', () => ({
BITBUCKET_REPOSITORY_LIST_AUDIENCE: 'bitbucket-repository-list',
TOKEN_EXPIRY: { fiveMinutes: 300 },
generateInternalServiceToken: () => 'service-token',
}));

import {
BitbucketRepositoryListResultSchema,
fetchBitbucketRepositoriesFromTokenService,
fetchBitbucketWorkspaceAccessTokenRepositoriesFromTokenService,
} from './token-service-client';

const fetchMock = jest.spyOn(globalThis, 'fetch');
const repository = {
id: '12345678-1234-4234-8234-123456789012',
workspaceUuid: '12345678-1234-4234-8234-123456789013',
name: 'repo',
fullName: 'workspace/repo',
private: true,
defaultBranch: 'main',
};
afterAll(() => fetchMock.mockRestore());
afterEach(() => jest.useRealTimers());

it.each([
fetchBitbucketRepositoriesFromTokenService,
fetchBitbucketWorkspaceAccessTokenRepositoriesFromTokenService,
])('bounds both authentication paths without changing legacy results', async read => {
fetchMock.mockImplementation(async (_url, init) => {
if (new Headers(init?.headers).get('authorization') !== 'Bearer service-token')
throw new Error('Missing authentication');
return Response.json({ status: 'available', repositories: Array(51).fill(repository) });
});
await expect(read('user', 'organization', { bounded: true })).resolves.toEqual({
status: 'available',
repositories: Array(50).fill(repository),
});
await expect(read('user', 'organization')).resolves.toEqual({
status: 'available',
repositories: Array(51).fill(repository),
});
});

it.each([
{ status: 'available', repositories: [] },
{ status: 'reconnect_required' },
{ status: 'temporarily_unavailable' },
])('preserves explicit provider states %j', async data => {
fetchMock.mockResolvedValue(Response.json(data));
await expect(
fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true })
).resolves.toEqual(data);
});

it.each([
{ status: 'available', repositories: [...Array(50).fill(repository), null] },
'invalid json',
])(
'rejects malformed or oversized bounded data %# but preserves the legacy fallback',
async data => {
fetchMock.mockImplementation(async () =>
typeof data === 'string'
? new Response(data, { headers: { 'content-type': 'application/json' } })
: Response.json(data)
);
await expect(
fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true })
).rejects.toThrow();
await expect(fetchBitbucketRepositoriesFromTokenService('user')).resolves.toEqual({
status: 'temporarily_unavailable',
});
}
);

it('rejects valid oversized JSON but keeps the legacy response unchanged', async () => {
const data = {
status: 'available',
repositories: [{ ...repository, name: 'x'.repeat(1048577) }],
};
fetchMock.mockImplementation(async () => Response.json(data));
await expect(
fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true })
).rejects.toThrow('size limit');
await expect(fetchBitbucketRepositoriesFromTokenService('user')).resolves.toEqual(data);
});

it('keeps network failures retryable without treating them as empty data', async () => {
fetchMock.mockRejectedValue(new Error('network unavailable'));
await expect(fetchBitbucketRepositoriesFromTokenService('user')).resolves.toEqual({
status: 'temporarily_unavailable',
});
await expect(
fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true })
).rejects.toThrow('network unavailable');
fetchMock.mockResolvedValue(Response.json({ status: 'available', repositories: [] }));
await expect(
fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true })
).resolves.toEqual({ status: 'available', repositories: [] });
});

it('cancels token-service response consumption at the deadline', async () => {
jest.useFakeTimers();
let cancelled = false;
fetchMock.mockResolvedValue(
new Response(
new ReadableStream({
cancel() {
cancelled = true;
},
}),
{ headers: { 'content-type': 'application/json' } }
)
);
const result = fetchBitbucketRepositoriesFromTokenService('user', undefined, { bounded: true });
const rejection = expect(result).rejects.toThrow('Repository fetch timed out');
await jest.advanceTimersByTimeAsync(30_000);
await rejection;
expect(cancelled).toBe(true);
});

describe('BitbucketRepositoryListResultSchema', () => {
it.each(['insufficient_permissions', 'invalid_request'] as const)(
Expand Down
Loading
Loading