commit 2485b59fc1bfdb5c0f4176497017b20b9e2d8e6c
parent ea2d3d04d126af72cf52b77dfe0d48a78e791d88
Author: I Mean I'm Just Saying <imeanimjustsaying@kiwifarms.st>
Date: Sat, 12 Sep 2026 02:46:00 -0400
common: one search pipeline — the pure half of lib/search/
S1 left five `export {}` stubs under `common/lib/search/`. This fills the
four that need no transport, plus one the stubs did not name, and adds a
test file per module. Nothing imports them yet — the two callers move in
the following commits — so this commit changes no behaviour at all.
policy.ts the four caps that were private constants of
mcp/src/search.ts, as data. MCP_POLICY holds them
verbatim, which is WHY the bench's structural counters
cannot move: the scanner is handed its limits instead
of re-deriving them. VIEWER_POLICY is the uncapped
counterpart — a human scrolling a page is not spending
an agent's token budget.
window.ts every excerpt shape, over lib/transcriptWindow's
primitives: the MCP's merged ±seconds window, the
viewer's one-line-per-cue row, one truncate, one clock.
The module says out loud that the two excerpt shapes
stay two: a sweep needs the paragraph around a claim, a
result card needs the line that matched.
evalTree.ts the per-record tree evaluator and the ONE filter
predicate, with the matcher compilation that feeds it.
rank.ts the two orderings this corpus has, as comparators.
collapse.ts mirror collapsing, now pure and synchronous — the
caller supplies the duplicate index, so the module
never reaches for a transport.
leafPipeline.ts the browser's streaming leaf scanner, with its three
fetches taken as parameters. Not in the stub list, and
it is the file that makes the lib -> components
inversion possible at all.
lib/search/leafPipeline.test.ts drives the real pipeline over fetchers
derived from an in-memory ArchiveReader, so "one page read serves every
slug on it" is counted rather than assumed.
Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Diffstat:
12 files changed, 2412 insertions(+), 101 deletions(-)
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,321 @@
+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 {
+ 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,
+ ...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);
+});
+
+test("snippetChars comes from the policy the caller passes", () => {
+ const long = "y".repeat(300);
+ const wide = evalLeaf(
+ newLeaf({ id: "l", query: "y" }),
+ { scope: "transcripts", test: () => true },
+ ctx({ cues: [cue(0, long)], snippetChars: Number.POSITIVE_INFINITY }),
+ );
+ assert.equal(wide.hits[0].text.length, 300);
+ const narrow = evalLeaf(
+ newLeaf({ id: "l", query: "y" }),
+ { scope: "transcripts", test: () => true },
+ ctx({ cues: [cue(0, long)] }),
+ );
+ assert.equal(narrow.hits[0].text.length, 240, "…defaulting to the MCP's 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";
+import { MCP_POLICY } from "./policy";
+
+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`. Defaults to the
+ // MCP's 240 so an existing caller's output is byte-identical.
+ 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 ?? MCP_POLICY.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,115 @@
+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";
+
+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 "), "a b c");
+ assert.equal(truncate("abcdef", 4), "abc…");
+ assert.equal(truncate("abcd", 4), "abcd");
+});
+
+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,
+ });
+ 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,
+ });
+ assert.deepEqual(bare.lines, ["hit"]);
+
+ const linked = windowedTranscript([cue(7, "hit")], () => true, {
+ stamp: (c, s) => `${c}|${s}`,
+ });
+ 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,200 @@
-// 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";
+import { MCP_POLICY } from "./policy";
+
+// 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` comes from a SearchPolicy (`snippetChars`); the default is the
+// MCP's, which is what every existing caller passed implicitly.
+export function truncate(text: string, max = MCP_POLICY.snippetChars): 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. `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 ?? MCP_POLICY.windowLineCap);
+ }
+ 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;
+}