import type { FileDiffMetadata } from "@pierre/diffs"; import { fetchSides, type RepoDiff, type RepoDiffDeps, type RepoDiffFile } from "$lib/api/repoDiff"; import { parseDiffInWorker, streamPatchInWorker } from "$lib/components/repo/diffRpc"; import type { DiffEntry } from "$lib/components/repo/fileDiff"; export type { DiffEntry }; export interface DiffQueueHost { deps: () => RepoDiffDeps; diff: () => RepoDiff; entry: (key: string) => DiffEntry | undefined; put: (parsed: Record) => void; /** true while a patch is still expected to deliver the whole diff */ patchExpected?: () => boolean; } export interface DiffQueue { focus: (key: string) => void; want: (key: string) => void; /** re-pumps files parked while a patch was expected; call when it settles */ drain: () => void; applyPatch: (url: string, signal: AbortSignal) => Promise; /** hands the queue files someone else already parsed, so it stops chasing them */ adopt: (parsed: Record) => void; destroy: () => void; } const CONCURRENCY = 6; export const createDiffQueue = (host: DiffQueueHost): DiffQueue => { let running = 0; const pending = new Set(); const inFlight = new Set(); // file order is monotonic down the page, so index distance stands in for // pixels let focusIndex = 0; let indexed: { files: RepoDiffFile[]; byKey: Map; }; const byKey = () => { const files = host.diff().files; if (indexed?.files !== files) indexed = { files, byKey: new Map(files.map((file, index) => [file.key, { file, index }])) }; return indexed.byKey; }; const nextKey = (): string | undefined => { const index = byKey(); let best: string | undefined; let bestDistance = Infinity; for (const key of pending) { const distance = Math.abs((index.get(key)?.index ?? 0) - focusIndex); if (distance < bestDistance) { best = key; bestDistance = distance; } } return best; }; const pump = () => { // one patch request beats two blob fetches per file, so while one is on its // way the wanted files wait for it. what the patch does not carry stays // queued and is fetched the moment it settles, either way if (host.patchExpected?.()) return; while (running < CONCURRENCY) { const key = nextKey(); if (!key) return; pending.delete(key); const file = byKey().get(key)?.file; if (file) void fill(file); } }; // the mirror rate-limits bursts, and a fling scroll is exactly a burst. a // retried file lands a moment later instead of showing a permanent error const RETRIES = 3; const backoff = (attempt: number) => new Promise((resolve) => setTimeout(resolve, 300 * 3 ** attempt + Math.random() * 200)); const fetchWithRetry = async (file: RepoDiffFile) => { for (let attempt = 0; ; attempt++) { try { return await fetchSides(host.deps(), host.diff().contents, file); } catch (cause) { if (attempt >= RETRIES) throw cause; await backoff(attempt); } } }; const fill = async (file: RepoDiffFile) => { running++; inFlight.add(file.key); try { if (host.entry(file.key)) return; const contents = await fetchWithRetry(file); if (!("oldFile" in contents)) { host.put({ [file.key]: { note: contents.note } }); } else { host.put({ [file.key]: await parseDiffInWorker( { ...contents.oldFile, cacheKey: `${file.key}:old` }, { ...contents.newFile, cacheKey: `${file.key}:new` } ) }); } } catch { // the patch may have delivered this file while the fetch was failing, and // a real diff beats an error note if (!host.entry(file.key)) { host.put({ [file.key]: { note: "Contents could not be loaded." } }); } } finally { inFlight.delete(file.key); running--; pump(); } }; return { destroy: () => { pending.clear(); }, focus: (key) => { focusIndex = byKey().get(key)?.index ?? focusIndex; }, want: (key) => { if (host.entry(key) || pending.has(key) || inFlight.has(key)) return; if (!byKey().get(key)?.file.sides) return; pending.add(key); pump(); }, drain: () => pump(), adopt: (parsed) => { const landing: Record = {}; for (const [key, value] of Object.entries(parsed)) { // a file already being fetched keeps its own result: that one carries // the whole file, so its context can still be expanded if (host.entry(key) || inFlight.has(key)) continue; landing[key] = value; pending.delete(key); } if (Object.keys(landing).length) host.put(landing); }, applyPatch: async (url, signal) => { // the worker streams the body and answers in batches, so the first files // land about a second in instead of after the whole patch has arrived const land = (files: FileDiffMetadata[]) => { if (signal.aborted) return; const byName: Record = {}; for (const parsed of files) { byName[parsed.name] = parsed; if (parsed.prevName) byName[parsed.prevName] = parsed; } const landed: Record = {}; for (const file of index.values()) { // keep sides already landed or on the wire: that diff carries the whole // file, so its context can still be expanded if (host.entry(file.file.key) || inFlight.has(file.file.key)) continue; const parsed = byName[file.file.name] ?? (file.file.oldName ? byName[file.file.oldName] : undefined); if (!parsed) continue; landed[file.file.key] = parsed; pending.delete(file.file.key); } if (Object.keys(landed).length) host.put(landed); }; const index = byKey(); await streamPatchInWorker(url, land, signal); } }; };