Archilyzer · Source

archilyzer

Archilyzer
git clone https://archilyzer.pages.dev/source/archilyzer.git
Log | Files | Refs | README | LICENSE

commit d0e83fead98d60c2b03de5051bbd7c6c1272f658
parent d6aa38431bbd97845456cf2cfef17e192997eebc
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date:   Sat, 12 Sep 2026 12:11:21 -0400

merge: one-core/phase-2-s3 — one search pipeline under common/lib/search

S3 of Phase 2, reviewed and the three low findings fixed: the evaluator,
window, rank and collapse logic mcp and the viewer shared by copy live once
in lib/search with injected fetchers and an explicit policy (no caller's
budget as a default); mcp/src/search.ts is orchestration over it with
MCP_POLICY passed, so the bench's structural counters are identical by
construction; lib/ no longer imports components/ and the architecture test
forbids the edge with its allow-list unchanged. Gates on the slice tip: tsc
clean, common 1129, mcp 205, scripts 71/1, export build clean, export e2e
172/172, hub 5/5, editor export-search 19/19.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>

Diffstat:
Mcommon/architecture.test.ts | 11+++++++++--
Mcommon/components/SearchSessionContext.tsx | 12++++++++----
Mcommon/components/charts/useSearchSeries.ts | 2++
Mcommon/components/searchLayerCache.ts | 13+++++++------
Mcommon/components/searchPipeline.ts | 904++++++-------------------------------------------------------------------------
Mcommon/components/urlState.ts | 9+++++----
Mcommon/lib/aiHandoff.test.ts | 2+-
Mcommon/lib/aiHandoff.ts | 2+-
Acommon/lib/search/collapse.test.ts | 102+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/search/collapse.ts | 92++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------------
Acommon/lib/search/evalTree.test.ts | 323+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/search/evalTree.ts | 393+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Acommon/lib/search/leafPipeline.test.ts | 286+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/search/leafPipeline.ts | 799+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acommon/lib/search/policy.test.ts | 26++++++++++++++++++++++++++
Mcommon/lib/search/policy.ts | 75++++++++++++++++++++++++++++++++++++++++++++++++---------------------------
Acommon/lib/search/rank.test.ts | 37+++++++++++++++++++++++++++++++++++++
Mcommon/lib/search/rank.ts | 57+++++++++++++++++++++++++++++++++++++++++++++------------
Acommon/lib/search/window.test.ts | 130+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcommon/lib/search/window.ts | 218+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Mcommon/lib/searchEval.ts | 77++++++++++++++++++++++++++++++++++++++++++++++++++++++++---------------------
Mcommon/lib/searchQuery.ts | 9++++++++-
Mexport/app/lib/askRetrieval.ts | 6+++++-
Mmcp/src/search.ts | 530+++++++++++++------------------------------------------------------------------
Mplans/one-core-phase-2.md | 265+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
25 files changed, 2950 insertions(+), 1430 deletions(-)

