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
4 changes: 3 additions & 1 deletion apps/web/src/lib/providers/plugin-proxy-http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -71,7 +71,8 @@ export class PluginProxyHttpAdapter implements IHttpService {
const abortSignals = [options?.signal, controller?.signal].filter(
(signal): signal is AbortSignal => signal !== undefined
);
const signal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined;
const mergedSignal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined;
const signal = mergedSignal?.signal;

try {
const proxyRes = await fetch(`${base}/api/plugins/${pluginId}/${action}`, {
Expand All @@ -92,6 +93,7 @@ export class PluginProxyHttpAdapter implements IHttpService {
return buildProxyResponse(proxyRes, proxyData);
} finally {
if (timeoutId) clearTimeout(timeoutId);
mergedSignal?.dispose();
}
}
}
2 changes: 1 addition & 1 deletion apps/web/src/lib/providers/providers.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -395,7 +395,7 @@ describe('Web Providers', () => {
timeoutMs: 30,
signal: new AbortController().signal
})
).rejects.toMatchObject({ name: 'AbortError' });
).rejects.toMatchObject({ name: 'TimeoutError' });

expect(fetchMock).toHaveBeenCalledTimes(1);
expect(fetchMock.mock.calls[0]?.[1]?.signal).toBeDefined();
Expand Down
249 changes: 249 additions & 0 deletions apps/web/src/lib/providers/web-http.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,249 @@
import { afterEach, describe, expect, it, vi } from 'vite-plus/test';
import { getEventListeners } from 'node:events';
import { WebHttpProxyProvider } from './web-http';
import { OfficialPluginInstallQueue } from '$lib/services/official-plugins/install-queue';

vi.mock('$app/paths', () => ({ base: '' }));
vi.mock('$lib/boot/plugin-proxy-meta.generated', () => ({
deploymentHasServerPlugins: () => false
}));

afterEach(() => {
vi.useRealTimers();
vi.unstubAllGlobals();
});

function pendingBody() {
let body!: ReadableStreamDefaultController<Uint8Array>;
let signal!: AbortSignal;
const cancel = vi.fn();
vi.stubGlobal(
'fetch',
vi.fn(async (_url: string, init: RequestInit) => {
signal = init.signal!;
const stream = new ReadableStream<Uint8Array>({
start(controller) {
body = controller;
signal.addEventListener(
'abort',
() => controller.error(new DOMException('Stream aborted', 'AbortError')),
{ once: true }
);
},
cancel
});
return new Response(stream, { headers: { 'Content-Length': '3' } });
})
);
return {
get body() {
return body;
},
get signal() {
return signal;
},
cancel
};
}

describe('Web HTTP response bodies', () => {
it.each(['complete', 'body failure', 'request failure', 'empty body'] as const)(
'removes merged signal listeners after %s',
async (outcome) => {
vi.useFakeTimers();
const source = new AbortController();
vi.stubGlobal(
'fetch',
vi.fn(async () => {
if (outcome === 'request failure') throw new Error('Offline');
return new Response(outcome === 'empty body' ? null : 'complete');
})
);
const request = new WebHttpProxyProvider().request('https://themes.test/asset', {
signal: source.signal,
timeoutMs: 100
});
if (outcome === 'request failure') await expect(request).rejects.toThrow('Offline');
else {
const response = await request;
if (outcome === 'body failure') await expect(response.json()).rejects.toThrow();
else if (outcome === 'complete') expect(await response.text()).toBe('complete');
}
expect(getEventListeners(source.signal, 'abort')).toHaveLength(0);
expect(vi.getTimerCount()).toBe(0);
}
);

it('forwards an explicit cache reload to fetch', async () => {
const fetchMock = vi.fn(async () => new Response('fresh'));
vi.stubGlobal('fetch', fetchMock);
const response = await new WebHttpProxyProvider().request('https://themes.test/bundle.js', {
cache: 'reload'
});
expect(await response.text()).toBe('fresh');
expect(fetchMock).toHaveBeenCalledWith(
'https://themes.test/bundle.js',
expect.objectContaining({ cache: 'reload' })
);
});

it.each(['text', 'json', 'bytes', 'bytes-progress'] as const)(
'keeps the request deadline while reading %s',
async (method) => {
vi.useFakeTimers();
const pending = pendingBody();
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100
});
await vi.advanceTimersByTimeAsync(60);
const reading = method === 'bytes-progress' ? response.bytes(vi.fn()) : response[method]();
const rejected = expect(reading).rejects.toMatchObject({ name: 'TimeoutError' });
await vi.advanceTimersByTimeAsync(40);
expect(pending.signal.aborted).toBe(true);
await rejected;
expect(vi.getTimerCount()).toBe(0);
}
);

