Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181import 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<string, DiffEntry>) => 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<void>; /** hands the queue files someone else already parsed, so it stops chasing them */ adopt: (parsed: Record<string, DiffEntry>) => void; destroy: () => void;}
const CONCURRENCY = 6;
export const createDiffQueue = (host: DiffQueueHost): DiffQueue => { let running = 0; const pending = new Set<string>(); const inFlight = new Set<string>(); // file order is monotonic down the page, so index distance stands in for // pixels let focusIndex = 0;
let indexed: { files: RepoDiffFile[]; byKey: Map<string, { file: RepoDiffFile; index: number }>; }; 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<string, DiffEntry> = {}; 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<string, FileDiffMetadata> = {}; for (const parsed of files) { byName[parsed.name] = parsed; if (parsed.prevName) byName[parsed.prevName] = parsed; } const landed: Record<string, DiffEntry> = {}; 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); } };};