"use client"; import { useCallback, useEffect, useMemo, useRef, useState } from "react"; import { useVirtualizer } from "@tanstack/react-virtual"; import { useCurrentUser } from "@/hooks/use-current-user"; import { cancelBackfillJob, pauseBackfillJob, resumeBackfillJob, createBackfillJob, getBackfillJobs, getBackfillRepos, getBackfillPdsSummary, flushBackfillDetails, flushAllBackfillDetails, getLexicons, } from "@/lib/api"; import type { BackfillJob, BackfillRepoEntry, PdsSummaryEntry, BackfillEvent, BlueskyProfile, } from "@/types/backfill"; import { CheckCircle2, ChevronRight, Circle, Loader2, PauseCircle } from "lucide-react"; import { SiteHeader } from "@/components/site-header"; import { AlertDialog, AlertDialogAction, AlertDialogCancel, AlertDialogContent, AlertDialogDescription, AlertDialogFooter, AlertDialogHeader, AlertDialogTitle, AlertDialogTrigger, } from "@/components/ui/alert-dialog"; import { Badge } from "@/components/ui/badge"; import { Collapsible, CollapsibleContent, CollapsibleTrigger, } from "@/components/ui/collapsible"; import { Button } from "@/components/ui/button"; import { Combobox, ComboboxContent, ComboboxEmpty, ComboboxInput, ComboboxItem, ComboboxList, } from "@/components/ui/combobox"; import { ResponsiveDialog, ResponsiveDialogClose, ResponsiveDialogContent, ResponsiveDialogDescription, ResponsiveDialogFooter, ResponsiveDialogHeader, ResponsiveDialogTitle, ResponsiveDialogTrigger, } from "@/components/ui/responsive-dialog"; import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Sheet, SheetContent, SheetFooter, SheetHeader, SheetTitle, } from "@/components/ui/sheet"; import { Table, TableBody, TableCell, TableHead, TableHeader, TableRow, } from "@/components/ui/table"; const PROGRESS_PHASES = [ "discovering_repos", "resolving_pds", "fetching_records", ] as const; function statusBadge(job: BackfillJob) { switch (job.status) { case "completed": return ( completed ); case "failed": return failed; case "cancelled": return ( cancelled ); case "cancelling": return ( cancelling ); case "pausing": return ( pausing ); case "paused": return ( paused ); case "running": return ( {job.stage === "pending" ? "starting" : job.stage.replace(/_/g, " ")} ); default: return {job.status}; } } function phaseIndex(stage: string): number { if (stage === "resolving_and_fetching") return 1; const idx = PROGRESS_PHASES.indexOf( stage as (typeof PROGRESS_PHASES)[number], ); return idx; } // SSE via Web Worker — events are batched off the main thread and flushed periodically function useBackfillSSE( jobId: string | null, active: boolean, onBatch: (events: BackfillEvent[]) => void, ) { const onBatchRef = useRef(onBatch); onBatchRef.current = onBatch; useEffect(() => { if (!jobId || !active) return; const worker = new Worker( new URL("@/workers/backfill-sse.worker.ts", import.meta.url), ); worker.addEventListener("message", (e: MessageEvent) => { if (e.data.type === "batch") { onBatchRef.current(e.data.events as BackfillEvent[]); } }); const basePath = process.env.NEXT_PUBLIC_BASE_PATH || ""; worker.postMessage({ type: "connect", jobId, baseUrl: `${location.origin}${basePath}` }); return () => { worker.postMessage({ type: "disconnect" }); worker.terminate(); }; }, [jobId, active]); } // Batch Bluesky profile resolution hook function useBlueskyProfiles(dids: string[]): Map { const [profiles, setProfiles] = useState>(new Map()); const resolvedRef = useRef>(new Set()); const pendingRef = useRef(false); useEffect(() => { const unresolved = dids.filter((d) => !resolvedRef.current.has(d)); if (unresolved.length === 0 || pendingRef.current) return; pendingRef.current = true; const batches: string[][] = []; for (let i = 0; i < unresolved.length; i += 25) { batches.push(unresolved.slice(i, i + 25)); } // Mark all as resolved immediately to prevent re-fetching for (const did of unresolved) { resolvedRef.current.add(did); } Promise.all( batches.map(async (batch) => { const params = batch.map((d) => `actors=${encodeURIComponent(d)}`).join("&"); try { const resp = await fetch( `https://public.api.bsky.app/xrpc/app.bsky.actor.getProfiles?${params}` ); if (!resp.ok) return []; const data = await resp.json(); return (data.profiles || []) as BlueskyProfile[]; } catch { return []; } }) ).then((results) => { setProfiles((prev) => { const newProfiles = new Map(prev); for (const batch of results) { for (const p of batch) { newProfiles.set(p.did, p); } } return newProfiles; }); pendingRef.current = false; }); }, [dids]); return profiles; } const BSKY_PDS_SUFFIX = ".bsky.network"; const BSKY_PDS_HOSTNAMES = ["bsky.social", "staging.bsky.dev"]; const failedFaviconUrls = new Set(); function isBskyPds(pdsEndpoint: string): boolean { try { const hostname = new URL(pdsEndpoint).hostname; return BSKY_PDS_HOSTNAMES.includes(hostname) || hostname.endsWith(BSKY_PDS_SUFFIX); } catch { return false; } } function PdsFavicon({ pdsEndpoint }: { pdsEndpoint: string }) { let hostname: string; try { hostname = new URL(pdsEndpoint).hostname; } catch { return ; } if (isBskyPds(pdsEndpoint)) { return ( ); } const faviconUrl = `https://twenty-icons.com/${hostname}`; if (failedFaviconUrls.has(faviconUrl)) { return ; } return ; } function PdsFaviconImg({ url }: { url: string }) { const [failed, setFailed] = useState(false); if (failed) return ; return ( { failedFaviconUrls.add(url); setFailed(true); }} /> ); } function PdsPlaceholderIcon() { return ( ); } function AnimatedNumber({ value }: { value: number }) { const targetRef = useRef(value); const displayedRef = useRef(value); const [displayed, setDisplayed] = useState(value); const rafRef = useRef(0); targetRef.current = value; useEffect(() => { cancelAnimationFrame(rafRef.current); function tick() { const current = displayedRef.current; const target = targetRef.current; const diff = target - current; if (Math.abs(diff) < 0.5) { displayedRef.current = target; setDisplayed(target); return; } const next = current + diff * 0.06; displayedRef.current = next; setDisplayed(next); rafRef.current = requestAnimationFrame(tick); } tick(); return () => cancelAnimationFrame(rafRef.current); }, [value]); return <>{Math.round(displayed).toLocaleString()}>; } function CompactRepoRow({ did, profile }: { did: string; profile?: BlueskyProfile; }) { return ( {profile?.avatar && ( )} {profile?.handle ? `@${profile.handle}` : {did}} ); } function ProfileRow({ did, profile, suffix }: { did: string; profile?: BlueskyProfile; suffix?: React.ReactNode; }) { return ( {profile?.avatar && ( )} {profile?.displayName || profile?.handle || did} {profile?.handle && ( @{profile.handle} )} {!profile?.handle && ( {did} )} {suffix && ( {suffix} )} ); } export default function BackfillPage() { const { hasPermission } = useCurrentUser(); const [jobs, setJobs] = useState([]); const [error, setError] = useState(null); const [selectedJobId, setSelectedJobId] = useState(null); const load = useCallback(() => { getBackfillJobs() .then(setJobs) .catch((e) => setError(e.message)); }, []); useEffect(() => { load(); }, [load]); const selectedJob = jobs.find((j) => j.id === selectedJobId) ?? null; const sseActive = selectedJob != null && (selectedJob.status === "running" || selectedJob.status === "cancelling" || selectedJob.status === "pausing"); useEffect(() => { if (sseActive) return; const interval = setInterval(load, 5000); return () => clearInterval(interval); }, [load, sseActive]); const canFlush = hasPermission("backfill:create"); return ( <> {error && {error}} Backfill Jobs {canFlush && ( Clear all details Clear all job details? This will permanently delete per-repo detail data for all backfill jobs. Cancel { await flushAllBackfillDetails(); setSelectedJobId(null); load(); }}>Clear )} {hasPermission("backfill:create") && ( )} ID Collection DID Status Started {jobs.length === 0 && ( No backfill jobs yet. )} {jobs.map((job) => ( setSelectedJobId(job.id)} onKeyDown={(e) => { if (e.key === "Enter" || e.key === " ") { e.preventDefault(); setSelectedJobId(job.id); } }} > {job.id.slice(0, 8)} {job.collection ?? "All"} {job.did ?? "All"} {statusBadge(job)} {job.started_at ? new Date(job.started_at).toLocaleString() : "--"} ))} { if (!open) { setSelectedJobId(null); load(); } }} > {selectedJob && ( { setJobs((prev) => prev.map((j) => j.id === selectedJob.id ? updater(j) : j, ), ); }} onCancel={async () => { await cancelBackfillJob(selectedJob.id); load(); }} onPause={async () => { await pauseBackfillJob(selectedJob.id); load(); }} onResume={async () => { await resumeBackfillJob(selectedJob.id); load(); }} /> )} > ); } function JobDetail({ job, canCancel, canFlush, onJobUpdate, onCancel, onPause, onResume, }: { job: BackfillJob; canCancel: boolean; canFlush: boolean; onJobUpdate: (updater: (job: BackfillJob) => BackfillJob) => void; onCancel: () => Promise; onPause: () => Promise; onResume: () => Promise; }) { const [cancelling, setCancelling] = useState(false); const [pausing, setPausing] = useState(false); const [resuming, setResuming] = useState(false); const current = phaseIndex(job.stage); const allDone = job.status === "completed"; const isActive = job.status === "running" || job.status === "cancelling" || job.status === "pausing"; const isPaused = job.status === "paused" || job.status === "pausing"; // Detail data state const [discoveredRepos, setDiscoveredRepos] = useState([]); const [discoveredCursor, setDiscoveredCursor] = useState(null); const [discoveredLoaded, setDiscoveredLoaded] = useState(false); const [pdsSummary, setPdsSummary] = useState([]); const [pdsLoaded, setPdsLoaded] = useState(false); const [fetchedRepos, setFetchedRepos] = useState([]); const [fetchedCursor, setFetchedCursor] = useState(null); const [fetchedLoaded, setFetchedLoaded] = useState(false); // Refs for open state and callbacks so the SSE callback doesn't need to re-bind on toggle const onJobUpdateRef = useRef(onJobUpdate); onJobUpdateRef.current = onJobUpdate; function hasReached(phase: (typeof PROGRESS_PHASES)[number]): boolean { if (allDone) return true; if (job.stage === "resolving_and_fetching") { return phase === "discovering_repos" || phase === "resolving_pds" || phase === "fetching_records"; } return current >= phaseIndex(phase); } function isPhasePaused(phase: (typeof PROGRESS_PHASES)[number]): boolean { if (!isPaused) return false; if (job.stage === "resolving_and_fetching") { return phase === "resolving_pds" || phase === "fetching_records"; } if (job.stage === "discovering_repos") { return phase === "discovering_repos"; } if (job.stage === "resolving_pds") { return phase === "resolving_pds"; } if (job.stage === "fetching_records") { return phase === "fetching_records"; } return false; } async function handleCancel() { setCancelling(true); try { await onCancel(); } finally { setCancelling(false); } } async function handlePause() { setPausing(true); try { await onPause(); } finally { setPausing(false); } } async function handleResume() { setResuming(true); try { await onResume(); } finally { setResuming(false); } } const discoveredReached = hasReached("discovering_repos"); const pdsReached = hasReached("resolving_pds"); const fetchedReached = hasReached("fetching_records") || job.stage === "resolving_and_fetching"; // Track which sections are expanded const [discoveredOpen, setDiscoveredOpen] = useState(false); const [pdsOpen, setPdsOpen] = useState(false); const [fetchedOpen, setFetchedOpen] = useState(false); // Lazy-load detail data only when sections are expanded useEffect(() => { if (discoveredOpen && discoveredReached && !discoveredLoaded) { getBackfillRepos(job.id, { phase: "discovered", limit: 50 }) .then((resp) => { setDiscoveredRepos(resp.repos); setDiscoveredCursor(resp.cursor); setDiscoveredLoaded(true); }) .catch(() => {}); } }, [job.id, discoveredOpen, discoveredReached, discoveredLoaded]); useEffect(() => { if (pdsOpen && pdsReached && !pdsLoaded) { getBackfillPdsSummary(job.id) .then((resp) => { setPdsSummary(resp.pds_endpoints); setPdsLoaded(true); }) .catch(() => {}); } }, [job.id, pdsOpen, pdsReached, pdsLoaded]); useEffect(() => { if (fetchedOpen && fetchedReached && !fetchedLoaded) { getBackfillRepos(job.id, { phase: "fetched", limit: 50 }) .then((resp) => { setFetchedRepos(resp.repos); setFetchedCursor(resp.cursor); setFetchedLoaded(true); }) .catch(() => {}); } }, [job.id, fetchedOpen, fetchedReached, fetchedLoaded]); const loadMoreDiscovered = useCallback(async () => { if (!discoveredCursor) return; try { const resp = await getBackfillRepos(job.id, { phase: "discovered", cursor: discoveredCursor, limit: 50 }); setDiscoveredRepos((prev) => [...prev, ...resp.repos]); setDiscoveredCursor(resp.cursor); } catch { /* ignore */ } }, [job.id, discoveredCursor]); const loadMoreFetched = useCallback(async () => { if (!fetchedCursor) return; try { const resp = await getBackfillRepos(job.id, { phase: "fetched", cursor: fetchedCursor, limit: 50 }); setFetchedRepos((prev) => [...prev, ...resp.repos]); setFetchedCursor(resp.cursor); } catch { /* ignore */ } }, [job.id, fetchedCursor]); // Process batched SSE events from the Web Worker. // Uses refs for open state so the callback identity is stable and doesn't // cause the worker to reconnect when sections are toggled. const handleSSEBatch = useCallback((events: BackfillEvent[]) => { const update = onJobUpdateRef.current; for (const e of events) { if (e.type === "job_counters") { update((j) => ({ ...j, ...(e.total_repos != null && { total_repos: e.total_repos }), ...(e.resolved_repos != null && { resolved_repos: e.resolved_repos }), ...(e.processed_repos != null && { processed_repos: e.processed_repos }), ...(e.total_records != null && { total_records: e.total_records }), })); } else if (e.type === "job_stage_changed" && e.stage) { update((j) => ({ ...j, stage: e.stage! })); } else if (e.type === "job_completed" && e.status) { update((j) => ({ ...j, status: e.status!, error: e.error ?? null })); } } }, []); useBackfillSSE(job.id, isActive, handleSSEBatch); // Track visible DIDs from virtualized lists (only items in viewport) const [visibleDiscoveredDids, setVisibleDiscoveredDids] = useState([]); const [visibleFetchedDids, setVisibleFetchedDids] = useState([]); const allVisibleDids = useMemo(() => { const dids = new Set(); for (const d of visibleDiscoveredDids) dids.add(d); for (const d of visibleFetchedDids) dids.add(d); return Array.from(dids); }, [visibleDiscoveredDids, visibleFetchedDids]); const profiles = useBlueskyProfiles(allVisibleDids); const sortedPdsSummary = useMemo( () => [...pdsSummary].sort((a, b) => b.total_repos - a.total_repos), [pdsSummary], ); const fetchedWithRecords = useMemo( () => fetchedRepos.filter((r) => r.records_fetched > 0), [fetchedRepos], ); return ( <> Backfill Details Job ID {job.id} Collection {job.collection ?? "All"} DID {job.did ?? "All"} Created {new Date(job.created_at).toLocaleString()} Started {job.started_at ? new Date(job.started_at).toLocaleString() : "--"} {job.completed_at && ( Completed {new Date(job.completed_at).toLocaleString()} )} {job.error && ( Error {job.error} )} Progress : undefined} suffix="repos found" loading={discoveredOpen && discoveredReached && !discoveredLoaded} open={discoveredOpen} onOpenChange={setDiscoveredOpen} > {discoveredRepos.length > 0 ? ( r.did} onVisibleKeysChange={setVisibleDiscoveredDids} hasMore={!!discoveredCursor} onLoadMore={loadMoreDiscovered} rowHeight={28} renderRow={(repo) => ( )} /> ) : discoveredLoaded ? ( No repos discovered yet. ) : null} / > : undefined } suffix="resolved" loading={pdsOpen && pdsReached && !pdsLoaded} open={pdsOpen} onOpenChange={setPdsOpen} > {sortedPdsSummary.length > 0 ? ( p.pds_endpoint} hasMore={false} onLoadMore={() => {}} rowHeight={32} renderRow={(pds) => ( {new URL(pds.pds_endpoint).hostname} / repos · records )} /> ) : pdsLoaded ? ( No PDS data yet. ) : null} / repos> : undefined } suffix={ hasReached("fetching_records") || job.stage === "resolving_and_fetching" ? <> records> : undefined } loading={fetchedOpen && fetchedReached && !fetchedLoaded} open={fetchedOpen} onOpenChange={setFetchedOpen} > {fetchedWithRecords.length > 0 ? ( r.did} onVisibleKeysChange={setVisibleFetchedDids} hasMore={!!fetchedCursor} onLoadMore={loadMoreFetched} rowHeight={40} renderRow={(repo) => ( records} /> )} /> ) : fetchedLoaded ? ( No repos fetched yet. ) : null} {canFlush && !isActive && ( Clear details Clear job details? This will permanently delete per-repo detail data for this backfill job. Cancel { await flushBackfillDetails(job.id); setDiscoveredRepos([]); setDiscoveredCursor(null); setDiscoveredLoaded(false); setPdsSummary([]); setPdsLoaded(false); setFetchedRepos([]); setFetchedCursor(null); setFetchedLoaded(false); }}>Clear )} {canCancel && (job.status === "running" || job.status === "pausing") && ( {pausing || job.status === "pausing" ? "Pausing…" : "Pause Job"} )} {canCancel && job.status === "paused" && ( {resuming ? "Resuming…" : "Resume Job"} )} {canCancel && isActive && ( {job.status === "cancelling" ? "Cancelling…" : "Cancel Job"} )} {canCancel && job.status === "paused" && ( Cancel Job )} > ); } function VirtualList({ items, getKey, hasMore, onLoadMore, rowHeight, renderRow, onVisibleKeysChange, }: { items: T[]; getKey: (item: T) => string; hasMore: boolean; onLoadMore: () => void; rowHeight: number; renderRow: (item: T) => React.ReactNode; onVisibleKeysChange?: (keys: string[]) => void; }) { const parentRef = useRef(null); const loadMoreTriggered = useRef(false); const virtualizer = useVirtualizer({ count: items.length, getScrollElement: () => parentRef.current, estimateSize: () => rowHeight, overscan: 5, }); const virtualItems = virtualizer.getVirtualItems(); const visibleKeysStr = virtualItems.map((item) => getKey(items[item.index])).filter(Boolean).join(","); useEffect(() => { onVisibleKeysChange?.(visibleKeysStr.split(",").filter(Boolean)); }, [visibleKeysStr, onVisibleKeysChange]); useEffect(() => { if (!hasMore) return; const lastItem = virtualItems[virtualItems.length - 1]; if (lastItem && lastItem.index >= items.length - 5 && !loadMoreTriggered.current) { loadMoreTriggered.current = true; onLoadMore(); } if (lastItem && lastItem.index < items.length - 5) { loadMoreTriggered.current = false; } }, [virtualItems, items.length, hasMore, onLoadMore]); return ( {virtualItems.map((virtualRow) => { const item = items[virtualRow.index]; if (!item) return null; return ( {renderRow(item)} ); })} ); } function ProgressRow({ label, active, reached, paused, value, suffix, loading, open, onOpenChange, children, }: { label: string; active: boolean; reached: boolean; paused?: boolean; value?: React.ReactNode; suffix?: React.ReactNode; loading?: boolean; open: boolean; onOpenChange: (open: boolean) => void; children?: React.ReactNode; }) { const done = reached && !active && !paused; const expandable = reached; return ( {active ? ( ) : paused ? ( ) : done ? ( ) : ( )} {label} {reached && value && ( {value} {suffix && <> · {suffix}>} )} {expandable && ( )} {expandable && ( {loading ? ( ) : children} )} ); } function CreateDialog({ onSuccess }: { onSuccess: () => void }) { const [collection, setCollection] = useState(null); const [did, setDid] = useState(""); const [error, setError] = useState(null); const [open, setOpen] = useState(false); const [recordLexicons, setRecordLexicons] = useState([]); useEffect(() => { if (open) { getLexicons() .then((lexicons) => setRecordLexicons( lexicons .filter((l) => l.lexicon_type === "record") .map((l) => l.id) .sort(), ), ) .catch(() => {}); } }, [open]); async function handleCreate() { setError(null); try { await createBackfillJob({ collection: collection || undefined, did: did || undefined, }); setCollection(null); setDid(""); setOpen(false); onSuccess(); } catch (e: unknown) { setError(e instanceof Error ? e.message : String(e)); } } return ( Create Backfill Job { const target = e.target as HTMLElement; if ( target.closest( "[data-slot='combobox-item'], [data-slot='combobox-content']", ) ) { e.preventDefault(); } }} > Create Backfill Job Start a backfill for a collection or specific DID. Leave both empty to backfill all collections. {error && {error}} Collection (optional) No matching lexicons. {(item: string) => ( {item} )} DID (optional) setDid(e.target.value)} placeholder="did:plc:..." /> Cancel Create ); }
{profile?.displayName || profile?.handle || did}
@{profile.handle}
{did}
{error}
{job.id}
{job.collection ?? "All"}
{job.did ?? "All"}
{new Date(job.created_at).toLocaleString()}
{job.started_at ? new Date(job.started_at).toLocaleString() : "--"}
{new Date(job.completed_at).toLocaleString()}
No repos discovered yet.
No PDS data yet.
No repos fetched yet.