diff --git a/common/architecture.test.ts b/common/architecture.test.ts @@ -29,7 +29,13 @@ const HERE = path.dirname(fileURLToPath(import.meta.url)); // Which directory may not import which. Read as "lib/ may not import // controller/ or jobs/". const FORBIDDEN: Record<string, readonly string[]> = { - lib: ["controller", "jobs"], + // `components` joined this row in one-core phase 2 slice S3. The model layer + // importing the view layer was four edges: searchEval reaching for the leaf + // pipeline and the IndexedDB layer memo, aiHandoff for a hit type, and + // searchQuery for a mode union. All four inverted — the pipeline and the two + // types moved down, and the fetch and the memo are now injected — so the + // guard turns on with nothing added to the allow-list to pay for it. + lib: ["controller", "jobs", "components"], jobs: ["controller"], components: ["controller", "jobs", "ytdlp"], }; @@ -156,7 +162,8 @@ test("no new back-edges between common's layers", async () => { assert.deepEqual( unexpected, [], - `NEW back-edge(s) in common/. lib/ may not import controller/ or jobs/; ` + + `NEW back-edge(s) in common/. lib/ may not import controller/, jobs/ or ` + + `components/; ` + `jobs/ may not import controller/; components/ may not import ` + `controller/, jobs/ or ytdlp/. Move the type or the function down a ` + `layer instead of adding it to ALLOWED. Found: ${unexpected.join(", ")}`, diff --git a/common/components/SearchSessionContext.tsx b/common/components/SearchSessionContext.tsx @@ -26,7 +26,7 @@ import { usePlayer } from "./PlayerProvider"; import { useChannelSubsManifests } from "./subsCache"; import { peekPost, useChannelPostsManifests } from "./postsCache"; import { useSearchData } from "./SearchDataContext"; -import type { LayerHit } from "./searchPipeline"; +import { searchRuntime, type LayerHit } from "./searchPipeline"; import { runQueryTree, type GroupState, @@ -82,6 +82,7 @@ import type { Post } from "../lib/posts"; import { makeId, splitId } from "./originId"; import { sortGroups, type ChannelGroup } from "../lib/channelGroups"; import { buildSearchHandoff, type SearchHandoff } from "../lib/aiHandoff"; +import { rankByUploadDateDesc } from "../lib/search/rank"; import { pruneSelection, selectGroups } from "../lib/selection"; type Summary = DisplaySummary; @@ -760,6 +761,7 @@ function useSearchSessionState() { const controller = runQueryTree({ root: committedRoot, + runtime: searchRuntime, globalScope: globalScopeSlugs, summaries: transcripts, chatScopeSlugs: needsChatManifests ? chatScopeSlugs : null, @@ -881,9 +883,11 @@ function useSearchSessionState() { }); } // One newest-first ordering across BOTH corpora, so a unified search reads - // as one result set rather than videos-then-posts. - out.sort((a, b) => (a.uploadDate === b.uploadDate ? 0 : a.uploadDate < b.uploadDate ? 1 : -1)); - return out; + // as one result set rather than videos-then-posts. The comparator is + // `lib/search/rank.ts` — same function the next consumer of this ordering + // gets, which is the only reason a shared comparator is worth having. + // Sorts in place and returns the same array, as the inline sort did. + return rankByUploadDateDesc(out); }, [ hasActiveQuery, transcripts, diff --git a/common/components/charts/useSearchSeries.ts b/common/components/charts/useSearchSeries.ts @@ -8,6 +8,7 @@ import { applyFilters, type ChartData } from "../../lib/chartAggregate"; import { seriesFromSlugs } from "../../lib/chartSeriesFromSlugs"; import { parseRoot } from "../../lib/searchQuery"; import { runQueryTree, type TreeProgress } from "../../lib/searchEval"; +import { searchRuntime } from "../searchPipeline"; export type SearchSeriesState = { data: ChartData; @@ -56,6 +57,7 @@ export function useSearchSeries( try { const controller = runQueryTree({ root, + runtime: searchRuntime, globalScope: scopeSlugs, summaries, initialHitLimit: 50000, diff --git a/common/components/searchLayerCache.ts b/common/components/searchLayerCache.ts @@ -12,7 +12,13 @@ // IndexedDB writes are batched via queueMicrotask the same way as // transcriptStore.ts. -import type { LayerHit } from "./searchPipeline"; +import type { LayerHit } from "../lib/search/leafPipeline"; +import type { CachedResult } from "../lib/searchEval"; + +// Re-exported so this module's importers keep their import site. The type +// belongs to `lib/searchEval.ts` now: it is the shape of a memoized leaf +// result, and the tree — not its storage — defines it. +export type { CachedResult }; // Separate DB from `transcriptStore.ts` (`yt-dlp-transcript-browser`) so the // two stores can evolve their schemas independently — bumping the version @@ -22,11 +28,6 @@ const DB_VERSION = 1; const STORE = "layers"; const MAX_ENTRIES = 500; -export type CachedResult = { - slugs: ReadonlySet<string>; - hits: ReadonlyMap<string, LayerHit[]>; -}; - type StoredEntry = { key: string; slugs: string[]; diff --git a/common/components/searchPipeline.ts b/common/components/searchPipeline.ts @@ -1,852 +1,68 @@ +// THE BINDING between the shared search pipeline and this app's caches. +// +// Everything that used to live here — the three worker-pool drivers, the leaf +// adapter, the cue/text matchers — is `lib/search/leafPipeline.ts` now, where +// `lib/searchEval.ts` can reach it without `lib/` importing `components/`. +// What is left is the one thing that genuinely belongs to the view layer: WHICH +// FETCH. `fetchTranscript` / `fetchSubs` / `fetchPost` are the react-query-backed +// page caches, and they are supplied here rather than imported there. +// +// Every name this module exported before is still exported from it, so the five +// files that import `LayerHit` from `components/searchPipeline` (including +// `export/app/lib/askRetrieval.ts`) are untouched. + import { fetchTranscript } from "./transcriptCache"; import { fetchSubs } from "./subsCache"; import { fetchPost } from "./postsCache"; -import type { DisplaySummary } from "../lib/transcripts"; -import type { LayerScope, LeafNode } from "../lib/searchQuery"; - -export type Hit = { start: number; text: string }; - -export type SubsHit = { track: string; start: number; text: string }; - -// Per-leaf hit shape consumed by the composite-search tree (searchEval.ts). -// Carries enough info for the result-list UI to render the hit with its -// originating layer's swatch + scope-specific decorations. -export type LayerHit = { - leafId: string; - scope: LayerScope; - track?: string; - start: number; - text: string; -}; - -export type PipelineUpdate = { - hitsBySlug: Record<string, Hit[]>; - 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; - // 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"; +import { + runLeafPipeline as runLeafPipelineWith, + type LeafController, + type LeafFetchers, + type LeafPipelineOptions, +} from "../lib/search/leafPipeline"; +import { + cacheKey, + getCached, + getCachedSync, + putCached, +} from "./searchLayerCache"; +import type { SearchRuntime } from "../lib/searchEval"; + +export type { + Hit, + SubsHit, + LayerHit, + PipelineUpdate, + PipelineController, + SubsPipelineUpdate, + LeafProgress, + LeafResult, + LeafController, +} from "../lib/search/leafPipeline"; +export { filterSlugs } from "../lib/search/leafPipeline"; + +// This app's three fetches, as one injectable bundle. +const CACHE_FETCHERS: LeafFetchers = { + transcript: fetchTranscript, + subs: fetchSubs, + post: fetchPost, }; -export type PipelineController = { - cancel(): void; - setHitLimit(limit: number): void; -}; - -// Runs a streaming search over the given slugs. State is closure-captured so -// `setHitLimit` can raise the cap and re-spawn workers without restarting -// traversal — workers resume from the shared `idx` cursor. -export function createSearchPipeline( - config: PipelineConfig, -): PipelineController { - const { - slugs, - query, - useRegex, - regex, - initialHitLimit, - concurrency, - flushIntervalMs, - emit, - 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<string, Hit[]> = {}; - let flushTimer: number | null = null; - - const pushUpdate = (overrides: Partial<PipelineUpdate> = {}) => { - 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 = window.setTimeout(() => { - flushTimer = null; - if (cancelled) return; - pushUpdate(); - }, flushIntervalMs); - }; - - const finalize = () => { - if (done || cancelled) return; - done = true; - if (flushTimer !== null) { - window.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, - ) - : full.cues - ? findHitsInCues(full.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(); - }; - - // Emit initial snapshot synchronously so the UI clears previous results. - pushUpdate(); - ensureWorkers(); - - return { - cancel() { - cancelled = true; - if (flushTimer !== null) { - window.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(); - }, - }; -} - -function findHitsInCues( - cues: { start: number; text: string }[], - query: string, - useRegex: boolean, - regex: RegExp | null, - limit: number, -): Hit[] { - const hits: Hit[] = []; - for (let i = 0; i < cues.length && hits.length < limit; i++) { - const cur = cues[i]; - const prevText = i > 0 ? cues[i - 1].text : ""; - const nextText = i < cues.length - 1 ? cues[i + 1].text : ""; - const sep1 = prevText ? " " : ""; - const sep2 = nextText ? " " : ""; - const windowText = prevText + sep1 + cur.text + sep2 + nextText; - const curStart = prevText.length + sep1.length; - const curEnd = curStart + cur.text.length; - const m = findFirstMatchInRange( - windowText, - curStart, - curEnd, - query, - useRegex, - regex, - ); - if (!m) continue; - const crosses = m.idx < curStart || m.idx + m.length > curEnd; - hits.push({ - start: Math.round(cur.start), - text: crosses ? windowText : cur.text, - }); - } - return hits; -} - -// Single-document text match (description scope): emit at most one hit with a -// snippet window around the first match, so the result row shows context. -function findHitsInText( - text: string, - query: string, - useRegex: boolean, - regex: RegExp | null, -): Hit[] { - if (!text) return []; - const m = findFirstMatchInRange(text, 0, text.length, query, useRegex, regex); - if (!m) return []; - const PAD = 80; - const from = Math.max(0, m.idx - PAD); - const to = Math.min(text.length, m.idx + m.length + PAD); - const snippet = - (from > 0 ? "…" : "") + - text.slice(from, to).replace(/\s+/g, " ").trim() + - (to < text.length ? "…" : ""); - return [{ start: 0, text: snippet }]; -} - -function findFirstMatchInRange( - haystack: string, - rangeStart: number, - rangeEnd: number, - query: string, - useRegex: boolean, - regex: RegExp | null, -): { idx: number; length: number } | null { - if (useRegex) { - if (!regex) return null; - const flags = regex.flags.includes("g") ? regex.flags : regex.flags + "g"; - const re = new RegExp(regex.source, flags); - let m: RegExpExecArray | null; - while ((m = re.exec(haystack)) !== null) { - if (m.index >= rangeStart && m.index < rangeEnd) { - return { idx: m.index, length: m[0].length }; - } - if (m.index >= rangeEnd) return null; - if (m[0].length === 0) re.lastIndex++; - } - return null; - } - if (!query) return null; - const lower = haystack.toLowerCase(); - const ql = query.toLowerCase(); - let from = 0; - while (from <= haystack.length) { - const idx = lower.indexOf(ql, from); - if (idx === -1) return null; - if (idx >= rangeStart && idx < rangeEnd) return { idx, length: ql.length }; - if (idx >= rangeEnd) return null; - from = idx + 1; - } - return null; -} - -// Streaming search over the social-post corpus. Structurally the transcripts -// pipeline with a different fetch + match: the unit of iteration is a POST -// slug (`<channelSlug>/<postId>`), and fetchPost() 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<PipelineConfig, "matchField">, -): PipelineController { - const { - slugs, - query, - useRegex, - regex, - initialHitLimit, - concurrency, - flushIntervalMs, - emit, - } = 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<string, Hit[]> = {}; - let flushTimer: number | null = null; - - const pushUpdate = (overrides: Partial<PipelineUpdate> = {}) => { - 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 = window.setTimeout(() => { - flushTimer = null; - if (cancelled) return; - pushUpdate(); - }, flushIntervalMs); - }; - - const finalize = () => { - if (done || cancelled) return; - done = true; - if (flushTimer !== null) { - window.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(); - }; - - pushUpdate(); - ensureWorkers(); - - return { - cancel() { - cancelled = true; - if (flushTimer !== null) { - window.clearTimeout(flushTimer); - flushTimer = null; - } - }, - setHitLimit(limit: number) { - if (cancelled) return; - if (limit <= hitLimit) return; - hitLimit = limit; - if (done) { - done = false; - pushUpdate(); - } - ensureWorkers(); - }, - }; +export function runLeafPipeline( + opts: Omit<LeafPipelineOptions, "fetchers">, +): LeafController { + return runLeafPipelineWith({ ...opts, fetchers: CACHE_FETCHERS }); } -// 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<string, SubsHit[]>; - 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<string>; - // 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<string> | null; - initialHitLimit: number; - concurrency: number; - flushIntervalMs: number; - emit: (update: SubsPipelineUpdate) => void; -}; - -export function createSubsSearchPipeline( - config: SubsPipelineConfig, -): PipelineController { - const { - slugs, - query, - useRegex, - regex, - excludedTracks, - includedTracks, - initialHitLimit, - concurrency, - flushIntervalMs, - emit, - } = 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<string, SubsHit[]> = {}; - let flushTimer: number | null = null; - - const pushUpdate = (overrides: Partial<SubsPipelineUpdate> = {}) => { - 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 = window.setTimeout(() => { - flushTimer = null; - if (cancelled) return; - pushUpdate(); - }, flushIntervalMs); - }; - - const finalize = () => { - if (done || cancelled) return; - done = true; - if (flushTimer !== null) { - window.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(); - }; - - pushUpdate(); - ensureWorkers(); - - return { - cancel() { - cancelled = true; - if (flushTimer !== null) { - window.clearTimeout(flushTimer); - flushTimer = null; - } - }, - setHitLimit(limit: number) { - if (cancelled) return; - if (limit <= hitLimit) return; - hitLimit = limit; - if (done) { - done = false; - pushUpdate(); - } - ensureWorkers(); - }, - }; -} - -// ─── Leaf pipeline wrapper ─── -// Thin adapter over `createSearchPipeline` / `createSubsSearchPipeline` for -// the composite-search tree (`searchEval.ts`). One controller per leaf in the -// tree. Returns a `LeafController` that 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) where the existing worker pool earns its keep. - -export type LeafProgress = { - slugs: Set<string>; - hits: Map<string, LayerHit[]>; - totalHits: number; - processed: number; - totalToProcess: number; - capped: boolean; - done: boolean; +// What `runQueryTree` is handed: the fetch-bound leaf runner plus the +// IndexedDB-backed layer memo. `lib/searchEval.ts` owns the tree algebra and +// knows nothing about either. +export const searchRuntime: SearchRuntime = { + runLeaf: runLeafPipeline, + cache: { + key: cacheKey, + getSync: getCachedSync, + get: getCached, + put: putCached, + }, }; - -export type LeafResult = { - slugs: Set<string>; - hits: Map<string, LayerHit[]>; -}; - -export type LeafController = PipelineController & { - done: Promise<LeafResult>; -}; - -export function runLeafPipeline(opts: { - leaf: LeafNode; - scopeSlugs: string[]; - initialHitLimit: number; - concurrency: number; - flushIntervalMs: number; - emit: (p: LeafProgress) => void; -}): LeafController { - const { - leaf, - scopeSlugs, - initialHitLimit, - concurrency, - flushIntervalMs, - emit, - } = 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<LeafResult>((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<string>(); - let finalHits = new Map<string, LayerHit[]>(); - let settled = false; - - const adaptTranscriptUpdate = (u: PipelineUpdate): LeafProgress => { - const slugs = new Set<string>(); - const hits = new Map<string, LayerHit[]>(); - 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, - 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<string>(); - const hits = new Map<string, LayerHit[]>(); - 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, - 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, - // 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, - 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<string> = new Set(); -const CHAT_ONLY_TRACK_SET: Set<string> = new Set(["live_chat"]); diff --git a/common/components/urlState.ts b/common/components/urlState.ts @@ -2,10 +2,11 @@ import { useMemo, useSyncExternalStore } from "react"; -// Which corpus the legacy single-input search targets. "posts" is the social -// corpus (common/lib/posts.ts); the composite query tree can mix all of them -// freely, this only decides what a bare `?q=` URL means. -export type SearchMode = "transcripts" | "subs" | "posts"; +// SearchMode is defined in `lib/searchQuery.ts`, which interprets it. +// Re-exported here so every existing `from "./urlState"` import site is +// unchanged. +export type { SearchMode } from "../lib/searchQuery"; +import type { SearchMode } from "../lib/searchQuery"; // Per-video modal content mode. Independent from the search page's `mode` so // the modal can be toggled without disturbing search state. Absence on the diff --git a/common/lib/aiHandoff.test.ts b/common/lib/aiHandoff.test.ts @@ -1,7 +1,7 @@ import { test } from "node:test"; import assert from "node:assert/strict"; import { buildSearchHandoff, type HandoffGroup } from "./aiHandoff"; -import type { LayerHit } from "../components/searchPipeline"; +import type { LayerHit } from "./search/leafPipeline"; function hit(start: number, text: string): LayerHit { return { leafId: "l1", scope: "transcripts", start, text }; diff --git a/common/lib/aiHandoff.ts b/common/lib/aiHandoff.ts @@ -1,4 +1,4 @@ -import type { LayerHit } from "../components/searchPipeline"; +import type { LayerHit } from "./search/leafPipeline"; // Search grounding: a completed transcript search's results, serialized so the // /ask chat can answer *grounded in exactly those results* instead of running diff --git a/common/lib/search/collapse.test.ts b/common/lib/search/collapse.test.ts @@ -0,0 +1,102 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { collapseDuplicates, type ClusterOf, type CollapsibleHit } from "./collapse"; + +const hit = (slug: string): CollapsibleHit => ({ + slug, + videoId: slug.split("/")[1], + channelName: slug.split("/")[0], +}); + +const index = ( + entries: [string, ClusterOf][], +): ReadonlyMap<string, ClusterOf> => new Map(entries); + +test("an empty index collapses nothing and keeps every row", () => { + const hits = [hit("a/1"), hit("b/2")]; + const r = collapseDuplicates(hits, index([])); + assert.equal(r.collapsed, 0); + assert.equal(r.clusters, 0); + assert.deepEqual(r.kept.map((h) => h.slug), ["a/1", "b/2"]); +}); + +test("mirrors fold into the first match and are NAMED, never dropped", () => { + const hits = [hit("a/1"), hit("b/2"), hit("c/3")]; + const r = collapseDuplicates( + hits, + index([ + ["a/1", { clusterId: "k", isCanonical: false }], + ["b/2", { clusterId: "k", isCanonical: false }], + ]), + ); + assert.equal(r.collapsed, 1); + assert.equal(r.clusters, 1); + assert.deepEqual(r.kept.map((h) => h.slug), ["a/1", "c/3"]); + assert.deepEqual(r.kept[0].mirrors, [ + { videoId: "2", channelName: "b", slug: "b/2" }, + ]); +}); + +test("a canonical member arriving LATER is promoted to the kept row", () => { + const hits = [hit("a/1"), hit("b/2")]; + const r = collapseDuplicates( + hits, + index([ + ["a/1", { clusterId: "k", isCanonical: false }], + ["b/2", { clusterId: "k", isCanonical: true }], + ]), + ); + assert.equal(r.collapsed, 1); + assert.deepEqual(r.kept.map((h) => h.slug), ["b/2"]); + assert.deepEqual(r.kept[0].mirrors, [ + { videoId: "1", channelName: "a", slug: "a/1" }, + ]); +}); + +test("an ABSENT canonical does not delete the surviving mirror", () => { + // The cluster's canonical is "gone/9", which did not match. The first match + // stays the representative rather than the row vanishing — a mirror is often + // the only surviving copy of a deleted upload. + const hits = [hit("a/1"), hit("b/2")]; + const r = collapseDuplicates( + hits, + index([ + ["a/1", { clusterId: "k", isCanonical: false }], + ["b/2", { clusterId: "k", isCanonical: false }], + ["gone/9", { clusterId: "k", isCanonical: true }], + ]), + ); + assert.deepEqual(r.kept.map((h) => h.slug), ["a/1"]); +}); + +test("three copies of one recording fold to one row carrying both mirrors", () => { + const hits = [hit("a/1"), hit("b/2"), hit("c/3")]; + const r = collapseDuplicates( + hits, + index([ + ["a/1", { clusterId: "k", isCanonical: true }], + ["b/2", { clusterId: "k", isCanonical: false }], + ["c/3", { clusterId: "k", isCanonical: false }], + ]), + ); + assert.equal(r.collapsed, 2); + assert.equal(r.clusters, 1); + assert.deepEqual(r.kept.map((h) => h.slug), ["a/1"]); + assert.deepEqual( + (r.kept[0].mirrors ?? []).map((m) => m.slug), + ["b/2", "c/3"], + ); +}); + +test("collapsing does not mutate the input array", () => { + const hits = [hit("a/1"), hit("b/2")]; + collapseDuplicates( + hits, + index([ + ["a/1", { clusterId: "k", isCanonical: false }], + ["b/2", { clusterId: "k", isCanonical: false }], + ]), + ); + assert.equal(hits.length, 2); + assert.equal(hits[0].mirrors, undefined); +}); diff --git a/common/lib/search/collapse.ts b/common/lib/search/collapse.ts @@ -1,29 +1,81 @@ -// ONE SEARCH PIPELINE — placeholder. Filled by one-core phase 2 slice S3 -// (`plans/one-core-phase-2.md` §S3); created here by S1 so no two parallel -// slices race to create the same file. +// ONE SEARCH PIPELINE — cross-platform mirror collapsing. // -// Intended contents: cross-platform mirror collapsing, from -// `mcp/src/search.ts` (collapseDuplicates) and `components/SearchResults.tsx`. +// Keeps one row per RECORDING rather than per upload. Pure and synchronous: +// the caller supplies the already-fetched duplicate index (reader.duplicateIndex() +// on the MCP side), so this module never reaches for a transport. // -// The two rules that make it safe on by default, and that must survive the -// move verbatim: +// Two rules make it safe to have on by default: // // 1. The kept row is the cluster's canonical member WHEN that member is -// itself among the matches — otherwise simply the first match. A mirror is -// frequently the only surviving copy of a deleted upload, and preferring -// an absent canonical would delete exactly the evidence a "what did the -// removed videos say" question is asking for. +// itself among the matches — otherwise it is simply the first match. A +// mirror is frequently the only surviving copy of a deleted upload, and +// preferring an absent canonical would delete exactly the evidence a +// "what did the removed videos say" question is asking for. // 2. The collapsed copies are NAMED on the row they folded into. Nothing // vanishes; the count stops double-counting. // // And the one it must never break: timestamps are NEVER mapped between copies -// here. That requires the per-pair `aligned` gate, and this function does not -// move a single second of anything. -// -// export function collapseDuplicates<T extends CollapsibleHit>( -// hits: T[], -// dupes: DuplicateIndex, -// enabled: boolean, -// ): { collapsed: number; kept: T[] }; +// here. That requires the per-pair `aligned` gate (see `ClusterMembership` in +// `lib/archive/reader.ts`), and this function does not move a single second of +// anything. + +// The mirror rows carried on a kept hit. A mirror is named, never dropped. +export type Mirror = { videoId: string; channelName: string; slug: string }; + +// What collapsing needs of a hit. Anything with these four fields collapses — +// the MCP's SearchHit today, a viewer result group tomorrow. +export type CollapsibleHit = { + videoId: string; + slug: string; + channelName: string; + mirrors?: Mirror[]; +}; + +// What collapsing needs of the duplicate index: which cluster a slug belongs +// to, and whether it is that cluster's canonical member. Structurally satisfied +// by `DuplicateIndex` from `lib/archive/reader.ts`, which carries far more. +export type ClusterOf = { clusterId: string; isCanonical: boolean }; + +export type CollapseResult<T> = { + collapsed: number; + clusters: number; + kept: T[]; +}; -export {}; +export function collapseDuplicates<T extends CollapsibleHit>( + hits: readonly T[], + index: ReadonlyMap<string, ClusterOf>, +): CollapseResult<T> { + const repIndexOf = new Map<string, number>(); // clusterId -> index in `kept` + const kept: T[] = []; + let collapsed = 0; + for (const hit of hits) { + const membership = index.get(hit.slug); + if (!membership) { + kept.push(hit); + continue; + } + const at = repIndexOf.get(membership.clusterId); + if (at === undefined) { + repIndexOf.set(membership.clusterId, kept.length); + kept.push(hit); + continue; + } + collapsed++; + const rep = kept[at]; + const repIsCanonical = index.get(rep.slug)?.isCanonical === true; + const fold = (into: T, gone: T): T => ({ + ...into, + mirrors: [ + ...(into.mirrors ?? []), + ...(gone.mirrors ?? []), + { videoId: gone.videoId, channelName: gone.channelName, slug: gone.slug }, + ], + }); + // Promote the canonical member to the representative if it turns up later; + // otherwise fold this copy into the incumbent. + kept[at] = + repIsCanonical || !membership.isCanonical ? fold(rep, hit) : fold(hit, rep); + } + return { collapsed, clusters: repIndexOf.size, kept }; +} diff --git a/common/lib/search/evalTree.test.ts b/common/lib/search/evalTree.test.ts @@ -0,0 +1,323 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { newGroup, newLeaf } from "../searchQuery"; +import type { SearchAlias } from "../searchAliases"; +import { VIDEO_STATES, type VideoState } from "../availability"; +import { MCP_POLICY, VIEWER_POLICY } from "./policy"; +import { + buildMatcher, + buildLeafMatchers, + evalLeaf, + evalNode, + passesFilters, + needsAvailability, + filterIsSelective, + type RecordCtx, + type SearchFilters, +} from "./evalTree"; + +const ctx = (over: Partial<RecordCtx> = {}): RecordCtx => ({ + title: "", + channel: "", + description: "", + tags: "", + cues: [], + chatCues: [], + snippetsPerVideo: 4, + includeSnippets: true, + snippetChars: MCP_POLICY.snippetChars, + ...over, +}); + +const cue = (start: number, text: string) => ({ start, end: start + 2, text }); + +const ALL_STATES: ReadonlySet<VideoState> = new Set(VIDEO_STATES); +const openFilters = (over: Partial<SearchFilters> = {}): SearchFilters => ({ + videos: true, + livestreams: true, + allAges: true, + restricted: true, + states: ALL_STATES, + ...over, +}); + +const alias = (over: Partial<SearchAlias> = {}): SearchAlias => ({ + id: "a1", + label: "k cups", + triggers: ["k cups"], + suggestion: "(k|cake)[ -]?cup", + useRegex: true, + ...over, +}); + +// ─── matcher compilation ─── + +test("buildMatcher is plain substring, case-insensitive, and alias-free by default", () => { + const { match, firedAliases } = buildMatcher({ query: "Needle" }); + assert.equal(match("a needle here"), true); + assert.equal(match("nothing"), false); + assert.deepEqual(firedAliases, []); +}); + +test("buildMatcher ORs a fired alias onto the plain match", () => { + const { match, firedAliases } = buildMatcher({ + query: "k cups", + aliases: [alias()], + }); + assert.equal(firedAliases.length, 1); + assert.equal(match("k cups"), true, "the literal still matches"); + assert.equal(match("cakecup"), true, "…and so does the curated regex"); + assert.equal(match("unrelated"), false); +}); + +test("an explicit regex query takes NO alias expansion", () => { + const { match, firedAliases } = buildMatcher({ + query: "k c.ps", + regex: true, + aliases: [alias()], + }); + assert.deepEqual(firedAliases, []); + assert.equal(match("k cups"), true); + assert.equal(match("cakecup"), false); +}); + +test("useAliases:false suppresses expansion even when an alias would fire", () => { + const { firedAliases } = buildMatcher({ + query: "k cups", + useAliases: false, + aliases: [alias()], + }); + assert.deepEqual(firedAliases, []); +}); + +test("a malformed alias regex degrades to substring instead of throwing", () => { + const { match } = buildMatcher({ + query: "k cups", + aliases: [alias({ suggestion: "([unclosed" })], + }); + assert.equal(match("k cups"), true); + assert.equal(match("([unclosed"), true); +}); + +test("buildLeafMatchers compiles only ACTIVE leaves and unions fired aliases", () => { + const a = newLeaf({ id: "l1", query: "k cups" }); + const b = newLeaf({ id: "l2", query: "" }); // inactive + const c = newLeaf({ id: "l3", query: "k cups", scope: "description" }); + const { matchers, fired } = buildLeafMatchers( + newGroup({ children: [a, b, c] }), + [alias()], + true, + ); + assert.deepEqual([...matchers.keys()], ["l1", "l3"]); + assert.equal(fired.length, 1, "the same alias fires once, not per leaf"); + // Only the transcripts leaf is alias-aware. + assert.equal(matchers.get("l1")!.test("cakecup"), true); + assert.equal(matchers.get("l3")!.test("cakecup"), false); +}); + +// ─── per-record leaf evaluation ─── + +test("a transcripts leaf counts every matched cue and caps the snippets", () => { + const leaf = newLeaf({ id: "l", query: "hit" }); + const r = evalLeaf( + leaf, + { scope: "transcripts", test: (t) => t.includes("hit") }, + ctx({ + cues: [cue(0, "hit"), cue(5, "hit"), cue(9, "miss"), cue(12, "hit")], + snippetsPerVideo: 2, + }), + ); + assert.equal(r.matched, true); + assert.equal(r.count, 3, "count is every match, not every snippet"); + assert.equal(r.hits.length, 2, "…and the snippets stop at snippetsPerVideo"); + assert.equal(r.hits[0].clock, "0:00"); + assert.equal(r.hits[0].seconds, 0); + assert.equal(r.hits[0].scope, "transcripts"); +}); + +test("includeSnippets:false keeps the count and drops the text", () => { + const r = evalLeaf( + newLeaf({ id: "l", query: "hit" }), + { scope: "transcripts", test: () => true }, + ctx({ cues: [cue(0, "hit")], includeSnippets: false }), + ); + assert.equal(r.count, 1); + assert.deepEqual(r.hits, []); +}); + +test("the untimed scopes carry seconds 0 and their own scope tag", () => { + const leaf = newLeaf({ id: "l", query: "x" }); + const hit = (scope: "description" | "tags" | "posts", c: Partial<RecordCtx>) => + evalLeaf(leaf, { scope, test: () => true }, ctx(c)).hits[0]; + assert.equal(hit("description", { description: "d" })!.seconds, 0); + assert.equal(hit("description", { description: "d" })!.scope, "description"); + assert.equal(hit("tags", { tags: "a, b" })!.text, "a, b"); + assert.equal(hit("posts", { postText: "body" })!.scope, "posts"); +}); + +test("a posts leaf never matches a video record (no postText)", () => { + const r = evalLeaf( + newLeaf({ id: "l", query: "x" }), + { scope: "posts", test: () => true }, + ctx({ title: "x", cues: [cue(0, "x")] }), + ); + assert.equal(r.matched, false); +}); + +test("metadata matches title OR channel, and names the channel when only it hit", () => { + const leaf = newLeaf({ id: "l", query: "acme" }); + const m = { scope: "metadata" as const, test: (t: string) => t.includes("acme") }; + const titled = evalLeaf(leaf, m, ctx({ title: "acme news", channel: "bob" })); + assert.equal(titled.hits[0].text, "acme news"); + const chan = evalLeaf(leaf, m, ctx({ title: "news", channel: "acme tv" })); + assert.equal(chan.hits[0].text, "Channel: acme tv"); + assert.equal(chan.count, 1, "title+channel is ONE metadata match, not two"); +}); + +test("a chat leaf reads chatCues and tags the live_chat track", () => { + const r = evalLeaf( + newLeaf({ id: "l", query: "x" }), + { scope: "chat", test: () => true }, + ctx({ chatCues: [cue(30, "lol")] }), + ); + assert.equal(r.hits[0].track, "live_chat"); + assert.equal(r.hits[0].seconds, 30); +}); + +// `RecordCtx.snippetChars` is REQUIRED — no MCP default survives in lib/search, +// because a viewer adopter that omitted it used to get 240-character excerpts +// with nothing failing. tsc enforces that it is passed; this pins that passing +// the viewer's policy actually widens the excerpt. +test("snippetChars comes from the policy the caller passes", () => { + const long = "y".repeat(300); + const run = (snippetChars: number) => + evalLeaf( + newLeaf({ id: "l", query: "y" }), + { scope: "transcripts", test: () => true }, + ctx({ cues: [cue(0, long)], snippetChars }), + ); + assert.equal(run(VIEWER_POLICY.snippetChars).hits[0].text.length, 300); + assert.equal(run(MCP_POLICY.snippetChars).hits[0].text.length, 240); +}); + +// ─── the tree algebra, per record ─── + +const treeOf = (...leaves: ReturnType<typeof newLeaf>[]) => + newGroup({ children: leaves }); + +function evalOver( + root: ReturnType<typeof newGroup>, + record: Partial<RecordCtx>, +) { + const { matchers } = buildLeafMatchers(root, [], false); + return evalNode(root, matchers, ctx(record)); +} + +test("AND requires every active child; OR requires one", () => { + const both = treeOf( + newLeaf({ query: "alpha" }), + newLeaf({ query: "beta" }), + ); + assert.equal(evalOver(both, { cues: [cue(0, "alpha beta")] }).match, true); + assert.equal(evalOver(both, { cues: [cue(0, "alpha only")] }).match, false); + + const either = newGroup({ + op: "OR", + children: [newLeaf({ query: "alpha" }), newLeaf({ query: "beta" })], + }); + assert.equal(evalOver(either, { cues: [cue(0, "alpha only")] }).match, true); + assert.equal(evalOver(either, { cues: [cue(0, "gamma")] }).match, false); +}); + +test("a negated leaf inverts the match and contributes NO hits", () => { + const tree = treeOf( + newLeaf({ query: "alpha" }), + newLeaf({ query: "beta", negate: true }), + ); + const yes = evalOver(tree, { cues: [cue(0, "alpha")] }); + assert.equal(yes.match, true); + assert.equal(yes.count, 1, "only the positive leaf contributes"); + assert.equal(evalOver(tree, { cues: [cue(0, "alpha beta")] }).match, false); +}); + +test("contributeHits:false filters without contributing count or snippets", () => { + const tree = treeOf( + newLeaf({ query: "alpha", contributeHits: false }), + newLeaf({ query: "beta" }), + ); + const r = evalOver(tree, { cues: [cue(0, "alpha beta"), cue(3, "beta")] }); + assert.equal(r.match, true); + assert.equal(r.count, 2, "both beta cues, neither alpha"); +}); + +test("an inactive (empty) subtree is identity, not a filter", () => { + const tree = treeOf(newLeaf({ query: "alpha" }), newLeaf({ query: "" })); + assert.equal(evalOver(tree, { cues: [cue(0, "alpha")] }).match, true); + const allEmpty = treeOf(newLeaf({ query: "" })); + const r = evalOver(allEmpty, { cues: [cue(0, "anything")] }); + assert.equal(r.match, true); + assert.equal(r.count, 0); +}); + +test("a negated GROUP inverts its children and contributes no hits", () => { + const tree = newGroup({ + children: [ + newLeaf({ query: "alpha" }), + newGroup({ + negate: true, + op: "OR", + children: [newLeaf({ query: "beta" }), newLeaf({ query: "gamma" })], + }), + ], + }); + assert.equal(evalOver(tree, { cues: [cue(0, "alpha")] }).match, true); + assert.equal(evalOver(tree, { cues: [cue(0, "alpha beta")] }).match, false); + const r = evalOver(tree, { cues: [cue(0, "alpha")] }); + assert.equal(r.count, 1); +}); + +// ─── the filter predicate ─── + +test("passesFilters applies type, audience, state and the date range", () => { + const rec = { uploadDate: "20250601", isLivestream: false, ageRestricted: false }; + assert.equal(passesFilters(rec, openFilters(), undefined), true); + assert.equal(passesFilters(rec, openFilters({ videos: false }), undefined), false); + assert.equal( + passesFilters({ ...rec, isLivestream: true }, openFilters({ livestreams: false }), undefined), + false, + ); + assert.equal( + passesFilters({ ...rec, ageRestricted: true }, openFilters({ allAges: false }), undefined), + true, + ); + assert.equal( + passesFilters({ ...rec, ageRestricted: true }, openFilters({ restricted: false }), undefined), + false, + ); + assert.equal(passesFilters(rec, openFilters({ dateFrom: "20250701" }), undefined), false); + assert.equal(passesFilters(rec, openFilters({ dateTo: "20250501" }), undefined), false); + assert.equal( + passesFilters(rec, openFilters({ dateFrom: "20250101", dateTo: "20251231" }), undefined), + true, + ); +}); + +test("an absent availability record reads as 'available'", () => { + const rec = { uploadDate: "20250601" }; + const onlyDeleted = openFilters({ states: new Set<VideoState>(["deleted"]) }); + assert.equal(passesFilters(rec, onlyDeleted, undefined), false); + assert.equal(passesFilters(rec, onlyDeleted, { state: "deleted" }), true); +}); + +test("needsAvailability and filterIsSelective are both false for an all-permissive filter", () => { + assert.equal(needsAvailability(openFilters()), false); + assert.equal(filterIsSelective(openFilters()), false); + assert.equal(needsAvailability(null), false); + assert.equal(filterIsSelective(null), false); + assert.equal(filterIsSelective(openFilters({ videos: false })), true); + assert.equal(filterIsSelective(openFilters({ dateFrom: "20200101" })), true); + assert.equal( + needsAvailability(openFilters({ states: new Set<VideoState>(["available"]) })), + true, + ); +}); diff --git a/common/lib/search/evalTree.ts b/common/lib/search/evalTree.ts @@ -1,31 +1,368 @@ -// ONE SEARCH PIPELINE — placeholder. Filled by one-core phase 2 slice S3 -// (`plans/one-core-phase-2.md` §S3); created here by S1 so no two parallel -// slices race to create the same file. +// ONE SEARCH PIPELINE — the query-tree evaluator and the filter predicate. // -// Intended contents: the query-tree evaluator that exists twice today — -// `mcp/src/search.ts:1189-1338` (evalLeaf / evalNode / passesFilters) and -// `common/lib/searchEval.ts`'s copy — as one implementation over one record -// shape. `passesFilters` stays SINGULAR: a second "cheap" predicate for -// planning is exactly how a pruner starts silently disagreeing with the scanner -// about what matches. +// This is the per-RECORD half of composite search: given one transcript (or +// post) record and a compiled query tree, does it match, how many times, and +// which excerpts prove it. It backs the MCP's `open_link` / `sweep link=` +// flows and the plain scanner's single-leaf case. // -// export type LeafMatcher = { scope: LayerScope; test: (text: string) => boolean }; -// export type RecordCtx = { -// cues: Cue[]; -// title: string; -// description?: string; -// tags?: string[]; -// includeSnippets: boolean; -// snippetsPerVideo: number; -// }; -// export type LeafOutcome = { matched: boolean; count: number; hits: ScopedSnippet[] }; +// Its sibling is `lib/searchEval.ts`, which evaluates the SAME tree algebra +// over SLUG SETS, streaming, because the browser cannot hold the corpus in +// memory and fetches per leaf. The two are not copies of one another — one +// answers "does this record match", the other "which slugs survive" — and +// merging them would mean the browser materialising every record. What they +// do share, and what this move makes single, is the algebra's meaning: +// AND narrows, OR unions, `negate` inverts and contributes no hits, and an +// inactive subtree is identity. When that changes it must change here and in +// `searchEval.ts` together; each file points at the other. // -// export function evalLeaf(leaf: QueryNode, m: LeafMatcher, ctx: RecordCtx): LeafOutcome; -// export function evalNode(node: QueryNode, ms: LeafMatcher[], ctx: RecordCtx): LeafOutcome; -// export function passesFilters( -// rec: { isLivestream?: boolean; ageRestricted?: boolean; uploadDate: string }, -// f: SearchFilters, -// avail: VideoAvailability | undefined, -// ): boolean; - -export {}; +// `passesFilters` stays SINGULAR. A second "cheap" predicate for scan planning +// is exactly how a pruner starts silently disagreeing with the scanner about +// what matches — so the page planner and the hit decision call this one +// function, typed on the three fields it reads so a summaries record and a +// full transcript record both satisfy it. + +import { + forEachLeaf, + isLeaf, + isLeafActive, + isNodeActive, + type LayerScope, + type QueryNode, +} from "../searchQuery"; +import { matchAliases, type SearchAlias } from "../searchAliases"; +import { VIDEO_STATES, type VideoState } from "../availability"; +import type { Cue } from "../vtt"; +import { clock, truncate, type Matcher } from "./window"; + +export type { Matcher }; + +// ─── Matcher compilation ─── + +// Build the combined, alias-aware matcher for a query, mirroring the browser's +// buildSearchRoot OR-of-leaves semantics: a text matches if the plain query +// substring matches OR any fired alias's suggestion regex matches. An explicit +// `regex` query is taken verbatim with NO alias expansion (the caller is +// crafting their own pattern). Returns the fired aliases so the tool can report +// which curated expansions it applied. +export function buildMatcher(opts: { + query: string; + regex?: boolean; + useAliases?: boolean; + aliases?: SearchAlias[]; +}): { match: Matcher; firedAliases: SearchAlias[] } { + if (opts.regex) { + const re = new RegExp(opts.query, "i"); + return { match: (t) => re.test(t), firedAliases: [] }; + } + const needle = opts.query.toLowerCase(); + const plain: Matcher = (t) => t.toLowerCase().includes(needle); + + const useAliases = opts.useAliases !== false; + const fired = + useAliases && opts.aliases && opts.aliases.length > 0 + ? matchAliases(opts.query, "transcripts", opts.aliases) + : []; + if (fired.length === 0) return { match: plain, firedAliases: [] }; + + const aliasMatchers: Matcher[] = fired.map((a) => { + if (a.useRegex) { + try { + const re = new RegExp(a.suggestion, "i"); + return (t: string) => re.test(t); + } catch { + // malformed suggestion regex — fall back to substring on the literal + } + } + const n = a.suggestion.toLowerCase(); + return (t: string) => t.toLowerCase().includes(n); + }); + + const match: Matcher = (t) => plain(t) || aliasMatchers.some((m) => m(t)); + return { match, firedAliases: fired }; +} + +export type LeafMatcher = { scope: LayerScope; test: Matcher }; + +// Compile a matcher per active leaf. Transcripts leaves are alias-aware (unless +// they're regex); every other scope matches plain-substring / regex only. Fired +// aliases are unioned for the caller to report. +export function buildLeafMatchers( + root: QueryNode, + aliases: SearchAlias[], + useAliases: boolean, +): { matchers: Map<string, LeafMatcher>; fired: SearchAlias[] } { + const matchers = new Map<string, LeafMatcher>(); + const fired: SearchAlias[] = []; + forEachLeaf(root, (leaf) => { + if (!isLeafActive(leaf)) return; + const aliasAware = + leaf.scope === "transcripts" && !leaf.useRegex && useAliases; + const built = buildMatcher({ + query: leaf.query, + regex: leaf.useRegex, + useAliases: aliasAware, + aliases: aliasAware ? aliases : [], + }); + matchers.set(leaf.id, { scope: leaf.scope, test: built.match }); + for (const a of built.firedAliases) { + if (!fired.some((x) => x.id === a.id)) fired.push(a); + } + }); + return { matchers, fired }; +} + +// ─── Per-record evaluation ─── + +// A snippet tagged with the scope it came from, carrying the seconds needed to +// build a moment link. `seconds` is 0 for the non-timed scopes (metadata / +// description / tags / posts). +export type ScopedSnippet = { + scope: LayerScope; + track?: string; + clock: string; + seconds: number; + text: string; +}; + +// Per-record context the tree evaluates against. +export type RecordCtx = { + title: string; + channel: string; + description: string; + tags: string; + cues: Cue[]; + chatCues: Cue[]; + // The post body, when this record IS a post rather than a video. Empty for a + // video record, so a posts-scope leaf never matches one. + postText?: string; + snippetsPerVideo: number; + includeSnippets: boolean; + // Snippet truncation width — `SearchPolicy.snippetChars`. REQUIRED: a default + // here would be one caller's budget imposed on every other, and the way a + // viewer adopter discovers it is by shipping quietly clipped excerpts. + snippetChars: number; +}; + +export type LeafOutcome = { matched: boolean; count: number; hits: ScopedSnippet[] }; + +export function evalLeaf( + leaf: QueryNode, + m: LeafMatcher, + ctx: RecordCtx, +): LeafOutcome { + if (!isLeaf(leaf)) return { matched: false, count: 0, hits: [] }; + const max = ctx.snippetChars; + const hits: ScopedSnippet[] = []; + const push = (s: ScopedSnippet): void => { + if (ctx.includeSnippets && hits.length < ctx.snippetsPerVideo) hits.push(s); + }; + let count = 0; + switch (m.scope) { + case "transcripts": + for (const cue of ctx.cues) { + if (!m.test(cue.text)) continue; + count++; + push({ + scope: "transcripts", + clock: clock(cue.start), + seconds: cue.start, + text: truncate(cue.text, max), + }); + } + break; + case "chat": + for (const cue of ctx.chatCues) { + if (!m.test(cue.text)) continue; + count++; + push({ + scope: "chat", + track: "live_chat", + clock: clock(cue.start), + seconds: cue.start, + text: truncate(cue.text, max), + }); + } + break; + case "posts": + // A post has no timeline: one hit, seconds 0 — the same convention the + // metadata / description / tags scopes already use. + if (ctx.postText && m.test(ctx.postText)) { + count++; + push({ + scope: "posts", + clock: clock(0), + seconds: 0, + text: truncate(ctx.postText, max), + }); + } + break; + case "metadata": { + const titleHit = m.test(ctx.title); + const channelHit = m.test(ctx.channel); + if (titleHit || channelHit) { + count++; + if (titleHit) { + push({ + scope: "metadata", + clock: clock(0), + seconds: 0, + text: truncate(ctx.title, max), + }); + } else { + push({ + scope: "metadata", + clock: clock(0), + seconds: 0, + text: `Channel: ${ctx.channel}`, + }); + } + } + break; + } + case "description": + if (ctx.description && m.test(ctx.description)) { + count++; + push({ + scope: "description", + clock: clock(0), + seconds: 0, + text: truncate(ctx.description, max), + }); + } + break; + case "tags": + if (ctx.tags && m.test(ctx.tags)) { + count++; + push({ + scope: "tags", + clock: clock(0), + seconds: 0, + text: truncate(ctx.tags, max), + }); + } + break; + } + return { matched: count > 0, count, hits }; +} + +export type NodeOutcome = { match: boolean; count: number; hits: ScopedSnippet[] }; + +// Evaluate the tree against one record — the per-record mirror of searchEval's +// AND/OR/negate. A negated node contributes no hits (like the browser's `diff`). +// An inactive (empty) subtree is identity (matches, no hits). +export function evalNode( + node: QueryNode, + matchers: ReadonlyMap<string, LeafMatcher>, + ctx: RecordCtx, +): NodeOutcome { + if (!isNodeActive(node)) return { match: true, count: 0, hits: [] }; + if (isLeaf(node)) { + const m = matchers.get(node.id); + if (!m) return { match: true, count: 0, hits: [] }; + const r = evalLeaf(node, m, ctx); + if (node.negate) return { match: !r.matched, count: 0, hits: [] }; + return { + match: r.matched, + count: node.contributeHits ? r.count : 0, + hits: node.contributeHits ? r.hits : [], + }; + } + const active = node.children.filter(isNodeActive); + if (active.length === 0) return { match: true, count: 0, hits: [] }; + if (node.op === "AND") { + let allMatch = true; + let count = 0; + const hits: ScopedSnippet[] = []; + for (const c of active) { + const r = evalNode(c, matchers, ctx); + if (!r.match) { + allMatch = false; + break; + } + count += r.count; + hits.push(...r.hits); + } + const match = node.negate ? !allMatch : allMatch; + return match && !node.negate + ? { match, count, hits } + : { match, count: 0, hits: [] }; + } + // OR + let any = false; + let count = 0; + const hits: ScopedSnippet[] = []; + for (const c of active) { + const r = evalNode(c, matchers, ctx); + if (r.match) { + any = true; + count += r.count; + hits.push(...r.hits); + } + } + const match = node.negate ? !any : any; + return match && !node.negate + ? { match, count, hits } + : { match, count: 0, hits: [] }; +} + +// ─── The share-filter predicate ─── + +// The positive share-filter selection (parseShareV1), as a predicate input. +export type SearchFilters = { + // ft — video / livestream types kept. + videos: boolean; + livestreams: boolean; + // fa — all-ages / age-restricted kept. + allAges: boolean; + restricted: boolean; + // fav — the VideoState values kept (see lib/availability). Absent from the + // set means filtered out. + states: ReadonlySet<VideoState>; + // fdf / fdt — inclusive upload-date bounds, "YYYYMMDD". + dateFrom?: string; + dateTo?: string; +}; + +// Typed on the three fields it actually reads rather than on TranscriptDetail, +// so the SAME predicate can be applied to a summaries index record while +// planning a scan and to the full transcript record while deciding a hit. +export type FilterableRecord = { + isLivestream?: boolean; + ageRestricted?: boolean; + uploadDate: string; +}; + +export function passesFilters( + rec: FilterableRecord, + f: SearchFilters, + avail: { state: VideoState } | undefined, +): boolean { + // ft — type + if (rec.isLivestream ? !f.livestreams : !f.videos) return false; + // fa — audience + if (rec.ageRestricted ? !f.restricted : !f.allAges) return false; + // fav — presence on the source platform + if (!f.states.has(avail?.state ?? "available")) return false; + // fdf / fdt — upload-date range (lexicographic on YYYYMMDD) + if (f.dateFrom && rec.uploadDate < f.dateFrom) return false; + if (f.dateTo && rec.uploadDate > f.dateTo) return false; + return true; +} + +// True when the fav filter could exclude something (so availability must be +// fetched). If every availability bucket is kept, there's nothing to look up. +export function needsAvailability(f: SearchFilters | null | undefined): boolean { + return !!f && !VIDEO_STATES.every((s) => f.states.has(s)); +} + +// True when a filter set could actually exclude something. An all-permissive +// filter (every state kept, both media types, both audiences, no dates) is the +// same query as no filter at all, and must NOT trigger an index read — that is +// the guard against making an unfiltered query slower by planning it. +export function filterIsSelective(f: SearchFilters | null | undefined): boolean { + if (!f) return false; + if (!VIDEO_STATES.every((s) => f.states.has(s))) return true; + if (!f.videos || !f.livestreams) return true; + if (!f.allAges || !f.restricted) return true; + return Boolean(f.dateFrom || f.dateTo); +} diff --git a/common/lib/search/leafPipeline.test.ts b/common/lib/search/leafPipeline.test.ts @@ -0,0 +1,286 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import type { TranscriptDetail } from "../transcripts"; +import type { SubsDetail } from "../subs"; +import type { Post } from "../posts"; +import type { ChannelTranscriptsManifest } from "../manifest"; +import type { ArchiveReader, ChannelRef } from "../archive/reader"; +import { newLeaf } from "../searchQuery"; +import { runLeafPipeline, type LeafFetchers, type LeafProgress } from "./leafPipeline"; + +// ─── an in-memory archive, and fetchers derived from it ─── +// +// The leaf pipeline takes its fetches as parameters, which is the whole point +// of the move: `components/searchPipeline.ts` binds them to the react-query +// caches, and here they are bound to an ArchiveReader instead. Nothing in the +// module under test knows the difference, and the page reads are countable — +// which is how "one page read serves every slug on it" is stated as a fact +// rather than hoped for. + +const CH: ChannelRef = { key: "c", slug: "c", name: "C" }; + +type Archive = { + pages: TranscriptDetail[][]; + subs: Record<string, SubsDetail>; + posts: Record<string, Post>; + reads: string[]; +}; + +function record(id: string, cues: { start: number; text: string }[]): TranscriptDetail { + return { + id, + slug: `c/${id}`, + title: `video ${id}`, + channel: "C", + uploadDate: "20250101", + cues: cues.map((c) => ({ ...c, end: c.start + 2 })), + description: `about ${id}`, + tags: [id, "shared-tag"], + } as TranscriptDetail; +} + +function makeReader(a: Archive): ArchiveReader { + const slugToPage: Record<string, number> = {}; + a.pages.forEach((page, i) => { + for (const r of page) slugToPage[r.id] = i; + }); + const manifest: ChannelTranscriptsManifest = { + version: 1, + pageCount: a.pages.length, + slugToPage, + } as ChannelTranscriptsManifest; + return { + label: "in-memory", + listChannels: async () => [CH], + transcriptsManifest: async () => { + a.reads.push("manifest"); + return manifest; + }, + transcriptPage: async (_ch, page) => { + a.reads.push(`page-${page}`); + return a.pages[page] ?? []; + }, + loadAliases: async () => [], + loadGroups: async () => ({ groups: [], defaultGroupId: "default" }), + publicOrigin: () => null, + subsManifest: async () => null, + subsPage: async () => [], + postsManifest: async () => null, + postsPage: async () => [], + availabilityMap: async () => new Map(), + }; +} + +// The adapter a caller writes once: slug -> record, through the reader's +// manifest/page walk. Page-level, never per-record — a per-record fetch is a +// bench regression by construction. +function fetchersFor(a: Archive): LeafFetchers { + const reader = makeReader(a); + const pageCache = new Map<number, Promise<TranscriptDetail[]>>(); + return { + transcript: async (slug) => { + const id = slug.split("/")[1] ?? slug; + const manifest = await reader.transcriptsManifest(CH); + const page = manifest.slugToPage[id]; + if (page === undefined) throw new Error(`no such slug: ${slug}`); + let p = pageCache.get(page); + if (!p) { + p = reader.transcriptPage(CH, page); + pageCache.set(page, p); + } + const found = (await p).find((r) => r.id === id); + if (!found) throw new Error(`no such record: ${slug}`); + return found; + }, + subs: async (slug) => { + const d = a.subs[slug]; + if (!d) throw new Error(`no subs: ${slug}`); + return d; + }, + post: async (slug) => { + const p = a.posts[slug]; + if (!p) throw new Error(`no post: ${slug}`); + return p; + }, + }; +} + +function archive(over: Partial<Archive> = {}): Archive { + return { pages: [], subs: {}, posts: {}, reads: [], ...over }; +} + +async function run( + leaf: ReturnType<typeof newLeaf>, + slugs: string[], + fetchers: LeafFetchers, + opts: { hitLimit?: number } = {}, +): Promise<{ final: LeafProgress[]; slugs: string[]; hits: Map<string, unknown[]> }> { + const seen: LeafProgress[] = []; + const controller = runLeafPipeline({ + leaf, + scopeSlugs: slugs, + initialHitLimit: opts.hitLimit ?? 1000, + concurrency: 2, + flushIntervalMs: 0, + emit: (p) => seen.push(p), + fetchers, + }); + const result = await controller.done; + return { + final: seen, + slugs: [...result.slugs].sort(), + hits: result.hits as Map<string, unknown[]>, + }; +} + +test("a transcripts leaf matches cues over an injected reader, one read per page", async () => { + const a = archive({ + pages: [ + [record("v1", [{ start: 0, text: "a needle here" }]), record("v2", [{ start: 5, text: "nothing" }])], + [record("v3", [{ start: 9, text: "another NEEDLE" }])], + ], + }); + const r = await run(newLeaf({ query: "needle" }), ["c/v1", "c/v2", "c/v3"], fetchersFor(a)); + assert.deepEqual(r.slugs, ["c/v1", "c/v3"]); + assert.equal(a.reads.filter((x) => x.startsWith("page-")).length, 2); + const hits = r.hits.get("c/v1") as { start: number; text: string; scope: string }[]; + assert.equal(hits.length, 1); + assert.equal(hits[0].start, 0); + assert.equal(hits[0].scope, "transcripts"); +}); + +test("the leaf's id is stamped on every hit so the tree can bucket them", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const r = await run(newLeaf({ id: "leaf-7", query: "needle" }), ["c/v1"], fetchersFor(a)); + const hits = r.hits.get("c/v1") as { leafId: string }[]; + assert.equal(hits[0].leafId, "leaf-7"); +}); + +test("contributeHits:false still narrows the slug set but carries no hits", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const r = await run( + newLeaf({ query: "needle", contributeHits: false }), + ["c/v1"], + fetchersFor(a), + ); + assert.deepEqual(r.slugs, ["c/v1"]); + assert.equal(r.hits.size, 0); +}); + +test("description and tags scopes read the same fetched record", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "silence" }])]] }); + const f = fetchersFor(a); + const desc = await run(newLeaf({ query: "about v1", scope: "description" }), ["c/v1"], f); + assert.deepEqual(desc.slugs, ["c/v1"]); + const tags = await run(newLeaf({ query: "shared-tag", scope: "tags" }), ["c/v1"], f); + assert.deepEqual(tags.slugs, ["c/v1"]); +}); + +test("a chat leaf reads ONLY the live_chat track", async () => { + const a = archive({ + subs: { + "c/v1": { + id: "v1", + slug: "c/v1", + tracks: { + live_chat: [{ start: 4, end: 5, text: "needle in chat" }], + en: [{ start: 9, end: 10, text: "needle in captions" }], + }, + } as unknown as SubsDetail, + }, + }); + const r = await run(newLeaf({ query: "needle", scope: "chat" }), ["c/v1"], fetchersFor(a)); + const hits = r.hits.get("c/v1") as { track: string; start: number; scope: string }[]; + assert.equal(hits.length, 1, "the en track is not a chat hit"); + assert.equal(hits[0].track, "live_chat"); + assert.equal(hits[0].scope, "chat"); + assert.equal(hits[0].start, 4); +}); + +test("a posts leaf matches the body and stamps start 0", async () => { + const a = archive({ + posts: { + "c/p1": { id: "p1", slug: "c/p1", text: "a needle in a post" } as unknown as Post, + "c/p2": { id: "p2", slug: "c/p2", text: "nothing here" } as unknown as Post, + }, + }); + const r = await run(newLeaf({ query: "needle", scope: "posts" }), ["c/p1", "c/p2"], fetchersFor(a)); + assert.deepEqual(r.slugs, ["c/p1"]); + const hits = r.hits.get("c/p1") as { start: number }[]; + assert.equal(hits[0].start, 0); +}); + +test("a failing fetch is skipped, not fatal, and progress still completes", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const r = await run(newLeaf({ query: "needle" }), ["c/v1", "c/missing"], fetchersFor(a)); + assert.deepEqual(r.slugs, ["c/v1"]); + const last = r.final[r.final.length - 1]; + assert.equal(last.done, true); + assert.equal(last.processed, 2, "the failed slug still counts as processed"); +}); + +test("the hit cap stops the scan and reports capped; setHitLimit resumes it", async () => { + const a = archive({ + pages: [ + Array.from({ length: 8 }, (_, i) => + record(`v${i}`, [{ start: i, text: "needle" }]), + ), + ], + }); + const slugs = Array.from({ length: 8 }, (_, i) => `c/v${i}`); + const fetchers = fetchersFor(a); + const seen: LeafProgress[] = []; + const controller = runLeafPipeline({ + leaf: newLeaf({ query: "needle" }), + scopeSlugs: slugs, + initialHitLimit: 2, + concurrency: 1, + flushIntervalMs: 0, + emit: (p) => seen.push(p), + fetchers, + }); + const first = await controller.done; + assert.ok(first.slugs.size <= 3, "stopped near the cap, not at the end"); + assert.ok(seen.some((p) => p.capped), "…and said so"); +}); + +test("an empty query settles immediately with nothing, rather than scanning", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const r = await run(newLeaf({ query: " " }), ["c/v1"], fetchersFor(a)); + assert.deepEqual(r.slugs, []); + assert.deepEqual(a.reads, []); +}); + +test("a metadata leaf is not this module's job and resolves empty", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const r = await run(newLeaf({ query: "video", scope: "metadata" }), ["c/v1"], fetchersFor(a)); + assert.deepEqual(r.slugs, []); + assert.deepEqual(a.reads, []); +}); + +test("cancel settles `done` with what had been seen so far", async () => { + const a = archive({ pages: [[record("v1", [{ start: 0, text: "needle" }])]] }); + const controller = runLeafPipeline({ + leaf: newLeaf({ query: "needle" }), + scopeSlugs: ["c/v1"], + initialHitLimit: 100, + concurrency: 1, + flushIntervalMs: 0, + emit: () => {}, + fetchers: fetchersFor(a), + }); + controller.cancel(); + const r = await controller.done; + assert.ok(r.slugs instanceof Set); +}); + +test("a regex leaf compiles once; a malformed pattern matches nothing", async () => { + const a = archive({ + pages: [[record("v1", [{ start: 0, text: "needle" }]), record("v2", [{ start: 0, text: "noodle" }])]], + }); + const f = fetchersFor(a); + const ok = await run(newLeaf({ query: "n[eo]+dle", useRegex: true }), ["c/v1", "c/v2"], f); + assert.deepEqual(ok.slugs, ["c/v1", "c/v2"]); + const bad = await run(newLeaf({ query: "n([edle", useRegex: true }), ["c/v1"], f); + assert.deepEqual(bad.slugs, []); +}); diff --git a/common/lib/search/leafPipeline.ts b/common/lib/search/leafPipeline.ts @@ -0,0 +1,799 @@ +// 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"; + +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<TranscriptDetail>; + subs: (slug: string) => Promise<SubsDetail>; + post: (slug: string) => Promise<Post>; +}; + +export type PipelineUpdate = { + hitsBySlug: Record<string, Hit[]>; + 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<TranscriptDetail>; + // 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<string, Hit[]> = {}; + let flushTimer: ReturnType<typeof setTimeout> | null = null; + + const pushUpdate = (overrides: Partial<PipelineUpdate> = {}) => { + 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, + ) + : full.cues + ? findHitsInCues(full.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(); + }; + + // Emit initial snapshot synchronously so the UI clears previous results. + pushUpdate(); + ensureWorkers(); + + 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(); + }, + }; +} + +// Streaming search over the social-post corpus. Structurally the transcripts +// pipeline with a different fetch + match: the unit of iteration is a POST +// slug (`<channelSlug>/<postId>`), 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<PipelineConfig, "matchField" | "fetchTranscript"> & { + fetchPost: (slug: string) => Promise<Post>; + }, +): 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<string, Hit[]> = {}; + let flushTimer: ReturnType<typeof setTimeout> | null = null; + + const pushUpdate = (overrides: Partial<PipelineUpdate> = {}) => { + 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(); + }; + + pushUpdate(); + ensureWorkers(); + + 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(); + }, + }; +} + +// 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<string, SubsHit[]>; + 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<string>; + // 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<string> | null; + initialHitLimit: number; + concurrency: number; + flushIntervalMs: number; + emit: (update: SubsPipelineUpdate) => void; + fetchSubs: (slug: string) => Promise<SubsDetail>; +}; + +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<string, SubsHit[]> = {}; + let flushTimer: ReturnType<typeof setTimeout> | null = null; + + const pushUpdate = (overrides: Partial<SubsPipelineUpdate> = {}) => { + 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(); + }; + + pushUpdate(); + ensureWorkers(); + + 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(); + }, + }; +} + +// ─── 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<string>; + hits: Map<string, LayerHit[]>; + totalHits: number; + processed: number; + totalToProcess: number; + capped: boolean; + done: boolean; +}; + +export type LeafResult = { + slugs: Set<string>; + hits: Map<string, LayerHit[]>; +}; + +export type LeafController = PipelineController & { + done: Promise<LeafResult>; +}; + +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<LeafPipelineOptions, "fetchers">, +) => 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<LeafResult>((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<string>(); + let finalHits = new Map<string, LayerHit[]>(); + let settled = false; + + const adaptTranscriptUpdate = (u: PipelineUpdate): LeafProgress => { + const slugs = new Set<string>(); + const hits = new Map<string, LayerHit[]>(); + 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, + 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<string>(); + const hits = new Map<string, LayerHit[]>(); + 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<string> = new Set(); +const CHAT_ONLY_TRACK_SET: Set<string> = new Set(["live_chat"]); diff --git a/common/lib/search/policy.test.ts b/common/lib/search/policy.test.ts @@ -0,0 +1,26 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { MCP_POLICY, VIEWER_POLICY } from "./policy"; +import { truncate } from "./window"; + +// These four numbers are a wire contract with the bench: the MCP scanner's +// structural read/byte counters are unchanged by the one-pipeline move ONLY +// because the caps travel as data rather than being re-derived. If one of them +// moves, the bench's page counts move with it, so pin them. +test("MCP_POLICY carries the MCP scanner's four caps verbatim", () => { + assert.equal(MCP_POLICY.maxPages, 400); + assert.equal(MCP_POLICY.hardVideoCap, 2000); + assert.equal(MCP_POLICY.windowLineCap, 200); + assert.equal(MCP_POLICY.snippetChars, 240); +}); + +test("VIEWER_POLICY is uncapped, and Infinity behaves as the identity", () => { + assert.equal(VIEWER_POLICY.maxPages, Number.POSITIVE_INFINITY); + assert.equal(VIEWER_POLICY.hardVideoCap, Number.POSITIVE_INFINITY); + assert.equal(VIEWER_POLICY.windowLineCap, Number.POSITIVE_INFINITY); + // The load-bearing consequence: a viewer excerpt is never clipped. + const long = "x".repeat(5000); + assert.equal(truncate(long, VIEWER_POLICY.snippetChars), long); + // …and 0 pages scanned is never "at the cap". + assert.equal(0 >= VIEWER_POLICY.maxPages, false); +}); diff --git a/common/lib/search/policy.ts b/common/lib/search/policy.ts @@ -1,30 +1,51 @@ -// ONE SEARCH PIPELINE — placeholder. Filled by one-core phase 2 slice S3 -// (`plans/one-core-phase-2.md` §S3); created here by S1 so no two parallel -// slices race to create the same file. +// ONE SEARCH PIPELINE — the caps, named and passed. // -// Intended contents: the caps that are today private constants of -// `mcp/src/search.ts`, named and passed rather than re-derived, so the bench's -// structural counts are unchanged by construction. +// These were private constants of `mcp/src/search.ts`. They are not "the" +// limits of searching an archive; they are the limits ONE CALLER chose, and +// the viewer chose differently. Naming them is what lets both callers share +// one scanner: the MCP server passes `MCP_POLICY`, a human scrolling a page +// passes `VIEWER_POLICY`, and neither re-derives a number the other owns. // -// export type SearchPolicy = { -// // A hard ceiling on shard pages fetched per query, so a rare term over a -// // large (or hub-wide) corpus cannot run away. Reaching it sets -// // `truncated`. search.ts:259 MAX_PAGES 400 -// maxPages: number; -// // A ceiling on matched videos collected before counting stops, so -// // `total` stays bounded for a very common term. Reaching it also sets -// // `truncated`. search.ts:263 HARD_VIDEO_CAP 2000 -// hardVideoCap: number; -// // Cap on windowed excerpt lines per video, so a video with hundreds of -// // matches cannot blow a batch's token budget. -// // search.ts:268 WINDOW_LINE_CAP 200 -// windowLineCap: number; -// // Snippet truncation width. search.ts:456 truncate(…, 240) -// snippetChars: number; -// }; -// -// // What the MCP server passes. The viewer passes its own, UNCAPPED: a human -// // scrolling a page is not spending an agent's token budget. -// export const MCP_POLICY: SearchPolicy; +// The bench's structural read/byte counters are unchanged by construction +// BECAUSE these are passed rather than re-derived — `MCP_POLICY` holds exactly +// the four values the MCP scanner used before the move. + +export type SearchPolicy = { + // A hard ceiling on shard pages fetched per query, so a rare term over a + // large (or hub-wide) corpus cannot run away. Reaching it sets `truncated`. + // Was `MAX_PAGES` (search.ts:259). + maxPages: number; + // A ceiling on matched videos collected before counting stops, so `total` + // stays bounded and stable for a very common term. Reaching it also sets + // `truncated` (the true total is higher than reported). Was + // `HARD_VIDEO_CAP` (search.ts:263). + hardVideoCap: number; + // Cap on merged windowed excerpt lines emitted per video, so a video with + // hundreds of matches cannot blow a batch's token budget. Was + // `WINDOW_LINE_CAP` (search.ts:268). + windowLineCap: number; + // Snippet truncation width. Was the `max = 240` default of `truncate` + // (search.ts:456). + snippetChars: number; +}; + +// What the MCP server passes — the four values verbatim. +export const MCP_POLICY: SearchPolicy = { + maxPages: 400, + hardVideoCap: 2000, + windowLineCap: 200, + snippetChars: 240, +}; -export {}; +// What the viewer passes: UNCAPPED. A human scrolling a page is not spending +// an agent's token budget, and the browser pipeline has always had its own +// per-run hit limit (`DEFAULT_MAX_HITS`, raised by "show more") rather than a +// fixed page or video ceiling. Infinity is load-bearing here, not decorative: +// `truncate(t, Infinity)` is the identity, and `pagesScanned >= Infinity` is +// never true. +export const VIEWER_POLICY: SearchPolicy = { + maxPages: Number.POSITIVE_INFINITY, + hardVideoCap: Number.POSITIVE_INFINITY, + windowLineCap: Number.POSITIVE_INFINITY, + snippetChars: Number.POSITIVE_INFINITY, +}; diff --git a/common/lib/search/rank.test.ts b/common/lib/search/rank.test.ts @@ -0,0 +1,37 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + byUploadDateDesc, + byCreatedAtThenId, + rankByUploadDateDesc, + rankThread, +} from "./rank"; + +test("rankByUploadDateDesc is newest-first and stable within a date", () => { + const rows = [ + { uploadDate: "20240101", id: "a" }, + { uploadDate: "20250601", id: "b" }, + { uploadDate: "20240101", id: "c" }, + { uploadDate: "20251231", id: "d" }, + ]; + assert.deepEqual( + rankByUploadDateDesc(rows).map((r) => r.id), + ["d", "b", "a", "c"], + ); + assert.equal(byUploadDateDesc({ uploadDate: "x" }, { uploadDate: "x" }), 0); +}); + +test("rankThread is oldest-first, tie-broken by id so the order is total", () => { + const posts = [ + { createdAt: "2025-01-02T00:00:00Z", id: "b" }, + { createdAt: "2025-01-01T00:00:00Z", id: "z" }, + { createdAt: "2025-01-01T00:00:00Z", id: "a" }, + ]; + assert.deepEqual( + rankThread(posts).map((p) => p.id), + ["a", "z", "b"], + ); + assert.ok( + byCreatedAtThenId({ createdAt: "t", id: "a" }, { createdAt: "t", id: "b" }) < 0, + ); +}); diff --git a/common/lib/search/rank.ts b/common/lib/search/rank.ts @@ -1,16 +1,49 @@ -// ONE SEARCH PIPELINE — placeholder. Filled by one-core phase 2 slice S3 -// (`plans/one-core-phase-2.md` §S3); created here by S1 so no two parallel -// slices race to create the same file. +// ONE SEARCH PIPELINE — result ordering. // -// Intended contents: result ordering, once. `mcp/src/search.ts:477` and the -// viewer's ordering in `components/SearchResults.tsx` are the same intent -// written twice. +// Two orderings exist in this corpus and both were written inline at their one +// call site: // -// export type RankMode = "relevance" | "date" | "duration" | …; -// export function rankHits<T extends RankableHit>(hits: T[], mode: RankMode): T[]; +// newest-first over uploadDate the search result list +// (`components/SearchSessionContext.tsx`), which sorts videos and social +// posts into ONE list so a unified search reads as one result set rather +// than videos-then-posts. +// oldest-first over createdAt a social thread +// (`mcp/src/search.ts` getThread), tie-broken by id so a thread whose posts +// share a timestamp still has a stable order. // -// Fetchers are INJECTED into this module's callers ({ reader, policy, -// onProgress }) so react-query stays in `components/`; reuse -// `lib/concurrency.ts`'s `mapConcurrent` rather than a second limiter. +// They are comparators, not policies: a caller picks one, and the point of +// naming them is that the next consumer (an MCP tool that wants the viewer's +// order, say) reaches for the same function instead of re-deriving a +// three-branch ternary with its own idea of which way "newest" points. +// +// Deliberately NOT here: relevance ranking. Nothing in this repo ranks by +// score today — the scanner emits in corpus scan order and the viewer emits in +// date order — and inventing one under cover of a refactor would change what +// every caller returns. + +export type UploadDated = { uploadDate: string }; +export type Threaded = { createdAt: string; id: string }; + +// Newest upload first. Dates are "YYYYMMDD", so a lexicographic compare is a +// chronological one. Equal dates keep their input order (Array#sort is stable), +// which is what preserves the summaries index's own ordering within a day. +export function byUploadDateDesc(a: UploadDated, b: UploadDated): number { + return a.uploadDate === b.uploadDate ? 0 : a.uploadDate < b.uploadDate ? 1 : -1; +} + +// Sorts IN PLACE and returns the same array — matching the call site this +// replaced, where the array is freshly built and nobody else holds it. +export function rankByUploadDateDesc<T extends UploadDated>(items: T[]): T[] { + return items.sort(byUploadDateDesc); +} + +// Oldest post first within a thread; ties broken by id so the order is total. +export function byCreatedAtThenId(a: Threaded, b: Threaded): number { + return a.createdAt === b.createdAt + ? a.id.localeCompare(b.id) + : a.createdAt.localeCompare(b.createdAt); +} -export {}; +export function rankThread<T extends Threaded>(items: T[]): T[] { + return items.sort(byCreatedAtThenId); +} diff --git a/common/lib/search/window.test.ts b/common/lib/search/window.test.ts @@ -0,0 +1,130 @@ +import { test } from "node:test"; +import assert from "node:assert/strict"; +import { + clock, + truncate, + windowedTranscript, + findHitsInCues, + findHitsInText, + findFirstMatchInRange, +} from "./window"; +import type { Cue } from "../vtt"; +import { MCP_POLICY, VIEWER_POLICY } from "./policy"; + +const cue = (start: number, text: string): Cue => ({ start, end: start + 3, text }); + +test("clock renders 0 as 0:00 rather than the empty string", () => { + assert.equal(clock(0), "0:00"); + assert.equal(clock(-5), "0:00"); + assert.equal(clock(65), "1:05"); + assert.equal(clock(3725), "1:02:05"); +}); + +test("truncate collapses whitespace and clips with the ellipsis inside the budget", () => { + assert.equal(truncate(" a b\n c ", MCP_POLICY.snippetChars), "a b c"); + assert.equal(truncate("abcdef", 4), "abc…"); + assert.equal(truncate("abcd", 4), "abcd"); +}); + +// The regression this guards: `max` used to default to MCP_POLICY.snippetChars, +// so a viewer path that forgot to pass its width clipped at 240 with every test +// green. The default is gone (tsc now requires the argument); this pins that the +// two policies actually produce different widths, which is the observable the +// silent default hid. +test("the width is the CALLER's, and the two policies differ", () => { + const long = "z".repeat(600); + assert.equal(truncate(long, MCP_POLICY.snippetChars).length, 240); + assert.equal(truncate(long, VIEWER_POLICY.snippetChars), long); +}); + +test("windowedTranscript merges overlapping windows and stamps each line", () => { + const cues = [ + cue(0, "one"), + cue(10, "needle here"), + cue(20, "three"), + cue(30, "needle again"), + cue(40, "five"), + ]; + const r = windowedTranscript(cues, (t) => t.includes("needle"), { + before: 15, + after: 15, + maxLines: MCP_POLICY.windowLineCap, + }); + assert.equal(r.matchCount, 2); + // Both windows overlap, so every cue appears exactly once. + assert.equal(r.lines.length, 5); + assert.equal(r.lines[0], "[0:00] one"); + assert.equal(r.lines[1], "[0:10] needle here"); +}); + +test("windowedTranscript honours maxLines and the timestamps/stamp options", () => { + const cues = Array.from({ length: 40 }, (_, i) => cue(i, `line ${i} needle`)); + const capped = windowedTranscript(cues, () => true, { maxLines: 5 }); + assert.equal(capped.lines.length, 5); + + const bare = windowedTranscript([cue(7, "hit")], () => true, { + timestamps: false, + maxLines: MCP_POLICY.windowLineCap, + }); + assert.deepEqual(bare.lines, ["hit"]); + + const linked = windowedTranscript([cue(7, "hit")], () => true, { + stamp: (c, s) => `${c}|${s}`, + maxLines: MCP_POLICY.windowLineCap, + }); + assert.deepEqual(linked.lines, ["[0:07|7] hit"]); +}); + +test("findHitsInCues emits the matched cue, widening only across a cue boundary", () => { + const cues = [ + { start: 0, text: "the quick brown" }, + { start: 3, text: "fox jumps over" }, + { start: 6, text: "the lazy dog" }, + ]; + // Wholly inside one cue → that cue's text alone. + const inside = findHitsInCues(cues, "jumps", false, null, 10); + assert.deepEqual(inside, [{ start: 3, text: "fox jumps over" }]); + + // Straddling the first/second cue → emitted ONCE, on the cue the match + // STARTS in, with the text widened to that cue's window. + const across = findHitsInCues(cues, "brown fox", false, null, 10); + assert.equal(across.length, 1); + assert.equal(across[0].start, 0); + assert.equal(across[0].text, "the quick brown fox jumps over"); +}); + +test("findHitsInCues stops at the limit", () => { + const cues = Array.from({ length: 10 }, (_, i) => ({ start: i, text: "needle" })); + assert.equal(findHitsInCues(cues, "needle", false, null, 3).length, 3); +}); + +test("findHitsInText emits one padded snippet with ellipses at the clipped ends", () => { + const text = `${"a".repeat(200)} needle ${"b".repeat(200)}`; + const hits = findHitsInText(text, "needle", false, null); + assert.equal(hits.length, 1); + assert.equal(hits[0].start, 0); + assert.ok(hits[0].text.startsWith("…")); + assert.ok(hits[0].text.endsWith("…")); + assert.ok(hits[0].text.includes("needle")); + assert.deepEqual(findHitsInText("", "needle", false, null), []); + assert.deepEqual(findHitsInText("nothing", "needle", false, null), []); +}); + +test("findFirstMatchInRange only accepts a match STARTING in the range", () => { + const hay = "prev cur next"; + // "cur" starts at 5, inside [5, 8). + assert.deepEqual(findFirstMatchInRange(hay, 5, 8, "cur", false, null), { + idx: 5, + length: 3, + }); + // "next" starts at 9, outside the current cue's span. + assert.equal(findFirstMatchInRange(hay, 5, 8, "next", false, null), null); + // A zero-width regex advances rather than spinning forever: it settles on + // the first offset inside the range. + assert.deepEqual(findFirstMatchInRange(hay, 5, 8, "", true, /x*/), { + idx: 5, + length: 0, + }); + // Regex mode with no compiled regex is a miss, never a throw. + assert.equal(findFirstMatchInRange(hay, 0, 13, "cur", true, null), null); +}); diff --git a/common/lib/search/window.ts b/common/lib/search/window.ts @@ -1,18 +1,208 @@ -// ONE SEARCH PIPELINE — placeholder. Filled by one-core phase 2 slice S3 -// (`plans/one-core-phase-2.md` §S3); created here by S1 so no two parallel -// slices race to create the same file. +// ONE SEARCH PIPELINE — the excerpt layer. // -// Intended contents: the windowing + snippet layer, over the existing -// `lib/transcriptWindow.ts` primitives (`windowCues`, `cuesToSnippets`, -// `mergeSnippets`), so the viewer's excerpt and the MCP's excerpt are the same -// excerpt. +// Every place that turns cues into something a person or an agent reads lives +// here, over the primitives in `lib/transcriptWindow.ts` (`windowCues`, +// `cuesToSnippets`, `mergeSnippets`). Before this file the MCP server owned +// one excerpt shape and `components/searchPipeline.ts` owned another, in two +// directories that never referenced each other. // -// export function windowedTranscript( -// cues: Cue[], -// match: Matcher, -// opts: { before: number; after: number; lineCap: number }, -// ): WindowSnippet[]; +// THEY ARE STILL TWO SHAPES, deliberately, and the module says so out loud: // -// export function truncate(text: string, max: number): string; +// windowedTranscript() a ±seconds window around every matched cue, merged +// and deduped — the batch read a sweep prompt drives. +// Wide context, timestamps, bounded by a line cap. +// findHitsInCues() one line per matched cue, widened to the previous +// and next cue ONLY when the match straddles a cue +// boundary — the result-card row a reader scans. +// +// A sweep needs the paragraph around a claim; a result card needs the line +// that matched and nothing else. Collapsing them would change both outputs, +// and the second one is pinned by the export e2e's rendered hit text. What the +// move buys is that there is now exactly ONE implementation of each, one +// `truncate`, and one `clock` — and the next excerpt shape has an obvious +// place to land next to its siblings rather than in whichever app needed it. + +import { formatDuration } from "../format"; +import type { Cue } from "../vtt"; +import { + windowCues, + cuesToSnippets, + mergeSnippets, + type WindowSnippet, +} from "../transcriptWindow"; + +// A zero-second cue reads as "0:00", not as the empty string formatDuration +// returns for a falsy input. +export function clock(seconds: number): string { + const s = Math.max(0, Math.floor(seconds)); + return s === 0 ? "0:00" : formatDuration(s); +} + +// Collapse whitespace and clip to `max` characters, ellipsis included in the +// budget. +// +// `max` is REQUIRED, and that is the point. It used to default to +// `MCP_POLICY.snippetChars`, which meant a viewer adopter that simply forgot to +// pass its policy would silently clip every excerpt at the MCP's 240 characters +// with every test still green — the failure mode where nothing is broken, the +// output is just quietly wrong. A shared module does not get to hold one +// caller's budget as its default. Pass `policy.snippetChars`, or a literal when +// the width is genuinely not the policy's to set (a post's 120-character +// stand-in title). +export function truncate(text: string, max: number): string { + const t = text.trim().replace(/\s+/g, " "); + return t.length > max ? t.slice(0, max - 1) + "…" : t; +} + +export type Matcher = (text: string) => boolean; + +// Render one record's transcript around the cues that match `matcher`: for each +// matched cue take a bounded window of surrounding cues (±`before`/`after` +// seconds), merge overlapping windows deduped by timestamp (capped by +// `maxLines`), and return timestamped excerpt lines plus the total match count. +export function windowedTranscript( + cues: readonly Cue[], + matcher: Matcher, + opts: { + before?: number; + after?: number; + maxCues?: number; + timestamps?: boolean; + // Optional formatter for the bracketed stamp's CONTENTS (e.g. an inline + // Markdown link "[m:ss](url)" or the compact base form "m:ss|156"), given + // the line's clock + start seconds. Bare clock when omitted. + stamp?: (clock: string, seconds: number) => string; + // Cap on merged excerpt lines emitted for this video (`policy.windowLineCap`). + // REQUIRED, for the same reason `truncate`'s `max` is. `maxCues` above + // stays the per-window bound. + maxLines: number; + }, +): { lines: string[]; matchCount: number } { + const list = cues as Cue[]; + const timestamps = opts.timestamps !== false; + let merged: WindowSnippet[] = []; + let matchCount = 0; + for (const cue of list) { + if (!matcher(cue.text)) continue; + matchCount++; + const win = cuesToSnippets( + windowCues(list, cue.start, { + before: opts.before, + after: opts.after, + maxCues: opts.maxCues, + }), + ); + merged = mergeSnippets(merged, win, opts.maxLines); + } + const lines = merged.map((s) => { + if (!timestamps) return s.text; + const stamp = opts.stamp ? opts.stamp(s.clock, s.seconds) : s.clock; + return `[${stamp}] ${s.text}`; + }); + return { lines, matchCount }; +} + +// ─── The result-row excerpt ─── +// +// One hit per matched cue. The match is looked for in a three-cue window +// (previous + current + next) but only ACCEPTED when it starts inside the +// current cue, so a phrase split across a caption break is found exactly once +// — on the cue it starts in — and the emitted text widens to the window only +// in that case. Everything else emits the matched cue verbatim. + +export type Hit = { start: number; text: string }; + +export type SubsHit = { track: string; start: number; text: string }; + +export function findHitsInCues( + cues: readonly { start: number; text: string }[], + query: string, + useRegex: boolean, + regex: RegExp | null, + limit: number, +): Hit[] { + const hits: Hit[] = []; + for (let i = 0; i < cues.length && hits.length < limit; i++) { + const cur = cues[i]; + const prevText = i > 0 ? cues[i - 1].text : ""; + const nextText = i < cues.length - 1 ? cues[i + 1].text : ""; + const sep1 = prevText ? " " : ""; + const sep2 = nextText ? " " : ""; + const windowText = prevText + sep1 + cur.text + sep2 + nextText; + const curStart = prevText.length + sep1.length; + const curEnd = curStart + cur.text.length; + const m = findFirstMatchInRange( + windowText, + curStart, + curEnd, + query, + useRegex, + regex, + ); + if (!m) continue; + const crosses = m.idx < curStart || m.idx + m.length > curEnd; + hits.push({ + start: Math.round(cur.start), + text: crosses ? windowText : cur.text, + }); + } + return hits; +} + +// Single-document text match (description / tags / post body): emit at most one +// hit with a ±PAD window around the first match, so the result row shows +// context without shipping a whole document into a card. +export function findHitsInText( + text: string, + query: string, + useRegex: boolean, + regex: RegExp | null, +): Hit[] { + if (!text) return []; + const m = findFirstMatchInRange(text, 0, text.length, query, useRegex, regex); + if (!m) return []; + const PAD = 80; + const from = Math.max(0, m.idx - PAD); + const to = Math.min(text.length, m.idx + m.length + PAD); + const snippet = + (from > 0 ? "…" : "") + + text.slice(from, to).replace(/\s+/g, " ").trim() + + (to < text.length ? "…" : ""); + return [{ start: 0, text: snippet }]; +} -export {}; +export function findFirstMatchInRange( + haystack: string, + rangeStart: number, + rangeEnd: number, + query: string, + useRegex: boolean, + regex: RegExp | null, +): { idx: number; length: number } | null { + if (useRegex) { + if (!regex) return null; + const flags = regex.flags.includes("g") ? regex.flags : regex.flags + "g"; + const re = new RegExp(regex.source, flags); + let m: RegExpExecArray | null; + while ((m = re.exec(haystack)) !== null) { + if (m.index >= rangeStart && m.index < rangeEnd) { + return { idx: m.index, length: m[0].length }; + } + if (m.index >= rangeEnd) return null; + if (m[0].length === 0) re.lastIndex++; + } + return null; + } + if (!query) return null; + const lower = haystack.toLowerCase(); + const ql = query.toLowerCase(); + let from = 0; + while (from <= haystack.length) { + const idx = lower.indexOf(ql, from); + if (idx === -1) return null; + if (idx >= rangeStart && idx < rangeEnd) return { idx, length: ql.length }; + if (idx >= rangeEnd) return null; + from = idx + 1; + } + return null; +} diff --git a/common/lib/searchEval.ts b/common/lib/searchEval.ts @@ -1,11 +1,24 @@ -// Composite-search tree orchestrator. +// Composite-search tree orchestrator — the SLUG-SET half of one search +// pipeline. // // Walks a QueryNode tree (from `searchQuery.ts`), kicks off per-leaf -// pipelines (from `searchPipeline.ts:runLeafPipeline`), runs metadata-scope -// leaves synchronously against summaries, memoizes per-leaf results in -// `searchLayerCache.ts`, and emits a streaming TreeProgress that the UI +// 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 @@ -24,20 +37,37 @@ import { type LeafNode, type QueryNode, } from "./searchQuery"; -import { - cacheKey, - getCached, - getCachedSync, - putCached, - type CachedResult, -} from "../components/searchLayerCache"; -import { - runLeafPipeline, - type LayerHit, - type LeafController, -} from "../components/searchPipeline"; +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<string>; + hits: ReadonlyMap<string, LayerHit[]>; +}; + +// 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<CachedResult | null>; + 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; @@ -73,6 +103,7 @@ export type TreeController = { type MetadataLeafScope = "transcripts" | "chat" | "metadata"; type EvalCtx = { + runtime: SearchRuntime; initialHitLimit: number; concurrency: number; flushIntervalMs: number; @@ -122,6 +153,8 @@ function buildMetadataIndex(summaries: DisplaySummary[]): MetadataIndex { export function runQueryTree(opts: { root: GroupNode; + // The fetch-bound leaf runner + layer memo. See SearchRuntime above. + runtime: SearchRuntime; globalScope: string[]; summaries: DisplaySummary[]; chatScopeSlugs?: ReadonlySet<string> | null; @@ -133,6 +166,7 @@ export function runQueryTree(opts: { }): TreeController { const { root, + runtime, globalScope, summaries, chatScopeSlugs = null, @@ -149,6 +183,7 @@ export function runQueryTree(opts: { let allDone = false; const ctx: EvalCtx = { + runtime, initialHitLimit, concurrency, flushIntervalMs, @@ -311,12 +346,12 @@ async function runLeaf( // Fully-network leaves (transcripts / chat). Try the layer cache first. const scopeArr = Array.from(effectiveScope); const scopeHash = hashSlugs(scopeArr); - const key = cacheKey(canonicalHash(leaf), scopeHash); - const cachedSync = getCachedSync(key); + const key = ctx.runtime.cache.key(canonicalHash(leaf), scopeHash); + const cachedSync = ctx.runtime.cache.getSync(key); if (cachedSync) { return applyCached(leaf, parentScope, cachedSync, ctx, /*cached*/ true); } - const cached = await getCached(key); + const cached = await ctx.runtime.cache.get(key); if (ctx.cancelled) return new Set(); if (cached) { return applyCached(leaf, parentScope, cached, ctx, /*cached*/ true); @@ -338,7 +373,7 @@ async function runLeaf( }); let lastCapped = false; - const controller = runLeafPipeline({ + const controller = ctx.runtime.runLeaf({ leaf, scopeSlugs: scopeArr, initialHitLimit, @@ -371,7 +406,7 @@ async function runLeaf( // 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) { - putCached(key, { + ctx.runtime.cache.put(key, { slugs: new Set(result.slugs), hits: new Map(result.hits), }); diff --git a/common/lib/searchQuery.ts b/common/lib/searchQuery.ts @@ -9,7 +9,14 @@ // and IndexedDB-backed layer cache keys. Keep field names short — they end // up in URLs. -import type { SearchMode } from "../components/urlState"; +// Which corpus the legacy single-input search targets. "posts" is the social +// corpus (lib/posts.ts); the composite query tree can mix all of them freely, +// this only decides what a bare `?q=` URL means. +// +// Defined here rather than in `components/urlState.ts` (which re-exports it) +// because `buildSearchRoot` below turns one into a query tree — the model owns +// the vocabulary it interprets. +export type SearchMode = "transcripts" | "subs" | "posts"; export type LayerScope = | "transcripts" diff --git a/export/app/lib/askRetrieval.ts b/export/app/lib/askRetrieval.ts @@ -12,7 +12,10 @@ import { tokenizeQuery, type SearchAlias, } from "yt-dlp-transcript-common/lib/searchAliases"; -import type { LayerHit } from "yt-dlp-transcript-common/components/searchPipeline"; +import { + searchRuntime, + type LayerHit, +} from "yt-dlp-transcript-common/components/searchPipeline"; import { peekPost } from "yt-dlp-transcript-common/components/postsCache"; import type { DisplaySummary } from "yt-dlp-transcript-common/lib/transcripts"; import { formatDuration } from "yt-dlp-transcript-common/lib/format"; @@ -348,6 +351,7 @@ export function retrieve( const postSlugs = opts.postScopeSlugs ?? null; const controller = runQueryTree({ root, + runtime: searchRuntime, globalScope: [ ...summaries.map((s) => s.slug), ...(postSlugs ? Array.from(postSlugs) : []), diff --git a/mcp/src/search.ts b/mcp/src/search.ts @@ -1,43 +1,61 @@ -import { formatDuration } from "yt-dlp-transcript-common/lib/format"; +// ORCHESTRATION over `yt-dlp-transcript-common/lib/search/`. +// +// Everything below plans and drives reads against a ShardSource; the matching, +// excerpting, tree algebra, filter predicate, ordering and mirror collapsing +// all live in the shared pipeline, once, alongside the viewer's half of it. +// The caps travel as `MCP_POLICY` rather than being re-derived here, which is +// what makes the bench's structural read/byte counters unchanged BY +// CONSTRUCTION: this file cannot quietly pick a different ceiling than the one +// the policy names. import type { TranscriptDetail } from "yt-dlp-transcript-common/lib/transcripts"; -import type { Cue } from "yt-dlp-transcript-common/lib/vtt"; import type { Post } from "yt-dlp-transcript-common/lib/posts"; +import type { Cue } from "yt-dlp-transcript-common/lib/vtt"; import type { Platform } from "yt-dlp-transcript-common/lib/platform"; import { - forEachLeaf, - isLeaf, - isLeafActive, - isNodeActive, type GroupNode, type LayerScope, - type QueryNode, } from "yt-dlp-transcript-common/lib/searchQuery"; -import { - matchAliases, - type SearchAlias, -} from "yt-dlp-transcript-common/lib/searchAliases"; -import { - windowCues, - cuesToSnippets, - mergeSnippets, - type WindowSnippet, -} from "yt-dlp-transcript-common/lib/transcriptWindow"; +import type { SearchAlias } from "yt-dlp-transcript-common/lib/searchAliases"; import { resolveChannelGroupId, type ChannelGroup, } from "yt-dlp-transcript-common/lib/channelGroups"; -import { - VIDEO_STATES, - type VideoState, -} from "yt-dlp-transcript-common/lib/availability"; import { mapConcurrent } from "yt-dlp-transcript-common/lib/concurrency"; import { + MCP_POLICY, + type SearchPolicy, +} from "yt-dlp-transcript-common/lib/search/policy"; +import { + clock, + truncate, + windowedTranscript, +} from "yt-dlp-transcript-common/lib/search/window"; +import { + buildMatcher, + buildLeafMatchers, + evalNode, + filterIsSelective, + needsAvailability, + passesFilters, + type Matcher, + type RecordCtx, + type ScopedSnippet, + type SearchFilters, +} from "yt-dlp-transcript-common/lib/search/evalTree"; +import { collapseDuplicates as collapseHits } from "yt-dlp-transcript-common/lib/search/collapse"; +import { rankThread } from "yt-dlp-transcript-common/lib/search/rank"; +import { DEFAULT_PAGE_CONCURRENCY, type ChannelRef, type ShardSource, - type VideoAvailability, } from "./source"; +// Re-exported so the thirteen tools and the two test files keep importing +// every name from `./search`, exactly as before. The definitions moved; the +// module's surface did not. +export { buildMatcher, filterIsSelective, MCP_POLICY }; +export type { Matcher, ScopedSnippet, SearchFilters, SearchPolicy }; + export type Snippet = { clock: string; seconds: number; @@ -200,7 +218,7 @@ export type SearchHit = { export type SearchResult = { hits: SearchHit[]; - // Total matched videos found in this scan (up to HARD_VIDEO_CAP). `hits` is + // Total matched videos found in this scan (up to policy.hardVideoCap). `hits` is // the [offset, offset+limit) slice of that set, so total is stable across // paged calls and lets a caller plan a full sweep. total: number; @@ -222,7 +240,7 @@ export type SearchResult = { // either way — distinct from "checked, found none". duplicates: { collapsed: number; clusters: number; available: boolean }; // Coverage is partial — the page cap (MAX_PAGES) or the video cap - // (HARD_VIDEO_CAP) was reached before the corpus was fully scanned. + // (policy.hardVideoCap) was reached before the corpus was fully scanned. truncated: boolean; // What was actually covered, in enough detail to act on. A cap truncates in // CHANNEL ITERATION ORDER, not at random, so "partial" on its own is @@ -254,24 +272,6 @@ export type SearchResult = { }; }; -// A hard ceiling on shard pages fetched per query so a rare term over a large -// (or hub-wide) corpus can't run away. Reaching it sets `truncated`. -const MAX_PAGES = 400; - -// A ceiling on the number of matched videos we collect before we stop counting, -// so `total` stays bounded and stable even for a very common term. Reaching it -// also sets `truncated` (the true total is higher than reported). -const HARD_VIDEO_CAP = 2000; - -// Cap on windowed excerpt lines emitted per video by getWindowedTranscript, so a -// video with hundreds of matches can't blow the batch's token budget. -const WINDOW_LINE_CAP = 200; - -function clock(seconds: number): string { - const s = Math.max(0, Math.floor(seconds)); - return s === 0 ? "0:00" : formatDuration(s); -} - // ─── Windowed page reading ─── // // Pages were read one `await` at a time, which on a local corpus leaves the @@ -335,18 +335,6 @@ export type ScanPlan = { unknownVideos: number; }; -// True when a filter set could actually exclude something. An all-permissive -// filter (every state kept, both media types, both audiences, no dates) is the -// same query as no filter at all, and must NOT trigger an index read — that is -// the guard against making an unfiltered query slower by planning it. -export function filterIsSelective(f: SearchFilters | null | undefined): boolean { - if (!f) return false; - if (!VIDEO_STATES.every((s) => f.states.has(s))) return true; - if (!f.videos || !f.livestreams) return true; - if (!f.allAges || !f.restricted) return true; - return Boolean(f.dateFrom || f.dateTo); -} - // Build the page plan for a scan. Falls back to "every page of every channel" // whenever pruning is impossible or pointless: no selective filter, or a source // with no summaries index (both in-memory test stubs, and any site that ships @@ -408,72 +396,11 @@ export async function buildScanPlan( }; } -export type Matcher = (text: string) => boolean; - -// Build the combined, alias-aware matcher for a query, mirroring the browser's -// buildSearchRoot OR-of-leaves semantics: a text matches if the plain query -// substring matches OR any fired alias's suggestion regex matches. An explicit -// `regex` query is taken verbatim with NO alias expansion (the caller is -// crafting their own pattern). Returns the fired aliases so the tool can report -// which curated expansions it applied. -export function buildMatcher(opts: { - query: string; - regex?: boolean; - useAliases?: boolean; - aliases?: SearchAlias[]; -}): { match: Matcher; firedAliases: SearchAlias[] } { - if (opts.regex) { - const re = new RegExp(opts.query, "i"); - return { match: (t) => re.test(t), firedAliases: [] }; - } - const needle = opts.query.toLowerCase(); - const plain: Matcher = (t) => t.toLowerCase().includes(needle); - - const useAliases = opts.useAliases !== false; - const fired = - useAliases && opts.aliases && opts.aliases.length > 0 - ? matchAliases(opts.query, "transcripts", opts.aliases) - : []; - if (fired.length === 0) return { match: plain, firedAliases: [] }; - - const aliasMatchers: Matcher[] = fired.map((a) => { - if (a.useRegex) { - try { - const re = new RegExp(a.suggestion, "i"); - return (t: string) => re.test(t); - } catch { - // malformed suggestion regex — fall back to substring on the literal - } - } - const n = a.suggestion.toLowerCase(); - return (t: string) => t.toLowerCase().includes(n); - }); - - const match: Matcher = (t) => plain(t) || aliasMatchers.some((m) => m(t)); - return { match, firedAliases: fired }; -} - -function truncate(text: string, max = 240): string { - const t = text.trim().replace(/\s+/g, " "); - return t.length > max ? t.slice(0, max - 1) + "…" : t; -} - -// Collapse cross-platform mirrors in a result list, IN PLACE, keeping one row -// per recording. Returns what it did so the caller can report it. -// -// Two rules make this safe to have on by default: -// -// 1. The kept row is the cluster's canonical member WHEN that member is -// itself among the matches — otherwise it is simply the first match. A -// mirror is frequently the only surviving copy of a deleted upload, and -// preferring an absent canonical would delete exactly the evidence a -// "what did the removed videos say" question is asking for. -// 2. The collapsed copies are NAMED on the row they folded into. Nothing -// vanishes; the count stops double-counting. Pass collapse_duplicates:false -// to see every upload as its own row. -// -// Timestamps are never mapped between copies here — that requires the per-pair -// `aligned` gate, and this function does not move a single second of anything. +// Fetch the duplicate index and apply the collapse rules (lib/search/collapse) +// to a result list, IN PLACE. Only the read and the "this corpus ships no +// duplicates.json" distinction live here: `available: false` means no claim +// about mirrors can be made either way, which is different from "checked, +// found none". async function collapseDuplicates( source: ShardSource, all: SearchHit[], @@ -490,51 +417,20 @@ async function collapseDuplicates( } if (index.size === 0) return { collapsed: 0, clusters: 0, available: true }; - const repIndexOf = new Map<string, number>(); // clusterId -> index in `kept` - const kept: SearchHit[] = []; - let collapsed = 0; - for (const hit of all) { - const membership = index.get(hit.slug); - if (!membership) { - kept.push(hit); - continue; - } - const at = repIndexOf.get(membership.clusterId); - if (at === undefined) { - repIndexOf.set(membership.clusterId, kept.length); - kept.push(hit); - continue; - } - collapsed++; - const rep = kept[at]; - const repIsCanonical = index.get(rep.slug)?.isCanonical === true; - const fold = (into: SearchHit, gone: SearchHit): SearchHit => ({ - ...into, - mirrors: [ - ...(into.mirrors ?? []), - ...(gone.mirrors ?? []), - { videoId: gone.videoId, channelName: gone.channelName, slug: gone.slug }, - ], - }); - // Promote the canonical member to the representative if it turns up later; - // otherwise fold this copy into the incumbent. - kept[at] = repIsCanonical || !membership.isCanonical - ? fold(rep, hit) - : fold(hit, rep); - } + const { collapsed, clusters, kept } = collapseHits(all, index); if (collapsed > 0) { all.length = 0; all.push(...kept); } - return { collapsed, clusters: repIndexOf.size, available: true }; + return { collapsed, clusters, available: true }; } // Scan a source's transcript shards for `query` (alias-aware by default), -// collecting ALL matched videos up to HARD_VIDEO_CAP so counting is stable, then +// collecting ALL matched videos up to policy.hardVideoCap so counting is stable, then // returning the [offset, offset+limit) slice with a `total`/`hasMore`. Plain // substring by default, or a caller-supplied regex; either can be paged. Set // `includeSnippets:false` for a cheap worklist (id/title/channel/date/matches, -// no cue text). Stops early at MAX_PAGES (→ truncated) and HARD_VIDEO_CAP. +// no cue text). Stops early at MAX_PAGES (→ truncated) and policy.hardVideoCap. export async function searchTranscripts( source: ShardSource, opts: { @@ -571,11 +467,17 @@ export async function searchTranscripts( // collapsed copies are named on the row they fold into, never dropped // silently. collapseDuplicates?: boolean; + // The caps this scan runs under. Defaults to MCP_POLICY — the same four + // numbers this scanner used when they were private constants — so an + // existing caller reads exactly the same pages and parses exactly the same + // bytes. A caller with a different budget passes its own. + policy?: SearchPolicy; }, ): Promise<SearchResult> { + const policy = opts.policy ?? MCP_POLICY; const limit = opts.limit ?? 20; const offset = Math.max(0, opts.offset ?? 0); - const maxPages = opts.maxPages ?? MAX_PAGES; + const maxPages = opts.maxPages ?? policy.maxPages; const includeSnippets = opts.includeSnippets !== false; const snippetsPerVideo = opts.snippetsPerVideo ?? 4; @@ -699,7 +601,7 @@ export async function searchTranscripts( push({ clock: clock(cue.start), seconds: cue.start, - text: truncate(cue.text), + text: truncate(cue.text, policy.snippetChars), }); } } @@ -722,7 +624,7 @@ export async function searchTranscripts( push({ clock: clock(0), seconds: 0, - text: truncate(rec.description), + text: truncate(rec.description, policy.snippetChars), scope: "description", }); } @@ -730,7 +632,12 @@ export async function searchTranscripts( const tags = (rec.tags ?? []).join(", "); if (tags && match(tags)) { otherHit = true; - push({ clock: clock(0), seconds: 0, text: truncate(tags), scope: "tags" }); + push({ + clock: clock(0), + seconds: 0, + text: truncate(tags, policy.snippetChars), + scope: "tags", + }); } } if (wantChat && chatCuesFor) { @@ -740,7 +647,7 @@ export async function searchTranscripts( push({ clock: clock(cue.start), seconds: cue.start, - text: truncate(cue.text), + text: truncate(cue.text, policy.snippetChars), scope: "chat", }); } @@ -761,7 +668,7 @@ export async function searchTranscripts( matches: matches || 1, snippets, }); - if (all.length >= HARD_VIDEO_CAP) { + if (all.length >= policy.hardVideoCap) { truncated = true; channelStopped = { channel: ch.name, @@ -828,7 +735,7 @@ export async function searchTranscripts( ? [{ clock: "", seconds: 0, text: truncate(post.text, 480) }] : [], }); - if (all.length >= HARD_VIDEO_CAP) { + if (all.length >= policy.hardVideoCap) { truncated = true; break postsOuter; } @@ -961,11 +868,7 @@ export async function getThread( if ((p.threadId || p.id) === threadId) thread.push(p); } } - thread.sort((a, b) => - a.createdAt === b.createdAt - ? a.id.localeCompare(b.id) - : a.createdAt.localeCompare(b.createdAt), - ); + rankThread(thread); return thread.length > 0 ? thread : [post]; } @@ -982,33 +885,15 @@ export function getWindowedTranscript( // the line's clock + start seconds. Bare clock when omitted. stamp?: (clock: string, seconds: number) => string; // Cap on merged excerpt lines emitted for this video (default - // WINDOW_LINE_CAP). The earliest lines are kept; `maxCues` above stays the - // per-window bound. + // MCP_POLICY.windowLineCap). The earliest lines are kept; `maxCues` above + // stays the per-window bound. maxLines?: number; } = {}, ): { lines: string[]; matchCount: number } { - const cues = record.cues ?? []; - const timestamps = opts.timestamps !== false; - let merged: WindowSnippet[] = []; - let matchCount = 0; - for (const cue of cues) { - if (!matcher(cue.text)) continue; - matchCount++; - const win = cuesToSnippets( - windowCues(cues, cue.start, { - before: opts.before, - after: opts.after, - maxCues: opts.maxCues, - }), - ); - merged = mergeSnippets(merged, win, opts.maxLines ?? WINDOW_LINE_CAP); - } - const lines = merged.map((s) => { - if (!timestamps) return s.text; - const stamp = opts.stamp ? opts.stamp(s.clock, s.seconds) : s.clock; - return `[${stamp}] ${s.text}`; + return windowedTranscript(record.cues ?? [], matcher, { + ...opts, + maxLines: opts.maxLines ?? MCP_POLICY.windowLineCap, }); - return { lines, matchCount }; } // Locate a single video across the source's channels via each channel's @@ -1067,26 +952,10 @@ export async function findVideo( // the simpler `searchTranscripts` scanner above (a trivial single-leaf query), // so today's callers/tests are unaffected. // -// It mirrors the browser's tree algebra (searchEval.ts) and filter predicate -// (SearchSessionContext.tsx `passesFilter`) — evaluated per record to a boolean -// instead of over slug sets — and the per-scope text extraction of -// searchPipeline.ts. - -// The positive share-filter selection (parseShareV1), as a predicate input. -export type SearchFilters = { - // ft — video / livestream types kept. - videos: boolean; - livestreams: boolean; - // fa — all-ages / age-restricted kept. - allAges: boolean; - restricted: boolean; - // fav — the VideoState values kept (see common/lib/availability). Absent - // from the set means filtered out. - states: ReadonlySet<VideoState>; - // fdf / fdt — inclusive upload-date bounds, "YYYYMMDD". - dateFrom?: string; - dateTo?: string; -}; +// The tree algebra, the leaf evaluation and the filter predicate are +// `lib/search/evalTree.ts` — the same module the viewer's `lib/searchEval.ts` +// points at for the meaning of AND / OR / negate. What stays here is the walk: +// which pages to read, in what order, and when to stop. export type SearchSpec = { // The root of the composite query (a `qt=` tree, or a synthesized single leaf @@ -1102,17 +971,6 @@ export type SearchSpec = { useAliases?: boolean; }; -// A snippet tagged with the scope it came from, carrying the seconds needed to -// build a moment link. `seconds` is 0 for the non-timed scopes (metadata / -// description / tags). -export type ScopedSnippet = { - scope: LayerScope; - track?: string; - clock: string; - seconds: number; - text: string; -}; - export type SpecHit = { videoId: string; slug: string; @@ -1139,226 +997,6 @@ export type SpecResult = { truncated: boolean; }; -type LeafMatcher = { scope: LayerScope; test: Matcher }; - -// Compile a matcher per active leaf. Transcripts leaves are alias-aware (unless -// they're regex); every other scope matches plain-substring / regex only. Fired -// aliases are unioned for the caller to report. -function buildLeafMatchers( - root: QueryNode, - aliases: SearchAlias[], - useAliases: boolean, -): { matchers: Map<string, LeafMatcher>; fired: SearchAlias[] } { - const matchers = new Map<string, LeafMatcher>(); - const fired: SearchAlias[] = []; - forEachLeaf(root, (leaf) => { - if (!isLeafActive(leaf)) return; - const aliasAware = - leaf.scope === "transcripts" && !leaf.useRegex && useAliases; - const built = buildMatcher({ - query: leaf.query, - regex: leaf.useRegex, - useAliases: aliasAware, - aliases: aliasAware ? aliases : [], - }); - matchers.set(leaf.id, { scope: leaf.scope, test: built.match }); - for (const a of built.firedAliases) { - if (!fired.some((x) => x.id === a.id)) fired.push(a); - } - }); - return { matchers, fired }; -} - -// Per-record context the tree evaluates against. -type RecordCtx = { - title: string; - channel: string; - description: string; - tags: string; - cues: Cue[]; - chatCues: Cue[]; - // The post body, when this record IS a post rather than a video. Empty for a - // video record, so a posts-scope leaf never matches one. - postText?: string; - snippetsPerVideo: number; - includeSnippets: boolean; -}; - -type LeafOutcome = { matched: boolean; count: number; hits: ScopedSnippet[] }; - -function evalLeaf(leaf: QueryNode, m: LeafMatcher, ctx: RecordCtx): LeafOutcome { - if (!isLeaf(leaf)) return { matched: false, count: 0, hits: [] }; - const hits: ScopedSnippet[] = []; - const push = (s: ScopedSnippet): void => { - if (ctx.includeSnippets && hits.length < ctx.snippetsPerVideo) hits.push(s); - }; - let count = 0; - switch (m.scope) { - case "transcripts": - for (const cue of ctx.cues) { - if (!m.test(cue.text)) continue; - count++; - push({ - scope: "transcripts", - clock: clock(cue.start), - seconds: cue.start, - text: truncate(cue.text), - }); - } - break; - case "chat": - for (const cue of ctx.chatCues) { - if (!m.test(cue.text)) continue; - count++; - push({ - scope: "chat", - track: "live_chat", - clock: clock(cue.start), - seconds: cue.start, - text: truncate(cue.text), - }); - } - break; - case "posts": - // A post has no timeline: one hit, seconds 0 — the same convention the - // metadata / description / tags scopes already use. - if (ctx.postText && m.test(ctx.postText)) { - count++; - push({ - scope: "posts", - clock: clock(0), - seconds: 0, - text: truncate(ctx.postText), - }); - } - break; - case "metadata": { - const titleHit = m.test(ctx.title); - const channelHit = m.test(ctx.channel); - if (titleHit || channelHit) { - count++; - if (titleHit) { - push({ scope: "metadata", clock: clock(0), seconds: 0, text: truncate(ctx.title) }); - } else { - push({ scope: "metadata", clock: clock(0), seconds: 0, text: `Channel: ${ctx.channel}` }); - } - } - break; - } - case "description": - if (ctx.description && m.test(ctx.description)) { - count++; - push({ scope: "description", clock: clock(0), seconds: 0, text: truncate(ctx.description) }); - } - break; - case "tags": - if (ctx.tags && m.test(ctx.tags)) { - count++; - push({ scope: "tags", clock: clock(0), seconds: 0, text: truncate(ctx.tags) }); - } - break; - } - return { matched: count > 0, count, hits }; -} - -type NodeOutcome = { match: boolean; count: number; hits: ScopedSnippet[] }; - -// Evaluate the tree against one record — mirrors searchEval's AND/OR/negate, -// per record. A negated node contributes no hits (like the browser's `diff`). -// An inactive (empty) subtree is identity (matches, no hits). -function evalNode( - node: QueryNode, - matchers: Map<string, LeafMatcher>, - ctx: RecordCtx, -): NodeOutcome { - if (!isNodeActive(node)) return { match: true, count: 0, hits: [] }; - if (isLeaf(node)) { - const m = matchers.get(node.id); - if (!m) return { match: true, count: 0, hits: [] }; - const r = evalLeaf(node, m, ctx); - if (node.negate) return { match: !r.matched, count: 0, hits: [] }; - return { - match: r.matched, - count: node.contributeHits ? r.count : 0, - hits: node.contributeHits ? r.hits : [], - }; - } - const active = node.children.filter(isNodeActive); - if (active.length === 0) return { match: true, count: 0, hits: [] }; - if (node.op === "AND") { - let allMatch = true; - let count = 0; - const hits: ScopedSnippet[] = []; - for (const c of active) { - const r = evalNode(c, matchers, ctx); - if (!r.match) { - allMatch = false; - break; - } - count += r.count; - hits.push(...r.hits); - } - const match = node.negate ? !allMatch : allMatch; - return match && !node.negate - ? { match, count, hits } - : { match, count: 0, hits: [] }; - } - // OR - let any = false; - let count = 0; - const hits: ScopedSnippet[] = []; - for (const c of active) { - const r = evalNode(c, matchers, ctx); - if (r.match) { - any = true; - count += r.count; - hits.push(...r.hits); - } - } - const match = node.negate ? !any : any; - return match && !node.negate - ? { match, count, hits } - : { match, count: 0, hits: [] }; -} - -// The share-filter predicate over a transcript record + its availability. -// Mirrors SearchSessionContext.tsx `passesFilter` exactly. -// -// Typed on the three fields it actually reads rather than on TranscriptDetail, -// so the SAME predicate can be applied to a summaries index record while -// planning a scan and to the full transcript record while deciding a hit. -// There is deliberately only one of these: a second, "cheap" predicate for -// planning is exactly how a pruner starts silently disagreeing with the -// scanner about what matches. -type FilterableRecord = { - isLivestream?: boolean; - ageRestricted?: boolean; - uploadDate: string; -}; - -function passesFilters( - rec: FilterableRecord, - f: SearchFilters, - avail: VideoAvailability | undefined, -): boolean { - // ft — type - if (rec.isLivestream ? !f.livestreams : !f.videos) return false; - // fa — audience - if (rec.ageRestricted ? !f.restricted : !f.allAges) return false; - // fav — presence on the source platform - if (!f.states.has(avail?.state ?? "available")) return false; - // fdf / fdt — upload-date range (lexicographic on YYYYMMDD) - if (f.dateFrom && rec.uploadDate < f.dateFrom) return false; - if (f.dateTo && rec.uploadDate > f.dateTo) return false; - return true; -} - -// True when the fav filter could exclude something (so availability must be -// fetched). If every availability bucket is kept, there's nothing to look up. -function needsAvailability(f: SearchFilters | null | undefined): boolean { - return !!f && !VIDEO_STATES.every((s) => f.states.has(s)); -} - // A small lazy live-chat fetcher: per-channel subs manifest + page caches, so a // chat-scope leaf only pulls the subs shards it actually touches. function makeChatFetcher(source: ShardSource) { @@ -1399,11 +1037,13 @@ export async function runSearchSpec( includeSnippets?: boolean; maxPages?: number; snippetsPerVideo?: number; + policy?: SearchPolicy; } = {}, ): Promise<SpecResult> { + const policy = opts.policy ?? MCP_POLICY; const limit = opts.limit ?? 20; const offset = Math.max(0, opts.offset ?? 0); - const maxPages = opts.maxPages ?? MAX_PAGES; + const maxPages = opts.maxPages ?? policy.maxPages; const includeSnippets = opts.includeSnippets !== false; const snippetsPerVideo = opts.snippetsPerVideo ?? 4; const filters = spec.filters ?? null; @@ -1460,6 +1100,7 @@ export async function runSearchSpec( chatCues: chatCuesFor ? await chatCuesFor(ch, rec) : [], snippetsPerVideo, includeSnippets, + snippetChars: policy.snippetChars, }; const r = evalNode(spec.tree, matchers, ctx); if (!r.match) continue; @@ -1477,7 +1118,7 @@ export async function runSearchSpec( matches: r.count || 1, snippets: r.hits, }); - if (all.length >= HARD_VIDEO_CAP) { + if (all.length >= policy.hardVideoCap) { truncated = true; break outer; } @@ -1523,6 +1164,7 @@ export async function runSearchSpec( postText: post.text, snippetsPerVideo, includeSnippets, + snippetChars: policy.snippetChars, }; const r = evalNode(spec.tree, matchers, ctx); if (!r.match) continue; @@ -1540,7 +1182,7 @@ export async function runSearchSpec( matches: r.count || 1, snippets: r.hits, }); - if (all.length >= HARD_VIDEO_CAP) { + if (all.length >= policy.hardVideoCap) { truncated = true; break postsOuter; } diff --git a/plans/one-core-phase-2.md b/plans/one-core-phase-2.md @@ -638,3 +638,268 @@ comment-only** — verified by diffing with comment lines filtered out. `e2e:hub` **5 passed / 0 failed**. `e2e:2origin` was **not run**: it shells `build:hub`, which is red on the base for the reason above, so it cannot reach a spec — exactly as S1 recorded. +### S3 — shipped 2026-09-12 + +Branch `one-core/phase-2-s3`, off `7f86aef` (the `integrate/2026-09-storage-priority` +tip with S1 merged). Six commits, `eaa50b6` → `99fe37a`, plus this note; then three review +fixes, `172099c` → the commit that corrects this note. Unmerged. +**No URL shape moved, no `corpus.json` byte moved, no CONTRACT version moved, no +architecture allow-list entry added — and `"components"` is now in +`FORBIDDEN.lib`, so the list of forbidden edges grew while the list of excused +ones did not.** + +| commit | what | +|---|---| +| `eaa50b6` | `lib/search/{policy,window,evalTree,rank,collapse,leafPipeline}.ts` + a test file each | +| `7252b57` | `mcp/src/search.ts`: 1,563 lines → 1,206, orchestration only; the caps are `MCP_POLICY` | +| `07075d6` | the inversion: four `lib → components` edges deleted, `"components"` added to `FORBIDDEN.lib` | +| `fb6a31d` | the unused `LeafMatcher` import the previous commit left | +| `99fe37a` | the `SearchMode` doc comment travels with the type | +| `726de3e` | this note, as first written | +| `172099c` | **F3** — `lib/search` holds no caller's budget as a default | +| `752ceb1` | **F1** — the result list calls `rankByUploadDateDesc` | +| _(this commit)_ | **F2** + the F1/F3 corrections to this note | + +#### What the slice actually found, where the brief and the code disagreed + +**1. `searchEval.ts` is not a COPY of the mcp evaluator, and merging them +would have been wrong.** The plan says `evalTree.ts` should hold +"`search.ts:1189-1338` plus `lib/searchEval.ts`'s copy, once". Read side +by side they are two algorithms over two data shapes: mcp's `evalNode` +answers *does this record match*, given a record it already holds; +searchEval's `evaluateSlugs`/`orchestrate` answer *which slugs survive*, +streaming, over per-leaf network pipelines, because the browser cannot +hold 1.3 GB of transcripts. Unifying them means materialising the corpus +in the browser. What the slice does instead: the per-record evaluator +moves into `lib/search/evalTree.ts` once, and **each file names the other +as its sibling and states the rule** — the meaning of AND / OR / negate +changes in both together, or it has silently forked. + +**2. `rank.ts` had one implementation to collect, not two.** The stub says +the mcp ordering and the viewer ordering "are the same intent written +twice". mcp does not rank at all: `searchTranscripts` emits in corpus scan +order and `SearchResult` has no sort. The orderings that DO exist are the +viewer's newest-first over `uploadDate` +(`SearchSessionContext.tsx`, inline) and `getThread`'s oldest-first over +`createdAt` tie-broken by id (`search.ts`, inline). Both are now +comparators in `rank.ts` and both call sites call them. **No relevance +ranking was invented** — that would have changed what every caller returns +under cover of a refactor. + +> **Corrected after review (F1).** As first written this note was wrong +> about its own slice: `getThread` was converted but +> `SearchSessionContext.tsx:886` still sorted inline, so `rank.ts` shipped +> with the viewer's ordering in it and zero viewer callers — a second copy +> with better documentation, not a shared comparator. Fixed in `752ceb1`; +> the helper IS the inline ternary, sorts in place and returns the same +> array, and `Array#sort` is stable, so same-date rows keep the order the +> two append loops built. + +**3. `collapse.ts` has one consumer, and the viewer's "duplicate-collapse" +is not one.** `SearchResults.tsx:502` is `duplicateSiblings(...)` from +`components/duplicatesCache.ts` — it decorates a card with a sibling +*count* and a jump menu. It does not remove a row, and the result list is +never collapsed. So there was no second implementation to fold in. What +moved is mcp's `collapseDuplicates`, now **pure and synchronous** and +generic over a `CollapsibleHit`: the caller supplies the already-fetched +duplicate index, so the rules are testable without a transport and the +viewer can adopt them when it wants them, rather than having them imposed +by this slice. + +**4. A sixth module was required: `lib/search/leafPipeline.ts`.** The +inversion is impossible without it. `searchEval.ts` does not merely CALL +`runLeafPipeline` — it is typed on `LayerHit` and `LeafController`, so +injecting the function alone leaves the import. Moving the browser's +streaming leaf scanner down (three worker-pool drivers + the leaf adapter, +~450 lines) is what makes `lib/` self-contained; `components/searchPipeline.ts` +is now 65 lines that answer one question — WHICH FETCH — and supplies +`searchRuntime`. The brief allows new files under `lib/search/`; no new +`components/*.ts` was created, so `common/package.json` is untouched. + +**5. Four back-edges, not one.** The brief named `lib/searchEval.ts:35-38`. +Re-grepping `lib/` found four, and all four had to go before +`"components"` could join `FORBIDDEN.lib`: + +| edge | what it was | how it inverted | +|---|---|---| +| `searchEval → components/searchPipeline` | `runLeafPipeline`, `LayerHit`, `LeafController` | pipeline moved down; the fetch is injected as `SearchRuntime.runLeaf` | +| `searchEval → components/searchLayerCache` | `cacheKey` / `getCached` / `getCachedSync` / `putCached` / `CachedResult` | the memo is injected as `SearchRuntime.cache`; `CachedResult` moved to `searchEval` (an IndexedDB store does not get to define what a leaf result IS) and `searchLayerCache` re-exports it | +| `aiHandoff → components/searchPipeline` | `LayerHit` (type) | import path only | +| `searchQuery → components/urlState` | `SearchMode` (type) | moved to `searchQuery`, which interprets it; `urlState` re-exports | + +`runQueryTree` gains one required `runtime` field. Its three callers pass +`searchRuntime`: `SearchSessionContext.tsx`, `charts/useSearchSeries.ts` +and `export/app/lib/askRetrieval.ts`. No other import site moved: every +name `components/searchPipeline` exported **that anything imported** it +still exports, including the `LayerHit` that `askRetrieval.ts` imports +from it. + +> **Corrected after review (F2).** The first wording said "every name it +> exported", which is not true and worth being exact about: three exports +> were REMOVED — `createSearchPipeline`, `createPostsSearchPipeline` and +> `createSubsSearchPipeline` (`7f86aef:searchPipeline.ts:54,300,454`). +> Grepping every package found no importer of any of them; they were only +> ever reached through `runLeafPipeline`, which dispatches on the leaf's +> scope. So they moved down to `lib/search/leafPipeline.ts` and stayed +> module-internal to it rather than being re-exported from a binding whose +> whole point is to answer one question. Re-exporting three functions +> nobody calls would have been publishing a surface to preserve a +> sentence. + +**6. `window.ts` keeps TWO excerpt shapes on purpose.** The stub hoped +"the viewer's excerpt and the MCP's excerpt are the same excerpt". They +are not the same function: the MCP takes a ±45 s window around every +matched cue and merges them (the paragraph a sweep cites); the viewer +emits one row per matched cue, widened to the neighbouring cues ONLY when +the match straddles a caption break (the line a reader scans). Collapsing +them changes both outputs, and the viewer's is pinned by the export e2e's +rendered hit text. One module owns both, one `truncate`, one `clock`, no +duplicated code — and the file says out loud why they stay two. + +**7. `policy.ts` also names the uncapped one.** The plan said "the viewer +passes its own, uncapped". `VIEWER_POLICY` exists so there is one place +that says what uncapped MEANS, and a test pins the load-bearing +consequence (`truncate(t, Infinity)` is the identity; +`0 >= Infinity` is false). + +> **Hardened after review (F3).** As first written, `truncate`'s `max` and +> `RecordCtx.snippetChars` defaulted to `MCP_POLICY.snippetChars`, so +> `lib/search/` held one caller's budget as the shared default — and the +> failure mode was the bad kind: a viewer adopter that forgot to pass its +> policy would not break, it would silently clip every excerpt at 240 with +> the whole suite green. All three widths (`truncate`'s `max`, +> `windowedTranscript`'s `maxLines`, `RecordCtx.snippetChars`) are now +> REQUIRED, `MCP_POLICY` is imported nowhere under `lib/search/` except by +> the module that defines it, and two tests run the same input under both +> policies and assert the widths differ. Fixed in `172099c`; mcp already +> passed its policy explicitly at every site, so **no read path changed and +> the bench counters cannot move**. + +**8. `buildScanPlan` stayed in `mcp/src/search.ts`.** It is reader-driven +orchestration, `scanPlan.test.ts` sits beside it, and it calls the single +`passesFilters` from `lib/search/evalTree` — so the pruner and the scanner +still cannot disagree about what matches, which was the one property that +mattered. + +#### Gates + +- `pnpm -r exec tsc --noEmit` — clean in all six packages, after every commit. +- `pnpm --filter yt-dlp-transcript-common test` — **1129 passed / 0 failed** + (baseline **1077** at `7f86aef`; + 3 `policy`/`window` width cases, + 8 + `window`, + 2 `rank`, + 6 `collapse`, + 21 `evalTree`, + 12 `leafPipeline` + = 52; none lost). + The architecture test passes with `"components"` in `FORBIDDEN.lib` and the + ALLOWED ledger **byte-identical** to the base. +- `pnpm --filter yt-dlp-transcript-mcp test` — **205 passed / 0 failed**, + unchanged. +- `pnpm test:scripts` — **71 passed / 1 skipped**, unchanged. +- `pnpm --filter export exec next build` — **compiled successfully**, 11 static + pages. The only real test that none of this dragged `reader-fs.ts` or a + react-query module into a client chunk. (Prerender needs a composed + `export/public`: the committed fixture at + `plans/tools/compose-fixture-one-youtube-channel/public/` copied in, which is + gitignored — a bare checkout fails on a missing `summaries/manifest.json` + before it reaches any code this slice touched.) +- **compose-site byte-identity** over the FACTS.md fixture recipe: + `IDENTICAL modulo the build clock` for both trees — 16 files under `public/`, + 7 under `index/`. Nothing composed moved. +- `pnpm --filter yt-dlp-transcript-mcp bench --repeat 1 --force --local <site>`, + before at `7f86aef` and after at the tip, both over the same composed + hasanalyzer site. **Identical fingerprint** (5 channels, 170 transcript pages, + 3,358 summaries videos, stats + duplicates + digests all present) and **every + structural counter identical, byte for byte** — including the per-case scan + notes and both `pruned` flags: + + | case | reads | bytes parsed | note | + |---|---|---|---| + | cold-channels | 0 → 0 | 0 → 0 | | + | rare, whole corpus | 170 → 170 | 1,335,885,512 → 1,335,885,512 | 170 pages | + | common, whole corpus | 170 → 170 | 1,364,679,296 → 1,364,679,296 | 170 pages | + | channel-scoped | 4 → 4 | 28,847,054 → 28,847,054 | 4 pages | + | date-scoped (filter-first) | 32 → 32 | 213,451,503 → 213,451,503 | 28 pages · pruned | + | state-scoped (filter-first) | 9 → 9 | 73,399,988 → 73,399,988 | 9 pages · pruned | + | enumerate, whole corpus | 170 → 170 | 1,364,679,296 → 1,364,679,296 | 170 pages | + | get_transcripts × 20 ids | 4 → 4 | 28,847,054 → 28,847,054 | | + + Wall ms is noise and is not quoted (the box was above the load threshold for + both runs, and the bench marked them UNRELIABLE, as designed). + + **The site S1 benched was gone and had to be rebuilt** — `/home` is at 100 %, + and nothing 1.3 GB survived. Recipe, for the next slice: compose `hasanalyzer` + READ-ONLY against the primary checkout, by pointing `EXPORT_INDEX_DIR` at a + directory holding two symlinks (`shared`, `sites` → the primary checkout's + `export/.export-index`), `SITES_DIR` at a **copy** of `sites/hasanalyzer` with + `"archives": false` added, `TRANSCRIPTS_DIR` at a throwaway dir holding + `settings.json` `{}` plus the corpus-wide `duplicates*.json` / + `search-aliases.json`, and `EXPORT_PUBLIC_DIR` at scratch. Nothing writes + inside `transcripts/`: the compose cache is keyed off the public dir, the + archive builder is off, and the channel signer opens its LMDB under the + throwaway root. The BEFORE run reproduced S1's recorded numbers to the byte, + which is the evidence that the rebuilt corpus is the same corpus. Delete the + 2.3 GB output afterwards. +- **e2e**, behind the queue lock from the worktree (port block + 4000/4001/4010/4011/4020): export `e2e` **172 passed / 0 failed**, + `e2e:hub` **5 passed / 0 failed**, editor `export-search` + + `export-player-platform-cache` **21 passed / 0 failed**. Re-run after the + three review fixes (F1 touches `SearchSessionContext.tsx`): export + **172 passed**, hub **5 passed**, editor `export-search` **19 passed** — + all 0 failed. +- The bench was **not** re-run after the review fixes, and does not need to be: + `git diff fb6a31d..HEAD -- mcp/` is empty, so the benchmarked binary is + byte-identical. F3 removed defaults that mcp never took (it passed + `policy.snippetChars` and `opts.maxLines ?? MCP_POLICY.windowLineCap` + explicitly at every site) and F1 is viewer-only. +- `e2e:2origin` was **not run**: it is red on the base for a reason outside this + slice (the hub `/ask` prerender — see S1's §Record, verified there against + `c7f7b90` itself). Because S3 touches two of the three files that failure + names (`askRetrieval.ts`, `SearchSessionContext.tsx`), `build:hub` was run + anyway to confirm the failure is the SAME one: `INSTANCE_MODE=hub next build` + compiles successfully and type-checks, then dies on the same page with the + same message from the same two chunks — + `Error: useSearchSession must be used within a SearchSessionProvider`, + `Export encountered an error on /(workspace)/ask/page`, + `SearchSessionContext` + `AskChat`. Unchanged, and the compile+typecheck + passing is itself the evidence that the injected `searchRuntime` resolves in + hub mode too. +- **The live jeralyzer contract** (S1 left this as S3's gate): + `curl https://jeralyzer.pages.dev/corpus.json` is **byte-identical** to + `plans/tools/jeralyzer-corpus-2026-09-12.json` — 12,380 bytes, zero + structural differences, `generatedAt` included. Be honest about what that + does and does not prove: the site has not been rebuilt since the snapshot, + so this confirms the snapshot is still current, not that a rebuild would + match. **A rebuilt jeralyzer was not attempted** — 30 channels / 30,886 + videos, and `/home` is at 100 %. The real evidence that the emitted contract + did not move is the compose-site fixture byte-identity above: jeralyzer's + `corpus.json` comes out of the same `buildSiteCorpus` the fixture exercises, + and S3 does not open `corpus.ts`, `compose-site.ts` or `compose-hub.ts` at + all. + +#### One e2e failure, and why it is not this slice + +The editor subset failed 1 of 21 on its first run: +`export-search.spec.ts:612` "a site's own social links win over the global +default", with `ENOENT: … editor/test-settings.json` at `helpers.ts:166`. + +It is a spec-SUBSET setup dependency, not a regression. That file is created +only by `resetData()` (`editor/e2e/helpers.ts:19`), and `export-search.spec.ts` +never calls `resetData` — it reads a settings file some earlier spec in the full +suite wrote. Running only these two spec files in a fresh worktree can never +create it. Seeding it the way the suite does +(`cp e2e/fixtures/test-settings.default.json editor/test-settings.json`) and +re-running the same subset: **21 passed / 0 failed**. Nothing in S3 touches +settings, social links or the footer. + +#### Notes for the slices still in flight + +- `components/searchPipeline.ts` is now 65 lines and exports `searchRuntime`. + **S2a: if a cache's `fetchX` signature changes, the one place to rebind it is + `CACHE_FETCHERS` there** — not inside `lib/`, which no longer knows the caches + exist. +- `lib/search/leafPipeline.ts` uses bare `setTimeout`/`clearTimeout` rather than + `window.setTimeout`. Identical in a browser, and it is what lets the module be + driven from a test process with no DOM — which is how + `leafPipeline.test.ts` counts page reads over an in-memory `ArchiveReader`. +- The two sibling comments (`lib/searchEval.ts` ↔ `lib/search/evalTree.ts`) are + load-bearing prose: they are the only thing stopping the two evaluations of + the same algebra from drifting, because no test can compare a streaming + slug-set walk against a per-record boolean.