// Composite-search tree orchestrator — the SLUG-SET half of one search // pipeline. // // Walks a QueryNode tree (from `searchQuery.ts`), kicks off per-leaf // pipelines, runs metadata-scope leaves synchronously against summaries, // memoizes per-leaf results, and emits a streaming TreeProgress that the UI // renders. // // Its sibling is `lib/search/evalTree.ts`, which evaluates the SAME algebra // per RECORD, for a caller that already holds the record (the MCP server). The // two are not copies: one answers "which slugs survive", streaming, without // the corpus in memory; the other "does this record match". When the meaning of // AND / OR / negate changes it changes in both, and each file says so. // // THE FETCH AND THE MEMO ARE INJECTED (`SearchRuntime`). They were imported // from `components/searchPipeline` and `components/searchLayerCache`, which // made the model layer depend on the view layer — the back-edge this slice // deletes. `components/searchPipeline.ts` now supplies both as // `searchRuntime`; a caller with different fetches (a test, a node-side sweep) // supplies its own. // // Evaluation rules (see plan): // AND group → narrow scope left-to-right between children // OR group → run children in parallel against the group's scope, union // NOT (flag) → child's result becomes `parent_scope \ child_result` // Empty / inactive subtrees act as identity (don't narrow) // The root is always a GroupNode (AND by default). import { canonicalHash, hashSlugs, isGroup, isLeaf, isLeafActive, isNodeActive, type GroupNode, type LeafNode, type QueryNode, } from "./searchQuery"; import type { LayerHit, LeafController, LeafRunner, } from "./search/leafPipeline"; import type { DisplaySummary } from "./transcripts"; // One leaf's memoized result: the slugs it matched and the hits that prove it. // Keyed by `${canonicalHash(node)}__${scopeHash}` — same node + same input // scope = same result, wherever in the tree it sits. export type CachedResult = { slugs: ReadonlySet; hits: ReadonlyMap; }; // The memo, as a dependency. `components/searchLayerCache.ts` implements it // over an in-memory Map plus IndexedDB; a caller with no browser backs it with // a plain Map, or with four no-ops for a run that should not memoize at all. export type LayerCache = { key(nodeHash: string, scopeHash: string): string; getSync(key: string): CachedResult | null; get(key: string): Promise; put(key: string, result: CachedResult): void; }; // Everything `runQueryTree` needs from the layer below it. export type SearchRuntime = { runLeaf: LeafRunner; cache: LayerCache; }; export type LeafState = { slugCount: number; totalHits: number; processed: number; totalToProcess: number; capped: boolean; cached: boolean; // Whether the leaf pipeline is still actively producing results. active: boolean; }; export type GroupState = { slugCount: number; }; export type TreeProgress = { slugs: ReadonlySet; // Per-video, ordered list of hits from contributing leaves. Stable // identity per slug so React can reuse rendered hit rows; rebuilt on // every progress emit since it's cheap. hits: ReadonlyMap; leafStates: ReadonlyMap; groupStates: ReadonlyMap; done: boolean; capped: boolean; }; export type TreeController = { cancel(): void; setHitLimit(limit: number): void; }; type MetadataLeafScope = "transcripts" | "chat" | "metadata"; type EvalCtx = { runtime: SearchRuntime; initialHitLimit: number; concurrency: number; flushIntervalMs: number; metadataIndex: MetadataIndex; // Slugs that have any subs (live_chat) content. Chat-scope leaves narrow // their input scope by this set before running, so we don't waste fetches // on videos that don't have a chat track at all. chatScopeSlugs: ReadonlySet | null; // Every post slug in the corpus. Posts live in a DISJOINT slug namespace from // videos (`/` vs a video id), so the two must be kept // apart: a posts leaf narrows to this set, and every video-shaped leaf // subtracts it. Without that, an AND against the global scope would // intersect the two namespaces to nothing, and video leaves would waste a // fetch per post. Null = no posts corpus on this site. postScopeSlugs: ReadonlySet | null; // Per-leaf live result state. Reused for cache-hit instant emission and // for incremental re-evaluation on each leaf progress. leafResults: Map; hits: Map }>; leafStates: Map; controllers: Map; cancelled: boolean; hitLimit: number; emitRoot: () => void; }; // Pre-built metadata index from summaries. Built once per tree run and // reused across all metadata-scope leaves so re-evaluating a metadata leaf // edit is instant. type MetadataIndex = { bySlug: Map; // Lower-cased title and channel name for fast case-insensitive contains. titleLowerBySlug: Map; channelLowerBySlug: Map; }; function buildMetadataIndex(summaries: DisplaySummary[]): MetadataIndex { const bySlug = new Map(); const titleLowerBySlug = new Map(); const channelLowerBySlug = new Map(); for (const s of summaries) { bySlug.set(s.slug, s); titleLowerBySlug.set(s.slug, (s.title ?? "").toLowerCase()); channelLowerBySlug.set(s.slug, (s.channel ?? "").toLowerCase()); } return { bySlug, titleLowerBySlug, channelLowerBySlug }; } export function runQueryTree(opts: { root: GroupNode; // The fetch-bound leaf runner + layer memo. See SearchRuntime above. runtime: SearchRuntime; globalScope: string[]; summaries: DisplaySummary[]; chatScopeSlugs?: ReadonlySet | null; postScopeSlugs?: ReadonlySet | null; initialHitLimit: number; concurrency: number; flushIntervalMs: number; emit: (p: TreeProgress) => void; }): TreeController { const { root, runtime, globalScope, summaries, chatScopeSlugs = null, postScopeSlugs = null, initialHitLimit, concurrency, flushIntervalMs, emit, } = opts; const metadataIndex = buildMetadataIndex(summaries); let cancelled = false; let allDone = false; const ctx: EvalCtx = { runtime, initialHitLimit, concurrency, flushIntervalMs, metadataIndex, chatScopeSlugs, postScopeSlugs, leafResults: new Map(), leafStates: new Map(), controllers: new Map(), cancelled: false, hitLimit: initialHitLimit, emitRoot: () => emitRoot(), }; // Latest emission cache — used so React doesn't see identity churn on // unchanged slug sets between progress events. let lastEmitToken = ""; function emitRoot(): void { if (cancelled) return; const r = evaluateTree(root, new Set(globalScope), ctx); const groupStates = new Map(); collectGroupStates(root, new Set(globalScope), ctx, groupStates); const capped = anyLeafCapped(ctx); const progress: TreeProgress = { slugs: r.slugs, hits: r.hits, leafStates: ctx.leafStates, groupStates, done: allDone, capped, }; // Cheap signature to avoid duplicate emits when re-running root after // a no-op cache touch. const token = `${r.slugs.size}|${allDone ? 1 : 0}|${capped ? 1 : 0}|${ctx.leafStates.size}`; if (token === lastEmitToken && !allDone) { // still emit to refresh hits ordering; cheap, fine. } lastEmitToken = token; emit(progress); } // Kick off all leaf evaluations on the next microtask so callers wiring // listeners have a chance to subscribe before the first emit fires. Promise.resolve().then(async () => { if (cancelled) return; try { await orchestrate(root, new Set(globalScope), ctx); } catch { /* per-leaf errors already swallowed inside */ } if (cancelled) return; allDone = true; emitRoot(); }); return { cancel() { cancelled = true; ctx.cancelled = true; for (const ctrl of ctx.controllers.values()) ctrl.cancel(); }, setHitLimit(limit: number) { if (cancelled) return; if (limit <= ctx.hitLimit) return; ctx.hitLimit = limit; for (const ctrl of ctx.controllers.values()) ctrl.setHitLimit(limit); }, }; } // Orchestrate evaluation: actually run leaf pipelines (or metadata sync) and // populate ctx.leafResults / leafStates. AND children run sequentially with // scope narrowed between them; OR children run in parallel. async function orchestrate( node: QueryNode, parentScope: Set, ctx: EvalCtx, ): Promise> { if (ctx.cancelled) return new Set(); if (!isNodeActive(node)) { // Inactive subtree → pass-through. Don't narrow anything. return parentScope; } if (isLeaf(node)) { const slugs = await runLeaf(node, parentScope, ctx); if (node.negate) return diff(parentScope, slugs); return intersect(parentScope, slugs); } // Group if (node.op === "AND") { let scope = parentScope; for (const child of node.children) { if (ctx.cancelled) break; if (!isNodeActive(child)) continue; scope = await orchestrate(child, scope, ctx); } if (node.negate) return diff(parentScope, scope); return scope; } // OR const childSlugSets: Set[] = []; await Promise.all( node.children.map(async (child) => { if (ctx.cancelled) return; if (!isNodeActive(child)) return; const r = await orchestrate(child, parentScope, ctx); childSlugSets.push(r); }), ); const unioned = unionAll(childSlugSets); if (node.negate) return diff(parentScope, unioned); return unioned; } // Run a single leaf against the given scope. Checks cache; on miss, runs the // appropriate scope-specific evaluator. Updates ctx.leafResults / leafStates // and pushes a root re-evaluation via ctx.emitRoot on every progress event. async function runLeaf( leaf: LeafNode, parentScope: Set, ctx: EvalCtx, ): Promise> { if (!isLeafActive(leaf)) { // Inactive (empty) leaf — record empty result, treat as identity. setLeafState(ctx, leaf.id, { slugCount: 0, totalHits: 0, processed: 0, totalToProcess: 0, capped: false, cached: false, active: false, }); ctx.leafResults.set(leaf.id, { slugs: new Set(), hits: new Map() }); return parentScope; } if (leaf.scope === "metadata") { return runMetadataLeaf(leaf, parentScope, ctx); } // Pre-filter chat-scope leaves to slugs that actually have a live_chat // track. Saves wasted fetches and makes the cache scope hash stable // against changes that only affect non-chat-having channels. let effectiveScope = parentScope; if (leaf.scope === "chat" && ctx.chatScopeSlugs) { effectiveScope = intersect(parentScope, ctx.chatScopeSlugs); } // Keep the two corpora out of each other's way (see ctx.postScopeSlugs). if (ctx.postScopeSlugs) { effectiveScope = leaf.scope === "posts" ? intersect(effectiveScope, ctx.postScopeSlugs) : diff(effectiveScope, ctx.postScopeSlugs); } else if (leaf.scope === "posts") { effectiveScope = new Set(); } // Fully-network leaves (transcripts / chat). Try the layer cache first. const scopeArr = Array.from(effectiveScope); // Nothing to read — a posts leaf on a site or selection with no posts, a // chat leaf where no video in scope has chat, a video leaf over posts only. // "Search in" makes such leaves from a plain query, so settle it here, at // once and without a cache lookup: matched nothing, done. (The drivers // settle an empty scope too; this skips the cache round trip.) if (scopeArr.length === 0) { return applyCached( leaf, 0, { slugs: new Set(), hits: new Map() }, ctx, /*cached*/ false, ); } const scopeHash = hashSlugs(scopeArr); const key = ctx.runtime.cache.key(canonicalHash(leaf), scopeHash); const cachedSync = ctx.runtime.cache.getSync(key); if (cachedSync) { return applyCached(leaf, scopeArr.length, cachedSync, ctx, /*cached*/ true); } const cached = await ctx.runtime.cache.get(key); if (ctx.cancelled) return new Set(); if (cached) { return applyCached(leaf, scopeArr.length, cached, ctx, /*cached*/ true); } // Cache miss — kick off a real pipeline. Pure-filter leaves (no hit // contribution) get an unbounded cap so a downstream AND chain can // continue to narrow without partial-scope artifacts. const initialHitLimit = leaf.contributeHits ? ctx.hitLimit : Infinity; setLeafState(ctx, leaf.id, { slugCount: 0, totalHits: 0, processed: 0, totalToProcess: scopeArr.length, capped: false, cached: false, active: true, }); let lastCapped = false; const controller = ctx.runtime.runLeaf({ leaf, scopeSlugs: scopeArr, initialHitLimit, concurrency: ctx.concurrency, flushIntervalMs: ctx.flushIntervalMs, emit: (p) => { if (ctx.cancelled) return; ctx.leafResults.set(leaf.id, { slugs: new Set(p.slugs), hits: new Map(p.hits), }); setLeafState(ctx, leaf.id, { slugCount: p.slugs.size, totalHits: p.totalHits, processed: p.processed, totalToProcess: p.totalToProcess, capped: p.capped, cached: false, active: !p.done, }); lastCapped = p.capped; ctx.emitRoot(); }, }); ctx.controllers.set(leaf.id, controller); const result = await controller.done; ctx.controllers.delete(leaf.id); if (ctx.cancelled) return new Set(); // Persist final result only on a natural completion (not cap-limited and // not cancelled). A capped result is partial — caching it would prevent a // later "show more" from extending into the rest of the scope. if (!lastCapped) { ctx.runtime.cache.put(key, { slugs: new Set(result.slugs), hits: new Map(result.hits), }); } return new Set(result.slugs); } // `scopeSize` is the size of the scope the leaf READ — its effective scope, // not its parent's — so its progress is "N of N" for the slugs it covered. The // parent scope over-counted: a chat leaf narrowed to the videos with chat, or // a "Search in" copy, reported the whole parent as processed, and the folded // plain leaf's "searched N/M" started full and then fell back as its other // copies streamed. An empty scope reports 0 of 0, which the session's // progress readout skips. function applyCached( leaf: LeafNode, scopeSize: number, cached: CachedResult, ctx: EvalCtx, isCacheHit: boolean, ): Set { const slugs = new Set(); for (const s of cached.slugs) slugs.add(s); const hits = new Map(); // Cached hits carry the leafId from whatever session wrote them. That id // is module-counter-scoped (searchQuery.ts) and won't survive a reload, // so rewrite to the current leaf's id — otherwise the renderer's // bucket-by-leafId step in TranscriptSearch drops every cached hit. for (const [slug, list] of cached.hits) { hits.set( slug, list.map((h) => ({ ...h, leafId: leaf.id })), ); } ctx.leafResults.set(leaf.id, { slugs, hits }); setLeafState(ctx, leaf.id, { slugCount: slugs.size, totalHits: countHits(hits), processed: scopeSize, totalToProcess: scopeSize, capped: false, cached: isCacheHit, active: false, }); ctx.emitRoot(); return slugs; } function runMetadataLeaf( leaf: LeafNode, parentScope: Set, ctx: EvalCtx, ): Set { const matched = new Set(); const hits = new Map(); const q = leaf.query.trim(); if (!q) { ctx.leafResults.set(leaf.id, { slugs: matched, hits }); setLeafState(ctx, leaf.id, { slugCount: 0, totalHits: 0, processed: parentScope.size, totalToProcess: parentScope.size, capped: false, cached: false, active: false, }); return parentScope; } let re: RegExp | null = null; if (leaf.useRegex) { try { re = new RegExp(q, "i"); } catch { re = null; } if (!re) { ctx.leafResults.set(leaf.id, { slugs: matched, hits }); setLeafState(ctx, leaf.id, { slugCount: 0, totalHits: 0, processed: parentScope.size, totalToProcess: parentScope.size, capped: false, cached: false, active: false, }); return new Set(); } } const lower = re ? "" : q.toLowerCase(); for (const slug of parentScope) { const title = ctx.metadataIndex.titleLowerBySlug.get(slug); const channel = ctx.metadataIndex.channelLowerBySlug.get(slug); if (title === undefined && channel === undefined) continue; const titleHit = re ? title !== undefined && re.test(title) : title !== undefined && title.includes(lower); const channelHit = re ? channel !== undefined && re.test(channel) : channel !== undefined && channel.includes(lower); if (!titleHit && !channelHit) continue; matched.add(slug); if (leaf.contributeHits) { const summary = ctx.metadataIndex.bySlug.get(slug); const out: LayerHit[] = []; if (titleHit && summary) { out.push({ leafId: leaf.id, scope: "metadata", start: 0, text: summary.title, }); } if (channelHit && !titleHit && summary) { // Channel-only hits don't have a meaningful timestamp; surface them // with start=0 too. Result rendering treats start=0 as "no jump". out.push({ leafId: leaf.id, scope: "metadata", start: 0, text: `Channel: ${summary.channel}`, }); } if (out.length > 0) hits.set(slug, out); } } ctx.leafResults.set(leaf.id, { slugs: matched, hits }); setLeafState(ctx, leaf.id, { slugCount: matched.size, totalHits: countHits(hits), processed: parentScope.size, totalToProcess: parentScope.size, capped: false, cached: false, active: false, }); ctx.emitRoot(); // Negate handled by the caller in `orchestrate` — return raw matches. return matched; } // ─── Pure-tree re-evaluation ─── // Computes the current root slug set + hit map from ctx.leafResults. Called // on every leaf progress so the UI sees streaming narrowing. This is the // "what does the tree currently produce" view, separate from the // orchestration that actually fills leafResults. function evaluateTree( node: QueryNode, parentScope: Set, ctx: EvalCtx, ): { slugs: Set; hits: Map } { const hits = new Map(); const slugs = evaluateSlugs(node, parentScope, ctx, hits); // Filter hits to surviving slugs, so a leaf earlier in an AND chain // doesn't leak hits for videos that a later child has since excluded. for (const slug of Array.from(hits.keys())) { if (!slugs.has(slug)) hits.delete(slug); } // Sort each video's hits by start time so transcripts + chat interleave // chronologically across layers. for (const list of hits.values()) list.sort((a, b) => a.start - b.start); return { slugs, hits }; } function evaluateSlugs( node: QueryNode, parentScope: Set, ctx: EvalCtx, hitSink: Map, ): Set { if (!isNodeActive(node)) return parentScope; if (isLeaf(node)) { const r = ctx.leafResults.get(node.id); if (!r) { // Pipeline hasn't started yet — be permissive so AND results don't // briefly collapse to empty before the leaf has reported. return parentScope; } if (node.contributeHits) { for (const [slug, list] of r.hits) { const existing = hitSink.get(slug); if (existing) existing.push(...list); else hitSink.set(slug, list.slice()); } } if (node.negate) return diff(parentScope, r.slugs); return intersect(parentScope, r.slugs); } // Group if (node.op === "AND") { let scope = parentScope; for (const child of node.children) { if (!isNodeActive(child)) continue; scope = evaluateSlugs(child, scope, ctx, hitSink); } if (node.negate) return diff(parentScope, scope); return scope; } // OR const acc = new Set(); for (const child of node.children) { if (!isNodeActive(child)) continue; const r = evaluateSlugs(child, parentScope, ctx, hitSink); for (const s of r) acc.add(s); } if (node.negate) return diff(parentScope, acc); return acc; } function collectGroupStates( node: QueryNode, parentScope: Set, ctx: EvalCtx, out: Map, ): Set { if (!isNodeActive(node)) { if (isGroup(node)) out.set(node.id, { slugCount: parentScope.size }); return parentScope; } if (isLeaf(node)) { const r = ctx.leafResults.get(node.id); if (!r) return parentScope; if (node.negate) return diff(parentScope, r.slugs); return intersect(parentScope, r.slugs); } let scopeOut: Set; if (node.op === "AND") { let scope = parentScope; for (const child of node.children) { if (!isNodeActive(child)) continue; scope = collectGroupStates(child, scope, ctx, out); } scopeOut = scope; } else { const acc = new Set(); for (const child of node.children) { if (!isNodeActive(child)) continue; const r = collectGroupStates(child, parentScope, ctx, out); for (const s of r) acc.add(s); } scopeOut = acc; } const finalScope = node.negate ? diff(parentScope, scopeOut) : scopeOut; out.set(node.id, { slugCount: finalScope.size }); return finalScope; } function anyLeafCapped(ctx: EvalCtx): boolean { for (const s of ctx.leafStates.values()) if (s.capped) return true; return false; } function setLeafState( ctx: EvalCtx, id: string, state: LeafState, ): void { ctx.leafStates.set(id, state); } function countHits(hits: ReadonlyMap): number { let n = 0; for (const list of hits.values()) n += list.length; return n; } function intersect(a: Set, b: ReadonlySet): Set { if (a.size === 0 || b.size === 0) return new Set(); // Iterate the smaller set for cheaper membership lookups. const [small, large]: [ReadonlySet, ReadonlySet] = a.size <= b.size ? [a, b] : [b, a]; const out = new Set(); for (const s of small) if (large.has(s)) out.add(s); return out; } function diff(a: Set, b: ReadonlySet): Set { const out = new Set(); for (const s of a) if (!b.has(s)) out.add(s); return out; } function unionAll(sets: Set[]): Set { if (sets.length === 1) return new Set(sets[0]); const out = new Set(); for (const s of sets) for (const v of s) out.add(v); return out; } // Re-export the leaf-scope type for UI consumers. export type { MetadataLeafScope };