import type { FileContents, FileDiffMetadata } from "@pierre/diffs"; import DiffWorker from "$lib/components/repo/diffWorker?worker"; import type { DiffWorkerRequest, DiffWorkerResponse } from "$lib/components/repo/diffWorker"; // Strip Svelte 5 reactive proxies; postMessage throws DataCloneError. let diffWorkerInstance: Worker | undefined; let nextTaskId = 1; interface PendingTask { resolve: (response: DiffWorkerResponse) => void; reject: (err: unknown) => void; /** set for streamed patches, which answer many times before they finish */ partial?: (files: FileDiffMetadata[]) => void; } const pendingTasks = new Map(); const DIFF_WORKER_TIMEOUT_MS = 30_000; const failPending = (error: Error) => { for (const [id, task] of pendingTasks) { pendingTasks.delete(id); task.reject(error); } }; const serializeFileContents = (file: FileContents): FileContents => ({ name: file.name, contents: typeof file.contents === "string" ? file.contents : "", ...(file.cacheKey !== undefined ? { cacheKey: file.cacheKey } : {}) }); const diffWorker = (): Worker => { if (diffWorkerInstance) return diffWorkerInstance; const worker = new DiffWorker(); worker.onmessage = (e: MessageEvent) => { const task = pendingTasks.get(e.data.id); if (!task) return; if (task.partial && !e.data.done) { if (e.data.files) task.partial(e.data.files); return; } pendingTasks.delete(e.data.id); task.resolve(e.data); }; worker.onerror = () => { failPending(new Error("diff worker failed")); diffWorkerInstance?.terminate(); diffWorkerInstance = undefined; }; diffWorkerInstance = worker; return worker; }; const runInDiffWorker = (build: (id: number) => DiffWorkerRequest): Promise => { const worker = diffWorker(); const id = nextTaskId++; return new Promise((resolve, reject) => { const timer = setTimeout(() => { if (!pendingTasks.delete(id)) return; reject(new Error("diff worker timed out")); }, DIFF_WORKER_TIMEOUT_MS); pendingTasks.set(id, { resolve: (response) => { clearTimeout(timer); resolve(response); }, reject: (err) => { clearTimeout(timer); reject(err); } }); worker.postMessage(build(id)); }); }; export const parseDiffInWorker = async ( oldFile: FileContents, newFile: FileContents ): Promise => { const response = await runInDiffWorker((id) => ({ id, kind: "file", oldFile: serializeFileContents(oldFile), newFile: serializeFileContents(newFile) })); if (response.fileDiff) return response.fileDiff; throw new Error(response.error ?? "diff worker returned no diff"); }; export const streamPatchInWorker = async ( url: string, onFiles: (files: FileDiffMetadata[]) => void, signal?: AbortSignal ): Promise => { // Fetch on main thread to reuse cache (worker has separate cache). let body: ReadableStream | undefined; try { const response = await fetch(url, { signal }); if (!response.ok) throw new Error(`patch ${response.status}`); body = response.body ?? undefined; } catch (cause) { // a failed handover is worth one retry from inside the worker, which has // its own fetch and no preload to match if (body === undefined && cause instanceof Error && cause.message.startsWith("patch ")) throw cause; } const worker = diffWorker(); const id = nextTaskId++; await new Promise((resolve, reject) => { pendingTasks.set(id, { partial: onFiles, resolve: (response) => { if (response.error) return reject(new Error(response.error)); if (response.files?.length) onFiles(response.files); resolve(); }, reject }); const request: DiffWorkerRequest = body ? { id, kind: "stream", body } : { id, kind: "stream", url }; worker.postMessage(request, body ? [body] : []); // leaving the page stops the parse too: without this the worker keeps // chewing through megabytes for a component that no longer exists signal?.addEventListener( "abort", () => { if (!pendingTasks.delete(id)) return; worker.postMessage({ id, kind: "cancel" } satisfies DiffWorkerRequest); resolve(); }, { once: true } ); }); };