it('reports received bytes and releases the timeout after a complete body', async () => {
vi.useFakeTimers();
const pending = pendingBody();
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100
});
const progress = vi.fn();
const reading = response.bytes(progress);
pending.body.enqueue(new Uint8Array([1]));
await vi.advanceTimersByTimeAsync(10);
expect(progress).toHaveBeenCalledWith({ receivedBytes: 1, totalBytes: 3 });
pending.body.enqueue(new Uint8Array([2, 3]));
pending.body.close();
expect(await reading).toEqual(new Uint8Array([1, 2, 3]));
expect(progress).toHaveBeenLastCalledWith({ receivedBytes: 3, totalBytes: 3 });
expect(vi.getTimerCount()).toBe(0);
await vi.advanceTimersByTimeAsync(100);
expect(pending.signal.aborted).toBe(false);
});

it('cancels a body read through an external signal', async () => {
vi.useFakeTimers();
pendingBody();
const controller = new AbortController();
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100,
signal: controller.signal
});
const reading = response.bytes(vi.fn());
const rejected = expect(reading).rejects.toMatchObject({ name: 'AbortError' });
controller.abort();
await rejected;
expect(vi.getTimerCount()).toBe(0);
});

it.each(['text', 'json', 'bytes'] as const)(
'releases the deadline after a failed %s read',
async (method) => {
vi.useFakeTimers();
const pending = pendingBody();
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100
});
const reading = response[method]();
const rejected = expect(reading).rejects.toThrow('Disconnected');
pending.body.error(new Error('Disconnected'));
await rejected;
expect(vi.getTimerCount()).toBe(0);
}
);

it('does not use a compressed content length as the decoded byte total', async () => {
vi.stubGlobal(
'fetch',
vi.fn(
async () =>
new Response(new Uint8Array([1, 2, 3]), {
headers: { 'Content-Length': '1', 'Content-Encoding': 'gzip' }
})
)
);
const response = await new WebHttpProxyProvider().request('https://themes.test/image');
const progress = vi.fn();
expect(await response.bytes(progress)).toEqual(new Uint8Array([1, 2, 3]));
expect(progress).toHaveBeenLastCalledWith({ receivedBytes: 3 });
});

it('keeps a timeout as a retryable installation failure', async () => {
vi.useFakeTimers();
pendingBody();
const http = new WebHttpProxyProvider();
const queue = new OfficialPluginInstallQueue({
runner: async (_manifest, _url, options) => {
const response = await http.request('https://themes.test/image', {
timeoutMs: 100,
signal: options?.signal
});
await response.bytes(vi.fn());
}
});
queue.enqueue({
id: 'theme',
name: { en: 'Theme' },
description: { en: 'Theme' },
version: '1.0.0',
author: 'Test',
type: 'theme',
bundleFormat: 'esm'
});
await vi.advanceTimersByTimeAsync(100);
expect(queue.getTask('theme')).toMatchObject({ status: 'failed', error: 'Request timed out' });
vi.stubGlobal(
'fetch',
vi.fn(async () => new Response(new Uint8Array([1, 2, 3])))
);
queue.retry('theme');
await vi.advanceTimersByTimeAsync(0);
expect(queue.getTask('theme')).toBeUndefined();
expect(vi.getTimerCount()).toBe(0);
queue.dispose();
});

it('clears the deadline when a progress callback fails and cancels its reader', async () => {
vi.useFakeTimers();
const pending = pendingBody();
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100
});
const reading = response.bytes(() => {
throw new Error('Progress failed');
});
const rejected = expect(reading).rejects.toThrow('Progress failed');
pending.body.enqueue(new Uint8Array([1]));
await rejected;
expect(pending.cancel).toHaveBeenCalledOnce();
expect(vi.getTimerCount()).toBe(0);
});

it('clears the deadline for a response without a body', async () => {
vi.useFakeTimers();
vi.stubGlobal(
'fetch',
vi.fn(async () => new Response(null, { status: 204 }))
);
const response = await new WebHttpProxyProvider().request('https://themes.test/image', {
timeoutMs: 100
});
expect(vi.getTimerCount()).toBe(0);
expect(await response.bytes()).toEqual(new Uint8Array());
});

it('reports bytes when the response does not advertise a total', async () => {
vi.stubGlobal(
'fetch',
vi.fn(async () => new Response(new Uint8Array([1, 2, 3])))
);
const response = await new WebHttpProxyProvider().request('https://themes.test/image');
const progress = vi.fn();
await response.bytes(progress);
expect(progress).toHaveBeenLastCalledWith({ receivedBytes: 3 });
});
});
Loading
Loading