From 6fb8c433a82bc7f0f74adffc494424a6311b8432 Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:10:49 +0800 Subject: [PATCH 1/8] =?UTF-8?q?=E2=9C=A8=20=E4=B8=BA=20HTTP=20=E5=93=8D?= =?UTF-8?q?=E5=BA=94=E5=A2=9E=E5=8A=A0=E4=B8=8B=E8=BD=BD=E8=BF=9B=E5=BA=A6?= =?UTF-8?q?=E6=8E=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- packages/core/src/types/services.ts | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/packages/core/src/types/services.ts b/packages/core/src/types/services.ts index fc78a0fe..9234249c 100644 --- a/packages/core/src/types/services.ts +++ b/packages/core/src/types/services.ts @@ -34,7 +34,13 @@ export interface HttpResponse { ok: boolean; text(): Promise; json(): Promise; - bytes(): Promise; + /** Reports decoded bytes while reading; totalBytes is omitted when unknown. */ + bytes(onProgress?: (progress: HttpDownloadProgress) => void): Promise; +} + +export interface HttpDownloadProgress { + receivedBytes: number; + totalBytes?: number; } /** Request options for a short-lived, cookie-backed upstream session. */ From 2e098f5c672aceceb0a275ebcabd8e68e0e3d2ed Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:11:07 +0800 Subject: [PATCH 2/8] =?UTF-8?q?=F0=9F=90=9B=20=E4=BF=AE=E5=A4=8D=E5=93=8D?= =?UTF-8?q?=E5=BA=94=E4=BD=93=E8=B6=85=E6=97=B6=E5=B9=B6=E6=94=AF=E6=8C=81?= =?UTF-8?q?=E4=B8=8B=E8=BD=BD=E8=BF=9B=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- apps/web/src/lib/providers/providers.test.ts | 2 +- apps/web/src/lib/providers/web-http.test.ts | 208 +++++++++++++++++++ apps/web/src/lib/providers/web-http.ts | 100 +++++++-- 3 files changed, 292 insertions(+), 18 deletions(-) create mode 100644 apps/web/src/lib/providers/web-http.test.ts diff --git a/apps/web/src/lib/providers/providers.test.ts b/apps/web/src/lib/providers/providers.test.ts index 3d10824c..0063ebea 100644 --- a/apps/web/src/lib/providers/providers.test.ts +++ b/apps/web/src/lib/providers/providers.test.ts @@ -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(); diff --git a/apps/web/src/lib/providers/web-http.test.ts b/apps/web/src/lib/providers/web-http.test.ts new file mode 100644 index 00000000..ff89f167 --- /dev/null +++ b/apps/web/src/lib/providers/web-http.test.ts @@ -0,0 +1,208 @@ +import { afterEach, describe, expect, it, vi } from 'vite-plus/test'; +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; + let signal!: AbortSignal; + const cancel = vi.fn(); + vi.stubGlobal( + 'fetch', + vi.fn(async (_url: string, init: RequestInit) => { + signal = init.signal!; + const stream = new ReadableStream({ + 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(['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 }); + }); +}); diff --git a/apps/web/src/lib/providers/web-http.ts b/apps/web/src/lib/providers/web-http.ts index 0d1e5786..9c5badab 100644 --- a/apps/web/src/lib/providers/web-http.ts +++ b/apps/web/src/lib/providers/web-http.ts @@ -1,5 +1,10 @@ import { base } from '$app/paths'; -import type { HttpRequestOptions, HttpResponse, IHttpService } from '@chronos/core'; +import type { + HttpDownloadProgress, + HttpRequestOptions, + HttpResponse, + IHttpService +} from '@chronos/core'; import { mergeAbortSignals } from '$lib/utils/abort-signal'; import { deploymentHasServerPlugins } from '$lib/boot/plugin-proxy-meta.generated'; @@ -49,6 +54,53 @@ function isDomainAllowed(hostname: string, allowedDomains: string[]): boolean { }); } +async function readResponseBytes( + response: Response, + onProgress?: (progress: HttpDownloadProgress) => void +): Promise { + if (!onProgress || !response.body) { + const bytes = new Uint8Array(await response.arrayBuffer()); + onProgress?.({ receivedBytes: bytes.length, totalBytes: bytes.length }); + return bytes; + } + const length = response.headers.get('Content-Length'); + const encoding = response.headers.get('Content-Encoding')?.trim().toLowerCase(); + const parsedLength = length === null ? undefined : Number(length); + let totalBytes = + (!encoding || encoding === 'identity') && + parsedLength !== undefined && + Number.isSafeInteger(parsedLength) && + parsedLength >= 0 + ? parsedLength + : undefined; + const reader = response.body.getReader(); + const chunks: Uint8Array[] = []; + let receivedBytes = 0; + try { + while (true) { + const { done, value } = await reader.read(); + if (done) break; + chunks.push(value); + receivedBytes += value.length; + // CORS may hide Content-Encoding, or the server may send an inaccurate length. + if (totalBytes !== undefined && receivedBytes > totalBytes) totalBytes = undefined; + onProgress({ receivedBytes, ...(totalBytes === undefined ? {} : { totalBytes }) }); + } + const bytes = new Uint8Array(receivedBytes); + let offset = 0; + for (const chunk of chunks) { + bytes.set(chunk, offset); + offset += chunk.length; + } + return bytes; + } catch (error) { + await reader.cancel().catch(() => {}); + throw error; + } finally { + reader.releaseLock(); + } +} + /** * WebHttpProxyProvider implements IHttpService for Web environments. * It manages direct fetches, CORS bypass proxy routing with whitelist validation, @@ -65,8 +117,18 @@ export class WebHttpProxyProvider implements IHttpService { const controller = options?.timeoutMs ? new AbortController() : undefined; const timeoutId = options?.timeoutMs && controller - ? setTimeout(() => controller.abort(), options.timeoutMs) + ? setTimeout( + () => controller.abort(new DOMException('Request timed out', 'TimeoutError')), + options.timeoutMs + ) : undefined; + const clearRequestTimeout = () => { + if (timeoutId !== undefined) clearTimeout(timeoutId); + }; + const abortSignals = [options?.signal, controller?.signal].filter( + (signal): signal is AbortSignal => signal !== undefined + ); + const requestSignal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined; try { if (options?.bypassCors && !deploymentHasServerPlugins()) { @@ -109,11 +171,6 @@ export class WebHttpProxyProvider implements IHttpService { } } - const abortSignals = [options?.signal, controller?.signal].filter( - (signal): signal is AbortSignal => signal !== undefined - ); - const requestSignal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined; - let requestUrl = url; if (base && url.startsWith('/official-plugins/')) requestUrl = `${base}${url}`; else if (base && typeof window !== 'undefined') { @@ -137,23 +194,32 @@ export class WebHttpProxyProvider implements IHttpService { response.headers.forEach((val, key) => { responseHeaders[key] = val; }); + if (!response.body) clearRequestTimeout(); + // fetch resolves at headers; the same deadline must also cover body consumption. + const consume = async (read: () => Promise): Promise => { + try { + requestSignal?.throwIfAborted(); + return await read(); + } catch (error) { + // Browsers may reject a timed-out body with AbortError; preserve the cause. + throw requestSignal?.aborted ? requestSignal.reason : error; + } finally { + clearRequestTimeout(); + } + }; return { status: response.status, statusText: response.statusText, headers: responseHeaders, ok: response.ok, - text: () => response.text(), - json: () => response.json() as Promise, - bytes: async () => { - const buf = await response.arrayBuffer(); - return new Uint8Array(buf); - } + text: () => consume(() => response.text()), + json: () => consume(() => response.json() as Promise), + bytes: (onProgress) => consume(() => readResponseBytes(response, onProgress)) }; - } finally { - if (timeoutId) { - clearTimeout(timeoutId); - } + } catch (error) { + clearRequestTimeout(); + throw requestSignal?.aborted ? requestSignal.reason : error; } } } From 6caedbb84df3b01f6a1fde80de83a0756fb4d993 Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:11:19 +0800 Subject: [PATCH 3/8] =?UTF-8?q?=F0=9F=90=9B=20=E4=BF=AE=E6=AD=A3=E6=8F=92?= =?UTF-8?q?=E4=BB=B6=E4=B8=8B=E8=BD=BD=E9=98=B6=E6=AE=B5=E4=B8=8E=E5=AE=89?= =?UTF-8?q?=E8=A3=85=E8=BF=9B=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- .../official-plugins/asset-pipeline.test.ts | 16 ++-- .../official-plugins/asset-pipeline.ts | 74 ++++++++++++------- .../official-plugins/theme-wallpaper.test.ts | 42 ++++++++++- 3 files changed, 101 insertions(+), 31 deletions(-) diff --git a/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts b/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts index 0d27b5d3..8a48c5da 100644 --- a/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts +++ b/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts @@ -100,14 +100,20 @@ describe('OfficialPluginAssetPipeline', () => { sha256: expected }; - httpRequest - .mockResolvedValueOnce(httpResponse({ text: async () => 'stale-bytes' })) - .mockResolvedValueOnce(httpResponse({ text: async () => SAMPLE_BUNDLE })); - - const assets = await pipeline.download(manifest); + httpRequest.mockResolvedValueOnce(httpResponse({ text: async () => 'stale-bytes' })); + const progress: { stage: string; percent: number }[] = []; + httpRequest.mockImplementationOnce(async () => { + expect(progress.at(-1)?.stage).toBe('downloading'); + return httpResponse({ text: async () => SAMPLE_BUNDLE }); + }); + + const assets = await pipeline.download(manifest, undefined, { + onProgress: (p) => progress.push(p) + }); expect(assets.code).toBe(SAMPLE_BUNDLE); expect(httpRequest).toHaveBeenCalledTimes(2); expect(httpRequest.mock.calls[1]?.[0]).toContain('?v='); + expect(progress.every((p, i) => i === 0 || p.percent >= progress[i - 1]!.percent)).toBe(true); }); it('rejects css sha256 mismatch', async () => { diff --git a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts index de75a9d0..007009f7 100644 --- a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts +++ b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts @@ -28,6 +28,11 @@ export class OfficialPluginAssetPipeline { let colorsJson: string | null = null; let iconThemeJson: string | null = null; let cssCode: string | null = null; + let percent = 0; + const reportProgress: NonNullable = (progress) => { + percent = Math.max(percent, progress.percent); + options?.onProgress?.({ ...progress, percent }); + }; const tasks: Array<{ url: string; @@ -75,22 +80,15 @@ export class OfficialPluginAssetPipeline { for (let i = 0; i < tasks.length; i++) { options?.signal?.throwIfAborted?.(); const task = tasks[i]!; - const downloadPercent = Math.round((i / total) * 70); - options?.onProgress?.({ - stage: 'downloading', - percent: downloadPercent, - label: task.label - }); - const text = await this.downloadTextAsset( task.url, task.sha256, task.label, options?.signal, - (verifyPercent) => { - options?.onProgress?.({ - stage: 'verifying', - percent: Math.round(70 + (i / total) * 15 + (verifyPercent / 100) * (15 / total)), + (progress) => { + reportProgress({ + stage: progress.stage, + percent: Math.round(((i + progress.percent / 100) / total) * 70), label: task.label }); } @@ -100,23 +98,27 @@ export class OfficialPluginAssetPipeline { } options?.signal?.throwIfAborted?.(); - options?.onProgress?.({ - stage: 'verifying', - percent: 85 - }); - const wallpaper = await this.downloadThemeWallpaper( colorsJson, resolvedManifest.colorsUrl, - options?.signal + options?.signal, + options?.onProgress + ? (progress) => + reportProgress({ + ...progress, + percent: Math.round(70 + (progress.percent / 100) * 15) + }) + : undefined ); + reportProgress({ stage: 'verifying', percent: 85 }); return { code, colorsJson, iconThemeJson, cssCode, ...(wallpaper ? { wallpaper } : {}) }; } async downloadThemeWallpaper( colorsJson: string | null, colorsUrl?: string, - signal?: AbortSignal + signal?: AbortSignal, + onProgress?: AssetDownloadOptions['onProgress'] ): Promise { if (!colorsJson) return undefined; const raw = JSON.parse(colorsJson); @@ -124,20 +126,41 @@ export class OfficialPluginAssetPipeline { const asset = parseColorThemeJson(raw).wallpaper!; if (!colorsUrl) throw new Error('Theme wallpaper requires a colors URL'); const url = resolveManifestAssetUrl(colorsUrl, asset.url); + signal?.throwIfAborted(); + onProgress?.({ stage: 'downloading', percent: 0, label: 'wallpaper' }); const response = await this.engine.http.request(withIntegrityBust(url, asset.sha256), { method: 'GET', timeoutMs: 20_000, signal }); if (!response.ok) throw new Error('Failed to download theme wallpaper'); - const bytes = await response.bytes(); + let downloadPercent = 0; + const bytes = await response.bytes( + onProgress + ? ({ receivedBytes, totalBytes }) => { + downloadPercent = Math.max( + downloadPercent, + totalBytes ? Math.min(90, (receivedBytes / totalBytes) * 90) : 0 + ); + onProgress({ + stage: 'downloading', + percent: downloadPercent, + label: 'wallpaper' + }); + } + : undefined + ); signal?.throwIfAborted(); + onProgress?.({ stage: 'verifying', percent: 90, label: 'wallpaper' }); const hash = await this.engine.runtime.sha256(bytes); + signal?.throwIfAborted(); if (hash.toLowerCase() !== asset.sha256.toLowerCase()) throw new Error('Theme wallpaper integrity check failed'); const blob = new Blob([new Uint8Array(bytes)]); + onProgress?.({ stage: 'verifying', percent: 95, label: 'wallpaper' }); await validateImage(blob); signal?.throwIfAborted(); + onProgress?.({ stage: 'verifying', percent: 100, label: 'wallpaper' }); return blob; } @@ -146,7 +169,7 @@ export class OfficialPluginAssetPipeline { expectedSha256: string | undefined, label = 'asset', signal?: AbortSignal, - onVerifyProgress?: (percent: number) => void + onProgress?: AssetDownloadOptions['onProgress'] ): Promise { const requestUrl = withIntegrityBust(url, expectedSha256); try { @@ -156,7 +179,7 @@ export class OfficialPluginAssetPipeline { expectedSha256, label, signal, - onVerifyProgress + onProgress ); } catch (err) { if (!isIntegrityMismatch(err) || requestUrl === url) throw err; @@ -168,7 +191,7 @@ export class OfficialPluginAssetPipeline { expectedSha256, label, signal, - onVerifyProgress + onProgress ); } } @@ -179,9 +202,10 @@ export class OfficialPluginAssetPipeline { expectedSha256: string | undefined, label: string, signal?: AbortSignal, - onVerifyProgress?: (percent: number) => void + onProgress?: AssetDownloadOptions['onProgress'] ): Promise { signal?.throwIfAborted?.(); + onProgress?.({ stage: 'downloading', percent: 0, label }); const response = await this.engine.http.request(requestUrl, { method: 'GET', timeoutMs: 20_000, @@ -193,7 +217,7 @@ export class OfficialPluginAssetPipeline { const text = await response.text(); signal?.throwIfAborted?.(); if (expectedSha256) { - onVerifyProgress?.(50); + onProgress?.({ stage: 'verifying', percent: 85, label }); const hash = await this.engine.runtime.sha256(text); signal?.throwIfAborted?.(); if (hash.toLowerCase() !== expectedSha256.toLowerCase()) { @@ -201,8 +225,8 @@ export class OfficialPluginAssetPipeline { `Plugin ${label} integrity check failed. Expected ${expectedSha256}, got ${hash}` ); } - onVerifyProgress?.(100); } + onProgress?.({ stage: 'verifying', percent: 100, label }); return text; } } diff --git a/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts b/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts index 09736b28..ea03acda 100644 --- a/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts +++ b/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts @@ -1,5 +1,5 @@ import { afterEach, describe, expect, it, vi } from 'vite-plus/test'; -import type { ChronosEngine } from '@chronos/core'; +import type { ChronosEngine, HttpResponse, PluginManifest } from '@chronos/core'; import type { ImageRepository } from '$lib/storage/image-repository'; import { OfficialPluginAssetPipeline } from './asset-pipeline'; import { OfficialPluginRuntimeActivator } from './runtime-activator'; @@ -42,6 +42,46 @@ function setup() { return { pipeline, request, sha256, decode, revoke }; } describe('theme wallpaper assets', () => { + it('reports wallpaper body reads as downloading and keeps overall progress monotonic', async () => { + const { pipeline, request, sha256, decode } = setup(); + const progress: { stage: string; percent: number; label?: string }[] = []; + request.mockResolvedValueOnce({ ok: true, text: async () => colors }); + request.mockResolvedValueOnce({ ok: true, text: async () => '{}' }); + request.mockResolvedValueOnce({ + ok: true, + bytes: async (onProgress: Parameters[0]) => { + expect(progress.at(-1)).toMatchObject({ stage: 'downloading', label: 'wallpaper' }); + onProgress?.({ receivedBytes: 1, totalBytes: 3 }); + expect(progress.at(-1)).toMatchObject({ stage: 'downloading', label: 'wallpaper' }); + onProgress?.({ receivedBytes: 2 }); + onProgress?.({ receivedBytes: 3, totalBytes: 3 }); + return new Uint8Array([1, 2, 3]); + } + }); + sha256.mockImplementation(async (data: string | Uint8Array) => { + if (data instanceof Uint8Array) + expect(progress.at(-1)).toMatchObject({ stage: 'verifying', label: 'wallpaper' }); + return hash; + }); + decode.mockImplementation(async () => { + expect(progress.at(-1)).toMatchObject({ stage: 'verifying', label: 'wallpaper' }); + }); + const manifest = { + colorsUrl: 'https://themes.test/colors.json', + colorsSha256: hash, + iconThemeUrl: 'https://themes.test/icons.json', + iconThemeSha256: hash + } as PluginManifest; + const assets = await pipeline.download(manifest, undefined, { + onProgress: (p) => progress.push(p) + }); + expect(assets.wallpaper?.size).toBe(3); + expect(progress.at(-1)).toMatchObject({ stage: 'verifying', percent: 85 }); + expect(progress.every((p, i) => i === 0 || p.percent >= progress[i - 1]!.percent)).toBe(true); + expect( + progress.filter((p) => p.stage === 'downloading' && p.label === 'wallpaper').length + ).toBeGreaterThan(1); + }); it('resolves against colors JSON and verifies bytes before decoding', async () => { const { pipeline, request, decode, revoke } = setup(); const blob = await pipeline.downloadThemeWallpaper( From 1098178937d830b2ef69c1084a8b0175ee160a90 Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:28:28 +0800 Subject: [PATCH 4/8] =?UTF-8?q?=F0=9F=90=9B=20=E9=81=BF=E5=85=8D=E5=AE=89?= =?UTF-8?q?=E8=A3=85=E6=8F=90=E4=BA=A4=E5=90=8E=E5=9B=A0=E5=B9=BF=E6=92=AD?= =?UTF-8?q?=E5=A4=B1=E8=B4=A5=E8=AF=AF=E5=9B=9E=E6=BB=9A?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- .../official-plugins/installed-store.test.ts | 63 ++++++++++++++++++- .../official-plugins/installed-store.ts | 17 +++-- 2 files changed, 74 insertions(+), 6 deletions(-) diff --git a/apps/web/src/lib/services/official-plugins/installed-store.test.ts b/apps/web/src/lib/services/official-plugins/installed-store.test.ts index a7667bf8..70f9ae45 100644 --- a/apps/web/src/lib/services/official-plugins/installed-store.test.ts +++ b/apps/web/src/lib/services/official-plugins/installed-store.test.ts @@ -1,8 +1,12 @@ -import { describe, expect, it, vi, beforeEach } from 'vite-plus/test'; +import { describe, expect, it, vi, beforeEach, afterEach } from 'vite-plus/test'; import { ChronosEngine } from '@chronos/core'; import type { ChronosEnv } from '@chronos/core'; import { DEFAULT_USER_PREFERENCES } from '@chronos/core'; -import { OfficialPluginInstalledStore, parseInstallationState } from './installed-store'; +import { + emptyInstallationState, + OfficialPluginInstalledStore, + parseInstallationState +} from './installed-store'; function createMockEnv() { const kv = new Map(); @@ -50,6 +54,61 @@ describe('OfficialPluginInstalledStore', () => { await engine.init(); store = new OfficialPluginInstalledStore(engine); }); + afterEach(() => { + vi.unstubAllGlobals(); + vi.restoreAllMocks(); + }); + + it.each(['notification failure', 'dispose during commit'])( + 'preserves a successful installation after %s', + async (scenario) => { + const gate = Promise.withResolvers(); + let persisted = emptyInstallationState(); + let closed = false; + const postMessage = vi.fn(() => { + throw new DOMException( + closed ? 'Channel closed' : 'Notification failed', + 'InvalidStateError' + ); + }); + vi.stubGlobal('window', {}); + vi.stubGlobal( + 'BroadcastChannel', + class { + postMessage = postMessage; + close() { + closed = true; + } + } + ); + vi.spyOn(console, 'error').mockImplementation(() => {}); + const installing = new OfficialPluginInstalledStore(engine, { + read: async () => structuredClone(persisted), + transaction: async (change) => { + await gate.promise; + const state = structuredClone(persisted); + change(state); + persisted = state; + return state; + } + }); + const listener = vi.fn(); + installing.onChanged(listener); + const pending = installing.upsert({ + manifest: { id: 'theme' } as never, + origin: { kind: 'user' }, + installedAt: 1, + wallpaperAssetId: 'new-wallpaper' + }); + if (scenario === 'dispose during commit') installing.dispose(); + gate.resolve(); + await expect(pending).resolves.toBeUndefined(); + expect(persisted.records[0]?.wallpaperAssetId).toBe('new-wallpaper'); + expect(installing.find('theme')?.wallpaperAssetId).toBe('new-wallpaper'); + expect(listener).toHaveBeenCalledTimes(scenario === 'notification failure' ? 1 : 0); + installing.dispose(); + } + ); it('rejects invalid development state without overwriting it', async () => { const invalid = [{ obsolete: true }]; diff --git a/apps/web/src/lib/services/official-plugins/installed-store.ts b/apps/web/src/lib/services/official-plugins/installed-store.ts index ff439b9c..fbd40337 100644 --- a/apps/web/src/lib/services/official-plugins/installed-store.ts +++ b/apps/web/src/lib/services/official-plugins/installed-store.ts @@ -122,6 +122,14 @@ export class OfficialPluginInstalledStore { } } } + private broadcast(message: Record): void { + // Notifications happen after commit and must not turn a saved installation into a failure. + try { + this.channel?.postMessage(message); + } catch (error) { + console.error(error); + } + } async load(): Promise { this.state = await this.repository.read(); return this.state.records; @@ -182,7 +190,7 @@ export class OfficialPluginInstalledStore { } }); this.generation = host.buildId; - this.channel?.postMessage({ generation: host.buildId }); + this.broadcast({ generation: host.buildId }); } private async mutate(change: (state: PluginInstallationState) => void): Promise { this.state = await this.repository.transaction((state) => { @@ -195,7 +203,7 @@ export class OfficialPluginInstalledStore { state.revision++; delete state.prepared; }); - this.channel?.postMessage({ revision: this.state.revision }); + this.broadcast({ revision: this.state.revision }); this.notify(); } async markSeeded() { @@ -243,14 +251,14 @@ export class OfficialPluginInstalledStore { throw new Error('Update is already running in another window'); state.prepared = update; }); - this.channel?.postMessage({ transition: update.target.buildId }); + this.broadcast({ transition: update.target.buildId }); this.notify(); } async cancelPreparation(token: string) { this.state = await this.repository.transaction((state) => { if (state.prepared?.token === token) delete state.prepared; }); - this.channel?.postMessage({ cancelled: token }); + this.broadcast({ cancelled: token }); this.notify(); } async persist() { @@ -264,6 +272,7 @@ export class OfficialPluginInstalledStore { } dispose() { this.channel?.close(); + this.channel = undefined; this.listeners.clear(); } } From beb18cf72c63cfc3544b9df945d429b768a6e60f Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:31:29 +0800 Subject: [PATCH 5/8] =?UTF-8?q?=F0=9F=90=9B=20=E9=98=B2=E6=AD=A2=E5=8F=96?= =?UTF-8?q?=E6=B6=88=E5=90=8E=E7=9A=84=E6=97=A7=E5=AE=89=E8=A3=85=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E5=88=A0=E9=99=A4=E9=87=8D=E8=AF=95=E4=BB=BB=E5=8A=A1?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- .../official-plugins/install-queue.test.ts | 59 +++++++++++++++++++ .../official-plugins/install-queue.ts | 23 +++++++- 2 files changed, 80 insertions(+), 2 deletions(-) diff --git a/apps/web/src/lib/services/official-plugins/install-queue.test.ts b/apps/web/src/lib/services/official-plugins/install-queue.test.ts index 48123928..37fad753 100644 --- a/apps/web/src/lib/services/official-plugins/install-queue.test.ts +++ b/apps/web/src/lib/services/official-plugins/install-queue.test.ts @@ -15,6 +15,65 @@ function createManifest(id: string, name = id): PluginManifest { } describe('OfficialPluginInstallQueue', () => { + it.each(['enqueue', 'retry'] as const)( + 'preserves a new attempt started by %s before the canceled runner settles', + async (restart) => { + const first = Promise.withResolvers(); + const second = Promise.withResolvers(); + let reportOldProgress: NonNullable[2]>['onProgress']; + const runner = vi + .fn() + .mockImplementationOnce((_manifest, _url, options) => { + reportOldProgress = options?.onProgress; + return first.promise; + }) + .mockImplementationOnce(() => second.promise); + const onTaskCanceled = vi.fn(); + const onTaskCompleted = vi.fn(); + const queue = new OfficialPluginInstallQueue({ runner, onTaskCanceled, onTaskCompleted }); + const manifest = createManifest('theme'); + queue.enqueue(manifest); + queue.cancel(manifest.id); + if (restart === 'enqueue') queue.enqueue(manifest); + else queue.retry(manifest.id); + reportOldProgress?.({ stage: 'verifying', percent: 80 }); + expect(queue.getTask(manifest.id)?.status).toBe('queued'); + expect(queue.getActiveTask()).toBeUndefined(); + first.reject(new DOMException('Aborted', 'AbortError')); + await Promise.resolve(); + await Promise.resolve(); + expect(runner).toHaveBeenCalledTimes(2); + expect(queue.getTask(manifest.id)?.status).toBe('downloading'); + expect(onTaskCanceled).not.toHaveBeenCalled(); + reportOldProgress?.({ stage: 'installing', percent: 90 }); + expect(queue.getTask(manifest.id)?.progress.percent).toBe(0); + second.resolve(); + await Promise.resolve(); + await Promise.resolve(); + expect(onTaskCompleted).toHaveBeenCalledTimes(1); + expect(queue.getTask(manifest.id)).toBeUndefined(); + } + ); + + it('does not resurrect a canceled task removed before its runner fails', async () => { + const pending = Promise.withResolvers(); + const queue = new OfficialPluginInstallQueue({ runner: () => pending.promise }); + const statuses: string[] = []; + queue.onChanged(() => { + const task = queue.getTask('theme'); + if (task) statuses.push(task.status); + }); + queue.enqueue(createManifest('theme')); + queue.cancel('theme'); + queue.clearFinished('theme'); + statuses.length = 0; + pending.reject(new DOMException('Aborted', 'AbortError')); + await Promise.resolve(); + await Promise.resolve(); + expect(statuses).toEqual([]); + expect(queue.getTask('theme')).toBeUndefined(); + }); + it('starts the first enqueued task immediately and marks it active', async () => { let resolveInstall!: () => void; const runner = vi.fn().mockImplementation(() => { diff --git a/apps/web/src/lib/services/official-plugins/install-queue.ts b/apps/web/src/lib/services/official-plugins/install-queue.ts index 7a9b6922..5c687218 100644 --- a/apps/web/src/lib/services/official-plugins/install-queue.ts +++ b/apps/web/src/lib/services/official-plugins/install-queue.ts @@ -37,10 +37,12 @@ export type { */ export class OfficialPluginInstallQueue implements Disposable { private readonly tasks = new Map(); + private readonly attempts = new Map(); private readonly listeners = new Set<(change: { kind: InstallQueueChangeKind }) => void>(); private isProcessing = false; private activeTaskId: string | null = null; private activeController: AbortController | null = null; + private activeAttempt: symbol | null = null; private disposed = false; constructor(private readonly deps: OfficialPluginInstallQueueDeps) {} @@ -59,12 +61,14 @@ export class OfficialPluginInstallQueue implements Disposable { return; } this.tasks.set(manifest.id, requeueTask(existing, manifest, manifestUrl)); + this.attempts.set(manifest.id, Symbol()); this.notifyState(); void this.processNext(); return; } this.tasks.set(manifest.id, createQueuedTask(manifest, manifestUrl)); + this.attempts.set(manifest.id, Symbol()); this.notifyState(); void this.processNext(); } @@ -103,6 +107,7 @@ export class OfficialPluginInstallQueue implements Disposable { if (task.status !== 'failed' && task.status !== 'canceled') return; this.tasks.set(pluginId, requeueTask(task, task.manifest, task.manifestUrl)); + this.attempts.set(pluginId, Symbol()); this.notifyState(); void this.processNext(); } @@ -147,6 +152,7 @@ export class OfficialPluginInstallQueue implements Disposable { const task = this.tasks.get(pluginId); if (task && (task.status === 'completed' || task.status === 'canceled')) { this.tasks.delete(pluginId); + this.attempts.delete(pluginId); this.notifyState(); } return; @@ -156,6 +162,7 @@ export class OfficialPluginInstallQueue implements Disposable { for (const [id, task] of this.tasks.entries()) { if (task.status === 'completed' || task.status === 'canceled') { this.tasks.delete(id); + this.attempts.delete(id); changed = true; } } @@ -174,6 +181,7 @@ export class OfficialPluginInstallQueue implements Disposable { getActiveTask(): PluginInstallTask | undefined { if (!this.activeTaskId) return undefined; + if (this.attempts.get(this.activeTaskId) !== this.activeAttempt) return undefined; return this.tasks.get(this.activeTaskId); } @@ -228,6 +236,8 @@ export class OfficialPluginInstallQueue implements Disposable { this.activeController = new AbortController(); const controller = this.activeController; const activeId = nextTask.pluginId; + const attempt = this.attempts.get(activeId)!; + this.activeAttempt = attempt; this.tasks.set(activeId, markTaskDownloading(nextTask)); this.notifyState(); @@ -237,6 +247,7 @@ export class OfficialPluginInstallQueue implements Disposable { await this.deps.runner(nextTask.manifest, nextTask.manifestUrl, { signal: controller.signal, onProgress: (prog) => { + if (this.attempts.get(activeId) !== attempt) return; const current = this.tasks.get(activeId); if (!current || current.status === 'canceled' || controller.signal.aborted) return; this.tasks.set(activeId, updateTaskProgress(current, prog)); @@ -245,7 +256,12 @@ export class OfficialPluginInstallQueue implements Disposable { }); const current = this.tasks.get(activeId); - if (current && current.status !== 'canceled' && !controller.signal.aborted) { + if ( + this.attempts.get(activeId) === attempt && + current && + current.status !== 'canceled' && + !controller.signal.aborted + ) { const completed = markTaskCompleted(current); this.tasks.set(activeId, completed); this.notifyState(); @@ -253,7 +269,9 @@ export class OfficialPluginInstallQueue implements Disposable { this.clearFinished(activeId); } } catch (err: unknown) { - const current = this.tasks.get(activeId) ?? nextTask; + // A canceled runner may settle after the same plugin has been queued again. + if (this.attempts.get(activeId) !== attempt) return; + const current = this.tasks.get(activeId)!; const outcome = resolveTaskErrorOutcome(err, { signalAborted: controller.signal.aborted, currentStatus: current.status @@ -275,6 +293,7 @@ export class OfficialPluginInstallQueue implements Disposable { this.isProcessing = false; this.activeTaskId = null; this.activeController = null; + this.activeAttempt = null; void this.processNext(); } } From 26a3a21c260efc4f2a03a2880309ad0731342e76 Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:35:23 +0800 Subject: [PATCH 6/8] =?UTF-8?q?=F0=9F=90=9B=20=E6=A0=A1=E9=AA=8C=E5=A4=B1?= =?UTF-8?q?=E8=B4=A5=E9=87=8D=E8=AF=95=E6=97=B6=E5=88=B7=E6=96=B0=E6=8F=92?= =?UTF-8?q?=E4=BB=B6=E8=B5=84=E6=BA=90=E7=BC=93=E5=AD=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- apps/web/src/lib/providers/web-http.test.ts | 13 +++++++++++++ apps/web/src/lib/providers/web-http.ts | 1 + .../official-plugins/asset-pipeline.test.ts | 11 +++++++---- .../lib/services/official-plugins/asset-pipeline.ts | 12 +++++++----- packages/core/src/types/services.ts | 2 ++ 5 files changed, 30 insertions(+), 9 deletions(-) diff --git a/apps/web/src/lib/providers/web-http.test.ts b/apps/web/src/lib/providers/web-http.test.ts index ff89f167..e585763b 100644 --- a/apps/web/src/lib/providers/web-http.test.ts +++ b/apps/web/src/lib/providers/web-http.test.ts @@ -46,6 +46,19 @@ function pendingBody() { } describe('Web HTTP response bodies', () => { + 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) => { diff --git a/apps/web/src/lib/providers/web-http.ts b/apps/web/src/lib/providers/web-http.ts index 9c5badab..3a03906c 100644 --- a/apps/web/src/lib/providers/web-http.ts +++ b/apps/web/src/lib/providers/web-http.ts @@ -187,6 +187,7 @@ export class WebHttpProxyProvider implements IHttpService { method: options?.method ?? 'GET', headers, body, + cache: options?.cache, signal: requestSignal }); diff --git a/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts b/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts index 8a48c5da..df845f7f 100644 --- a/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts +++ b/apps/web/src/lib/services/official-plugins/asset-pipeline.test.ts @@ -86,7 +86,7 @@ describe('OfficialPluginAssetPipeline', () => { await expect(pipeline.download(manifest)).rejects.toThrow(/integrity check failed/); }); - it('retries once with an integrity-busted URL after sha256 mismatch', async () => { + it('reloads cached bytes once after sha256 mismatch', async () => { const expected = await engine.env.runtime.sha256(SAMPLE_BUNDLE); const manifest: PluginManifest = { id: 'test', @@ -100,11 +100,12 @@ describe('OfficialPluginAssetPipeline', () => { sha256: expected }; - httpRequest.mockResolvedValueOnce(httpResponse({ text: async () => 'stale-bytes' })); const progress: { stage: string; percent: number }[] = []; - httpRequest.mockImplementationOnce(async () => { + httpRequest.mockImplementation(async (_url, options) => { expect(progress.at(-1)?.stage).toBe('downloading'); - return httpResponse({ text: async () => SAMPLE_BUNDLE }); + return httpResponse({ + text: async () => (options?.cache === 'reload' ? SAMPLE_BUNDLE : 'stale-bytes') + }); }); const assets = await pipeline.download(manifest, undefined, { @@ -113,6 +114,8 @@ describe('OfficialPluginAssetPipeline', () => { expect(assets.code).toBe(SAMPLE_BUNDLE); expect(httpRequest).toHaveBeenCalledTimes(2); expect(httpRequest.mock.calls[1]?.[0]).toContain('?v='); + expect(httpRequest.mock.calls[0]?.[1]?.cache).toBeUndefined(); + expect(httpRequest.mock.calls[1]?.[1]?.cache).toBe('reload'); expect(progress.every((p, i) => i === 0 || p.percent >= progress[i - 1]!.percent)).toBe(true); }); diff --git a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts index 007009f7..745386d9 100644 --- a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts +++ b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts @@ -183,15 +183,15 @@ export class OfficialPluginAssetPipeline { ); } catch (err) { if (!isIntegrityMismatch(err) || requestUrl === url) throw err; - // Stale SW/runtime cache served old bytes (e.g. slow-network fallback): - // the busted key already differs per content version, retry once. + // Even a hash-specific URL may have cached corrupt bytes; refresh it once. return await this.fetchVerifiedText( requestUrl, url, expectedSha256, label, signal, - onProgress + onProgress, + 'reload' ); } } @@ -202,14 +202,16 @@ export class OfficialPluginAssetPipeline { expectedSha256: string | undefined, label: string, signal?: AbortSignal, - onProgress?: AssetDownloadOptions['onProgress'] + onProgress?: AssetDownloadOptions['onProgress'], + cache?: 'reload' ): Promise { signal?.throwIfAborted?.(); onProgress?.({ stage: 'downloading', percent: 0, label }); const response = await this.engine.http.request(requestUrl, { method: 'GET', timeoutMs: 20_000, - signal + signal, + cache }); if (!response.ok) { throw new Error(`Failed to download plugin ${label} from ${url}`); diff --git a/packages/core/src/types/services.ts b/packages/core/src/types/services.ts index 9234249c..43a73859 100644 --- a/packages/core/src/types/services.ts +++ b/packages/core/src/types/services.ts @@ -23,6 +23,8 @@ export interface HttpRequestOptions { headers?: Record; body?: string | Uint8Array; bypassCors?: boolean; + /** Overrides the HTTP cache policy, for example to reload after an integrity failure. */ + cache?: RequestCache; timeoutMs?: number; signal?: AbortSignal; } From e1cf515b00262412610f889246c306141a1c89da Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:38:37 +0800 Subject: [PATCH 7/8] =?UTF-8?q?=F0=9F=90=9B=20=E4=B8=BA=E6=8F=92=E4=BB=B6?= =?UTF-8?q?=E5=A3=81=E7=BA=B8=E8=A7=A3=E7=A0=81=E5=A2=9E=E5=8A=A0=E5=8F=96?= =?UTF-8?q?=E6=B6=88=E5=92=8C=E8=B6=85=E6=97=B6=E5=A4=84=E7=90=86?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- .../official-plugins/asset-pipeline.ts | 2 +- .../official-plugins/theme-wallpaper.test.ts | 39 ++++++++++++++++++- apps/web/src/lib/wallpaper/validate-image.ts | 26 +++++++++++-- 3 files changed, 62 insertions(+), 5 deletions(-) diff --git a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts index 745386d9..fe5b8c81 100644 --- a/apps/web/src/lib/services/official-plugins/asset-pipeline.ts +++ b/apps/web/src/lib/services/official-plugins/asset-pipeline.ts @@ -158,7 +158,7 @@ export class OfficialPluginAssetPipeline { throw new Error('Theme wallpaper integrity check failed'); const blob = new Blob([new Uint8Array(bytes)]); onProgress?.({ stage: 'verifying', percent: 95, label: 'wallpaper' }); - await validateImage(blob); + await validateImage(blob, signal); signal?.throwIfAborted(); onProgress?.({ stage: 'verifying', percent: 100, label: 'wallpaper' }); return blob; diff --git a/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts b/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts index ea03acda..caaf383a 100644 --- a/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts +++ b/apps/web/src/lib/services/official-plugins/theme-wallpaper.test.ts @@ -3,6 +3,7 @@ import type { ChronosEngine, HttpResponse, PluginManifest } from '@chronos/core' import type { ImageRepository } from '$lib/storage/image-repository'; import { OfficialPluginAssetPipeline } from './asset-pipeline'; import { OfficialPluginRuntimeActivator } from './runtime-activator'; +import { OfficialPluginInstallQueue, type PluginInstallRunner } from './install-queue'; const hash = 'a'.repeat(64); const colors = JSON.stringify({ id: 'image-theme', @@ -10,7 +11,10 @@ const colors = JSON.stringify({ variants: { light: { colors: {} }, dark: { colors: {} } }, wallpaper: { url: './image.png', sha256: hash } }); -afterEach(() => vi.unstubAllGlobals()); +afterEach(() => { + vi.unstubAllGlobals(); + vi.useRealTimers(); +}); function setup() { const request = vi .fn() @@ -42,6 +46,39 @@ function setup() { return { pipeline, request, sha256, decode, revoke }; } describe('theme wallpaper assets', () => { + it.each(['cancel', 'timeout'] as const)( + 'releases the install queue after %s during stalled image decoding', + async (outcome) => { + vi.useFakeTimers(); + const { pipeline, decode, revoke } = setup(); + const pendingDecode = Promise.withResolvers(); + decode.mockReturnValue(pendingDecode.promise); + const runner = vi.fn(async (manifest, _url, options) => { + if (manifest.id === 'theme') + await pipeline.downloadThemeWallpaper( + colors, + 'https://themes.test/colors.json', + options?.signal + ); + }); + const queue = new OfficialPluginInstallQueue({ runner }); + queue.enqueue({ id: 'theme' } as PluginManifest); + queue.enqueue({ id: 'next' } as PluginManifest); + await vi.advanceTimersByTimeAsync(0); + expect(decode).toHaveBeenCalledOnce(); + if (outcome === 'cancel') queue.cancel('theme'); + await vi.advanceTimersByTimeAsync(outcome === 'timeout' ? 20_000 : 0); + expect(runner).toHaveBeenCalledTimes(2); + expect(queue.getTask('next')).toBeUndefined(); + if (outcome === 'cancel') expect(queue.getTask('theme')).toBeUndefined(); + else expect(queue.getTask('theme')).toMatchObject({ status: 'failed' }); + expect(revoke).toHaveBeenCalledWith('blob:validate'); + expect(vi.getTimerCount()).toBe(0); + pendingDecode.reject(new Error('Late decode rejection')); + await vi.advanceTimersByTimeAsync(0); + } + ); + it('reports wallpaper body reads as downloading and keeps overall progress monotonic', async () => { const { pipeline, request, sha256, decode } = setup(); const progress: { stage: string; percent: number; label?: string }[] = []; diff --git a/apps/web/src/lib/wallpaper/validate-image.ts b/apps/web/src/lib/wallpaper/validate-image.ts index 5e89838f..720cfd07 100644 --- a/apps/web/src/lib/wallpaper/validate-image.ts +++ b/apps/web/src/lib/wallpaper/validate-image.ts @@ -1,12 +1,32 @@ /** Decode before installing an image so a valid hash cannot hide unusable bytes. */ -export async function validateImage(blob: Blob): Promise { +export async function validateImage(blob: Blob, signal?: AbortSignal): Promise { + signal?.throwIfAborted(); + const image = new Image(); const url = URL.createObjectURL(blob); + let timeout: ReturnType | undefined; + let onAbort: (() => void) | undefined; try { - const image = new Image(); image.src = url; - await image.decode(); + await new Promise((resolve, reject) => { + onAbort = () => reject(signal?.reason); + signal?.addEventListener('abort', onAbort, { once: true }); + if (signal?.aborted) { + onAbort(); + return; + } + timeout = setTimeout( + () => reject(new DOMException('Image decoding timed out', 'TimeoutError')), + 20_000 + ); + // Keep a rejection handler attached even after cancellation or timeout wins. + image.decode().then(resolve, reject); + }); + signal?.throwIfAborted(); if (!image.naturalWidth || !image.naturalHeight) throw new Error('Invalid wallpaper image'); } finally { + if (timeout !== undefined) clearTimeout(timeout); + if (onAbort) signal?.removeEventListener('abort', onAbort); + image.src = ''; URL.revokeObjectURL(url); } } From 0627538da6d0676d10bb263aeefd60090fa0517c Mon Sep 17 00:00:00 2001 From: UE-DND <100979820+UE-DND@users.noreply.github.com> Date: Thu, 1 Oct 2026 00:42:38 +0800 Subject: [PATCH 8/8] =?UTF-8?q?=F0=9F=90=9B=20=E5=9C=A8=E6=8F=92=E4=BB=B6?= =?UTF-8?q?=E5=AE=89=E8=A3=85=E5=92=8C=E8=AF=B7=E6=B1=82=E7=BB=93=E6=9D=9F?= =?UTF-8?q?=E5=90=8E=E6=B8=85=E7=90=86=E5=8F=96=E6=B6=88=E7=9B=91=E5=90=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-authored-by: Codex --- .../src/lib/providers/plugin-proxy-http.ts | 4 +- apps/web/src/lib/providers/web-http.test.ts | 28 +++++++++++ apps/web/src/lib/providers/web-http.ts | 12 +++-- .../official-plugin-service.ts | 13 +++-- apps/web/src/lib/utils/abort-signal.test.ts | 47 +++++++++++++++++++ apps/web/src/lib/utils/abort-signal.ts | 29 ++++++++---- 6 files changed, 114 insertions(+), 19 deletions(-) create mode 100644 apps/web/src/lib/utils/abort-signal.test.ts diff --git a/apps/web/src/lib/providers/plugin-proxy-http.ts b/apps/web/src/lib/providers/plugin-proxy-http.ts index e2337b29..7273220b 100644 --- a/apps/web/src/lib/providers/plugin-proxy-http.ts +++ b/apps/web/src/lib/providers/plugin-proxy-http.ts @@ -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}`, { @@ -92,6 +93,7 @@ export class PluginProxyHttpAdapter implements IHttpService { return buildProxyResponse(proxyRes, proxyData); } finally { if (timeoutId) clearTimeout(timeoutId); + mergedSignal?.dispose(); } } } diff --git a/apps/web/src/lib/providers/web-http.test.ts b/apps/web/src/lib/providers/web-http.test.ts index e585763b..ca65cfb6 100644 --- a/apps/web/src/lib/providers/web-http.test.ts +++ b/apps/web/src/lib/providers/web-http.test.ts @@ -1,4 +1,5 @@ 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'; @@ -46,6 +47,33 @@ function pendingBody() { } 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); diff --git a/apps/web/src/lib/providers/web-http.ts b/apps/web/src/lib/providers/web-http.ts index 3a03906c..7e408580 100644 --- a/apps/web/src/lib/providers/web-http.ts +++ b/apps/web/src/lib/providers/web-http.ts @@ -122,13 +122,15 @@ export class WebHttpProxyProvider implements IHttpService { options.timeoutMs ) : undefined; - const clearRequestTimeout = () => { + const finishRequest = () => { if (timeoutId !== undefined) clearTimeout(timeoutId); + mergedSignal?.dispose(); }; const abortSignals = [options?.signal, controller?.signal].filter( (signal): signal is AbortSignal => signal !== undefined ); - const requestSignal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined; + const mergedSignal = abortSignals.length > 0 ? mergeAbortSignals(abortSignals) : undefined; + const requestSignal = mergedSignal?.signal; try { if (options?.bypassCors && !deploymentHasServerPlugins()) { @@ -195,7 +197,7 @@ export class WebHttpProxyProvider implements IHttpService { response.headers.forEach((val, key) => { responseHeaders[key] = val; }); - if (!response.body) clearRequestTimeout(); + if (!response.body) finishRequest(); // fetch resolves at headers; the same deadline must also cover body consumption. const consume = async (read: () => Promise): Promise => { try { @@ -205,7 +207,7 @@ export class WebHttpProxyProvider implements IHttpService { // Browsers may reject a timed-out body with AbortError; preserve the cause. throw requestSignal?.aborted ? requestSignal.reason : error; } finally { - clearRequestTimeout(); + finishRequest(); } }; @@ -219,7 +221,7 @@ export class WebHttpProxyProvider implements IHttpService { bytes: (onProgress) => consume(() => readResponseBytes(response, onProgress)) }; } catch (error) { - clearRequestTimeout(); + finishRequest(); throw requestSignal?.aborted ? requestSignal.reason : error; } } diff --git a/apps/web/src/lib/services/official-plugins/official-plugin-service.ts b/apps/web/src/lib/services/official-plugins/official-plugin-service.ts index dc4153eb..ba416782 100644 --- a/apps/web/src/lib/services/official-plugins/official-plugin-service.ts +++ b/apps/web/src/lib/services/official-plugins/official-plugin-service.ts @@ -673,10 +673,8 @@ export class OfficialPluginService implements Disposable { }) => void; } ): Promise { - const signal = options?.signal - ? mergeAbortSignals([options.signal, this.lifecycle.signal]) - : this.lifecycle.signal; - signal?.throwIfAborted?.(); + options?.signal?.throwIfAborted(); + this.lifecycle.signal.throwIfAborted(); validatePluginManifest(manifest); const sourceUrl = manifestUrl ?? @@ -696,8 +694,13 @@ export class OfficialPluginService implements Disposable { await this.installedStore.load(); const existingSnapshot = this.installedStore.find(manifest.id); const expectedRevision = existingSnapshot?.revision ?? -1; + const mergedSignal = mergeAbortSignals( + options?.signal ? [options.signal, this.lifecycle.signal] : [this.lifecycle.signal] + ); + const { signal } = mergedSignal; try { + signal.throwIfAborted(); const assets = await this.assetPipeline.download(manifest, manifestUrl, { signal, onProgress: options?.onProgress @@ -748,6 +751,8 @@ export class OfficialPluginService implements Disposable { throw new DOMException('Aborted', 'AbortError'); } throw err; + } finally { + mergedSignal.dispose(); } } diff --git a/apps/web/src/lib/utils/abort-signal.test.ts b/apps/web/src/lib/utils/abort-signal.test.ts new file mode 100644 index 00000000..448a350e --- /dev/null +++ b/apps/web/src/lib/utils/abort-signal.test.ts @@ -0,0 +1,47 @@ +import { getEventListeners } from 'node:events'; +import { describe, expect, it } from 'vite-plus/test'; +import { mergeAbortSignals } from './abort-signal'; + +describe('merged abort signal lifetime', () => { + it('removes listeners after successful operations sharing a long-lived signal', () => { + const lifecycle = new AbortController(); + for (let i = 0; i < 5; i++) { + const operation = new AbortController(); + const merged = mergeAbortSignals([lifecycle.signal, operation.signal]); + expect(merged.signal.aborted).toBe(false); + merged.dispose(); + merged.dispose(); + expect(getEventListeners(operation.signal, 'abort')).toHaveLength(0); + } + expect(getEventListeners(lifecycle.signal, 'abort')).toHaveLength(0); + }); + + it('preserves the abort reason and removes listeners from all sources', () => { + const first = new AbortController(); + const second = new AbortController(); + const merged = mergeAbortSignals([first.signal, second.signal, first.signal]); + const reason = new DOMException('Stopped', 'AbortError'); + second.abort(reason); + expect(merged.signal.reason).toBe(reason); + expect(getEventListeners(first.signal, 'abort')).toHaveLength(0); + expect(getEventListeners(second.signal, 'abort')).toHaveLength(0); + }); + + it('does not attach listeners when a later input is already aborted', () => { + const first = new AbortController(); + const second = new AbortController(); + second.abort('already stopped'); + const merged = mergeAbortSignals([first.signal, second.signal]); + expect(merged.signal.reason).toBe('already stopped'); + expect(getEventListeners(first.signal, 'abort')).toHaveLength(0); + merged.dispose(); + }); + + it('does not detach consumers of a single source signal', () => { + const source = new AbortController(); + const merged = mergeAbortSignals([source.signal]); + merged.dispose(); + source.abort('direct signal'); + expect(merged.signal.reason).toBe('direct signal'); + }); +}); diff --git a/apps/web/src/lib/utils/abort-signal.ts b/apps/web/src/lib/utils/abort-signal.ts index 8c2ccf5e..e76f60c9 100644 --- a/apps/web/src/lib/utils/abort-signal.ts +++ b/apps/web/src/lib/utils/abort-signal.ts @@ -1,23 +1,34 @@ /** * Merges multiple AbortSignals into one that aborts when any source aborts. * Does not rely on AbortSignal.any for broader runtime compatibility. + * Dispose after the operation settles to release listeners on non-aborted sources. */ -export function mergeAbortSignals(signals: AbortSignal[]): AbortSignal { - const active = signals.filter(Boolean); +export function mergeAbortSignals(signals: AbortSignal[]): { + signal: AbortSignal; + dispose(): void; +} { + const active = [...new Set(signals.filter(Boolean))]; if (active.length === 0) { throw new Error('mergeAbortSignals requires at least one signal'); } - if (active.length === 1) { - return active[0]!; + const aborted = active.find((signal) => signal.aborted); + if (aborted || active.length === 1) { + return { signal: aborted ?? active[0]!, dispose() {} }; } const controller = new AbortController(); + const listeners = new Map void>(); + const dispose = () => { + for (const [signal, listener] of listeners) signal.removeEventListener('abort', listener); + listeners.clear(); + }; for (const signal of active) { - if (signal.aborted) { + const listener = () => { controller.abort(signal.reason); - return controller.signal; - } - signal.addEventListener('abort', () => controller.abort(signal.reason), { once: true }); + dispose(); + }; + listeners.set(signal, listener); + signal.addEventListener('abort', listener, { once: true }); } - return controller.signal; + return { signal: controller.signal, dispose }; }