// ONE SEARCH PIPELINE — the streaming leaf scanner. // // Moved out of `components/searchPipeline.ts`, which is now the binding that // supplies this module's fetchers (`transcriptCache` / `subsCache` / // `postsCache`) and nothing else. Nothing here touches react-query, a cache, a // DOM node or `window` — the fetch is a parameter — which is what lets // `lib/searchEval.ts` drive a leaf without `lib/` importing `components/`. // // Three shapes of leaf, one worker-pool driver each, because the unit of // iteration genuinely differs: a transcript record (cues / description / tags), // a subs record (per-track cue lists), a post (one body). `runLeafPipeline` // picks between them from the leaf's scope and adapts all three onto one // `LayerHit` stream. // // State is closure-captured so `setHitLimit` can raise the cap and re-spawn // workers without restarting traversal — workers resume from the shared `idx` // cursor. // // TIMERS: plain `setTimeout`/`clearTimeout`, not `window.setTimeout`. Same // function in a browser; the bare form is also callable from a test process, // which is what makes this module directly testable without a DOM. import type { TranscriptDetail, DisplaySummary } from "../transcripts"; import type { SubsDetail } from "../subs"; import type { Post } from "../posts"; import type { LayerScope, LeafNode } from "../searchQuery"; import { findHitsInCues, findHitsInText, type Hit, type SubsHit } from "./window"; import { hitsAcrossTracks } from "../captionTracks"; export type { Hit, SubsHit }; // Per-leaf hit shape consumed by the composite-search tree (`lib/searchEval.ts`) // and rendered by the result cards. Carries enough info for the UI to draw the // hit with its originating layer's swatch + scope-specific decorations. export type LayerHit = { leafId: string; scope: LayerScope; track?: string; start: number; text: string; }; // The three fetches a leaf can need, injected. `slug` is the member-local video // or post slug the scope list is built from. export type LeafFetchers = { transcript: (slug: string) => Promise; subs: (slug: string) => Promise; post: (slug: string) => Promise; }; export type PipelineUpdate = { hitsBySlug: Record; totalHits: number; processed: number; totalToProcess: number; capped: boolean; done: boolean; }; type PipelineConfig = { slugs: string[]; query: string; useRegex: boolean; regex: RegExp | null; initialHitLimit: number; concurrency: number; flushIntervalMs: number; emit: (update: PipelineUpdate) => void; fetchTranscript: (slug: string) => Promise; // Which field of the fetched transcript document to match. "cues" (default) // is the transcripts scope; "description"/"tags" match those metadata fields. // All reuse the same per-channel transcript-page fetch + worker pool. matchField?: "cues" | "description" | "tags"; }; export type PipelineController = { cancel(): void; setHitLimit(limit: number): void; }; export function createSearchPipeline( config: PipelineConfig, ): PipelineController { const { slugs, query, useRegex, regex, initialHitLimit, concurrency, flushIntervalMs, emit, fetchTranscript, matchField = "cues", } = config; let cancelled = false; // `idx` is the next-to-claim cursor (leading edge); `completed` counts // slugs that workers have finished scanning (trailing edge). We report // `completed` to the UI so progress reflects actual work done, not // in-flight claims — and it never exceeds slugs.length. let idx = 0; let completed = 0; let totalSoFar = 0; let hitLimit = initialHitLimit; let activeWorkers = 0; let done = false; const localHits: Record = {}; let flushTimer: ReturnType | null = null; const pushUpdate = (overrides: Partial = {}) => { emit({ hitsBySlug: { ...localHits }, totalHits: totalSoFar, processed: completed, totalToProcess: slugs.length, // `capped` keys off `idx` (claimed), not `completed` — raising the cap // only helps if there are still unclaimed slugs a worker can pick up. capped: totalSoFar >= hitLimit && idx < slugs.length, done, ...overrides, }); }; const scheduleFlush = () => { if (flushTimer !== null || cancelled) return; flushTimer = setTimeout(() => { flushTimer = null; if (cancelled) return; pushUpdate(); }, flushIntervalMs); }; const finalize = () => { if (done || cancelled) return; done = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } pushUpdate(); }; const worker = async () => { activeWorkers++; try { while (!cancelled) { if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const my = idx++; const slug = slugs[my]; try { const full = await fetchTranscript(slug); if (cancelled) return; if (totalSoFar < hitLimit) { const remaining = hitLimit - totalSoFar; const hits = matchField === "description" ? findHitsInText(full.description ?? "", query, useRegex, regex) : matchField === "tags" ? findHitsInText( (full.tags ?? []).join(", "), query, useRegex, regex, ) : // Every English track of the record (lib/captionTracks.ts). hitsAcrossTracks(full, (cues) => findHitsInCues(cues, query, useRegex, regex, remaining), ); if (hits.length > 0) { localHits[slug] = hits; totalSoFar += hits.length; } } } catch { // ignore per-transcript failures } completed++; scheduleFlush(); } } finally { activeWorkers--; if (activeWorkers === 0 && !cancelled) { // Either we hit the cap or ran out of work — either way, settle. finalize(); } } }; const ensureWorkers = () => { if (cancelled || done) return; if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const needed = Math.min(concurrency - activeWorkers, slugs.length - idx); for (let i = 0; i < needed; i++) worker(); }; // `worker()` counts itself in synchronously, before its first await, so // right after ensureWorkers a zero here means none is running or will. const settleIfIdle = () => { if (activeWorkers === 0 && !done) finalize(); }; // Emit initial snapshot synchronously so the UI clears previous results. pushUpdate(); ensureWorkers(); // No worker started (nothing to scan, or the cap is already met): nobody is // left to call finalize, so settle here. Without it an empty scope — which // "Search in" makes from a plain query, e.g. a posts copy under a tag filter // — never reported done and the search read "searching" for ever. settleIfIdle(); return { cancel() { cancelled = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } }, setHitLimit(limit: number) { if (cancelled) return; if (limit <= hitLimit) return; hitLimit = limit; // The prior run may have finalized because we were capped. Resume. if (done) { done = false; pushUpdate(); } ensureWorkers(); settleIfIdle(); }, }; } // Streaming search over the social-post corpus. Structurally the transcripts // pipeline with a different fetch + match: the unit of iteration is a POST // slug (`/`), and the injected fetch resolves through the // page cache, so the first post on a page warms every other post on it. // // Post bodies are short but unbounded, and unlike the cue path there is no // natural per-line unit to clip to — so hits are truncated to a ±80-char // window (the same one the description/tags scopes use) rather than shipping // the whole body into a result card. export function createPostsSearchPipeline( config: Omit & { fetchPost: (slug: string) => Promise; }, ): PipelineController { const { slugs, query, useRegex, regex, initialHitLimit, concurrency, flushIntervalMs, emit, fetchPost, } = config; let cancelled = false; let idx = 0; let completed = 0; let totalSoFar = 0; let hitLimit = initialHitLimit; let activeWorkers = 0; let done = false; const localHits: Record = {}; let flushTimer: ReturnType | null = null; const pushUpdate = (overrides: Partial = {}) => { emit({ hitsBySlug: { ...localHits }, totalHits: totalSoFar, processed: completed, totalToProcess: slugs.length, capped: totalSoFar >= hitLimit && idx < slugs.length, done, ...overrides, }); }; const scheduleFlush = () => { if (flushTimer !== null || cancelled) return; flushTimer = setTimeout(() => { flushTimer = null; if (cancelled) return; pushUpdate(); }, flushIntervalMs); }; const finalize = () => { if (done || cancelled) return; done = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } pushUpdate(); }; const worker = async () => { activeWorkers++; try { while (!cancelled) { if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const my = idx++; const slug = slugs[my]; try { const post = await fetchPost(slug); if (cancelled) return; if (totalSoFar < hitLimit) { const hits = findHitsInText(post.text, query, useRegex, regex); if (hits.length > 0) { localHits[slug] = hits; totalSoFar += hits.length; } } } catch { // ignore per-post failures } completed++; scheduleFlush(); } } finally { activeWorkers--; if (activeWorkers === 0 && !cancelled) finalize(); } }; const ensureWorkers = () => { if (cancelled || done) return; if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const needed = Math.min(concurrency - activeWorkers, slugs.length - idx); for (let i = 0; i < needed; i++) worker(); }; // `worker()` counts itself in synchronously, before its first await, so // right after ensureWorkers a zero here means none is running or will. const settleIfIdle = () => { if (activeWorkers === 0 && !done) finalize(); }; pushUpdate(); ensureWorkers(); settleIfIdle(); // see createSearchPipeline return { cancel() { cancelled = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } }, setHitLimit(limit: number) { if (cancelled) return; if (limit <= hitLimit) return; hitLimit = limit; if (done) { done = false; pushUpdate(); } ensureWorkers(); settleIfIdle(); }, }; } // Helper to build slug list from summaries + filter predicate. export function filterSlugs( summaries: DisplaySummary[], passes: (t: DisplaySummary) => boolean, ): string[] { const slugs: string[] = []; for (const t of summaries) if (passes(t)) slugs.push(t.slug); return slugs; } export type SubsPipelineUpdate = { hitsBySlug: Record; totalHits: number; processed: number; totalToProcess: number; capped: boolean; done: boolean; }; type SubsPipelineConfig = { slugs: string[]; query: string; useRegex: boolean; regex: RegExp | null; // Tracks to exclude. Empty set means "search all tracks". excludedTracks: Set; // Optional inclusion filter. When provided, ONLY tracks in this set are // scanned (and `excludedTracks` is irrelevant). Used by the composite- // search "chat" scope to limit scanning to live_chat regardless of which // other tracks a video has. includedTracks?: Set | null; initialHitLimit: number; concurrency: number; flushIntervalMs: number; emit: (update: SubsPipelineUpdate) => void; fetchSubs: (slug: string) => Promise; }; export function createSubsSearchPipeline( config: SubsPipelineConfig, ): PipelineController { const { slugs, query, useRegex, regex, excludedTracks, includedTracks, initialHitLimit, concurrency, flushIntervalMs, emit, fetchSubs, } = config; let cancelled = false; let idx = 0; let completed = 0; let totalSoFar = 0; let hitLimit = initialHitLimit; let activeWorkers = 0; let done = false; const localHits: Record = {}; let flushTimer: ReturnType | null = null; const pushUpdate = (overrides: Partial = {}) => { emit({ hitsBySlug: { ...localHits }, totalHits: totalSoFar, processed: completed, totalToProcess: slugs.length, capped: totalSoFar >= hitLimit && idx < slugs.length, done, ...overrides, }); }; const scheduleFlush = () => { if (flushTimer !== null || cancelled) return; flushTimer = setTimeout(() => { flushTimer = null; if (cancelled) return; pushUpdate(); }, flushIntervalMs); }; const finalize = () => { if (done || cancelled) return; done = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } pushUpdate(); }; const worker = async () => { activeWorkers++; try { while (!cancelled) { if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const my = idx++; const slug = slugs[my]; try { const detail = await fetchSubs(slug); if (cancelled) return; if (totalSoFar < hitLimit) { const slugHits: SubsHit[] = []; for (const [track, cueList] of Object.entries(detail.tracks)) { if (includedTracks) { if (!includedTracks.has(track)) continue; } else if (excludedTracks.has(track)) continue; if (totalSoFar + slugHits.length >= hitLimit) break; const remaining = hitLimit - totalSoFar - slugHits.length; const trackHits = findHitsInCues( cueList, query, useRegex, regex, remaining, ); for (const h of trackHits) { slugHits.push({ track, start: h.start, text: h.text }); } } if (slugHits.length > 0) { // Order hits by start time so live_chat and language tracks // interleave naturally instead of being grouped per-track. slugHits.sort((a, b) => a.start - b.start); localHits[slug] = slugHits; totalSoFar += slugHits.length; } } } catch { // ignore per-video failures } completed++; scheduleFlush(); } } finally { activeWorkers--; if (activeWorkers === 0 && !cancelled) finalize(); } }; const ensureWorkers = () => { if (cancelled || done) return; if (totalSoFar >= hitLimit) return; if (idx >= slugs.length) return; const needed = Math.min(concurrency - activeWorkers, slugs.length - idx); for (let i = 0; i < needed; i++) worker(); }; // `worker()` counts itself in synchronously, before its first await, so // right after ensureWorkers a zero here means none is running or will. const settleIfIdle = () => { if (activeWorkers === 0 && !done) finalize(); }; pushUpdate(); ensureWorkers(); settleIfIdle(); // see createSearchPipeline return { cancel() { cancelled = true; if (flushTimer !== null) { clearTimeout(flushTimer); flushTimer = null; } }, setHitLimit(limit: number) { if (cancelled) return; if (limit <= hitLimit) return; hitLimit = limit; if (done) { done = false; pushUpdate(); } ensureWorkers(); settleIfIdle(); }, }; } // ─── Leaf pipeline wrapper ─── // Thin adapter over the three drivers above for the composite-search tree // (`lib/searchEval.ts`). One controller per leaf in the tree. Returns a // `LeafController` the tree orchestrator can cancel / resize, plus a `done` // promise that resolves with the final LeafResult. // // Scope=metadata isn't handled here — `searchEval.ts` evaluates it // synchronously over the summaries cache. This wrapper only deals with the // network-backed scopes (transcripts, chat, posts, description, tags) where the // worker pool earns its keep. export type LeafProgress = { slugs: Set; hits: Map; totalHits: number; processed: number; totalToProcess: number; capped: boolean; done: boolean; }; export type LeafResult = { slugs: Set; hits: Map; }; export type LeafController = PipelineController & { done: Promise; }; export type LeafPipelineOptions = { leaf: LeafNode; scopeSlugs: string[]; initialHitLimit: number; concurrency: number; flushIntervalMs: number; emit: (p: LeafProgress) => void; fetchers: LeafFetchers; }; // What `lib/searchEval.ts` is handed instead of importing a bound pipeline from // `components/`: a leaf runner with its fetchers already closed over. export type LeafRunner = ( opts: Omit, ) => LeafController; export function runLeafPipeline(opts: LeafPipelineOptions): LeafController { const { leaf, scopeSlugs, initialHitLimit, concurrency, flushIntervalMs, emit, fetchers, } = opts; const trimmed = leaf.query.trim(); // Caller guarantees query is non-empty before invoking us (see // isLeafActive). Stay defensive: empty queries produce an immediate done // with no hits. if (!trimmed) { const empty: LeafResult = { slugs: new Set(), hits: new Map() }; queueMicrotask(() => { emit({ slugs: empty.slugs, hits: empty.hits, totalHits: 0, processed: 0, totalToProcess: 0, capped: false, done: true, }); }); return { cancel() { /* no-op */ }, setHitLimit() { /* no-op */ }, done: Promise.resolve(empty), }; } const regex = compileLeafRegex(leaf); let resolveDone!: (r: LeafResult) => void; const donePromise = new Promise((resolve) => { resolveDone = resolve; }); // Latest seen progress, kept locally so `done` settles with the final // payload that the consumer already saw on the last emit. let finalSlugs = new Set(); let finalHits = new Map(); let settled = false; const adaptTranscriptUpdate = (u: PipelineUpdate): LeafProgress => { const slugs = new Set(); const hits = new Map(); for (const [slug, list] of Object.entries(u.hitsBySlug)) { if (!list || list.length === 0) continue; slugs.add(slug); if (leaf.contributeHits) { hits.set( slug, list.map((h) => ({ leafId: leaf.id, scope: leaf.scope, ...(h.track ? { track: h.track } : {}), start: h.start, text: h.text, })), ); } } return { slugs, hits, totalHits: u.totalHits, processed: u.processed, totalToProcess: u.totalToProcess, capped: u.capped, done: u.done, }; }; const adaptSubsUpdate = (u: SubsPipelineUpdate): LeafProgress => { const slugs = new Set(); const hits = new Map(); for (const [slug, list] of Object.entries(u.hitsBySlug)) { if (!list || list.length === 0) continue; slugs.add(slug); if (leaf.contributeHits) { hits.set( slug, list.map((h) => ({ leafId: leaf.id, scope: "chat" as const, track: h.track, start: h.start, text: h.text, })), ); } } return { slugs, hits, totalHits: u.totalHits, processed: u.processed, totalToProcess: u.totalToProcess, capped: u.capped, done: u.done, }; }; const onProgress = (p: LeafProgress) => { finalSlugs = p.slugs; finalHits = p.hits; emit(p); if (p.done && !settled) { settled = true; resolveDone({ slugs: finalSlugs, hits: finalHits }); } }; let controller: PipelineController; if ( leaf.scope === "transcripts" || leaf.scope === "description" || leaf.scope === "tags" ) { controller = createSearchPipeline({ slugs: scopeSlugs, query: trimmed, useRegex: leaf.useRegex, regex, initialHitLimit, concurrency, flushIntervalMs, fetchTranscript: fetchers.transcript, matchField: leaf.scope === "description" ? "description" : leaf.scope === "tags" ? "tags" : "cues", emit: (u) => onProgress(adaptTranscriptUpdate(u)), }); } else if (leaf.scope === "posts") { controller = createPostsSearchPipeline({ slugs: scopeSlugs, query: trimmed, useRegex: leaf.useRegex, regex, initialHitLimit, concurrency, flushIntervalMs, fetchPost: fetchers.post, // A post has no timeline, so every hit is `start: 0` — the same // convention description/tags/metadata hits already use. emit: (u) => onProgress(adaptTranscriptUpdate(u)), }); } else if (leaf.scope === "chat") { controller = createSubsSearchPipeline({ slugs: scopeSlugs, query: trimmed, useRegex: leaf.useRegex, regex, excludedTracks: EMPTY_TRACK_SET, includedTracks: CHAT_ONLY_TRACK_SET, initialHitLimit, concurrency, flushIntervalMs, fetchSubs: fetchers.subs, emit: (u) => onProgress(adaptSubsUpdate(u)), }); } else { // Metadata scope is handled by searchEval.ts directly. We shouldn't be // invoked here; resolve immediately as a safety net. const empty: LeafResult = { slugs: new Set(), hits: new Map() }; queueMicrotask(() => { onProgress({ slugs: empty.slugs, hits: empty.hits, totalHits: 0, processed: 0, totalToProcess: 0, capped: false, done: true, }); }); return { cancel() { /* no-op */ }, setHitLimit() { /* no-op */ }, done: donePromise, }; } return { cancel() { controller.cancel(); if (!settled) { settled = true; resolveDone({ slugs: finalSlugs, hits: finalHits }); } }, setHitLimit(limit: number) { controller.setHitLimit(limit); }, done: donePromise, }; } function compileLeafRegex(leaf: LeafNode): RegExp | null { if (!leaf.useRegex) return null; try { return new RegExp(leaf.query, "i"); } catch { return null; } } const EMPTY_TRACK_SET: Set = new Set(); const CHAT_ONLY_TRACK_SET: Set = new Set(["live_chat"]);