// 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 { postConversation, type Post } from "yt-dlp-transcript-common/lib/posts"; import type { Cue } from "yt-dlp-transcript-common/lib/vtt"; import { hitsAcrossTracks } from "yt-dlp-transcript-common/lib/captionTracks"; import type { Platform } from "yt-dlp-transcript-common/lib/platform"; import { type GroupNode, type LayerScope, } from "yt-dlp-transcript-common/lib/searchQuery"; import type { SearchAlias } from "yt-dlp-transcript-common/lib/searchAliases"; import { resolveChannelGroupId, type ChannelGroup, } from "yt-dlp-transcript-common/lib/channelGroups"; 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, } 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; text: string; // Which layer this snippet came from, when it isn't the transcript. The // non-timed layers (metadata / description / tags) carry seconds 0, and a // reader has to be able to tell "the word appears in the description" from // "the word was said at 0:00" — they license completely different citations. scope?: LayerScope; // A transcript hit from one of the record's ALTERNATE English tracks // (lib/captionTracks.ts) — words its primary does not have there. Absent for // a hit in the primary. track?: string; }; // A scope selector for a search/sweep: any mix of channel handles (slug / key / // name) and group handles (id / name). All fields are optional and additive — // the resolved scope is the UNION of every matched channel and every matched // group's members. No tokens at all means "all channels" (the primitive's // default). export type ChannelSelector = { channel?: string; channels?: string[]; group?: string; groups?: string[]; }; // The outcome of resolving a ChannelSelector against a source: the deduped set // of channels to scan, which group tokens matched (for a scope note), and any // tokens that matched nothing (surfaced as a warning so a typo isn't silently a // whole-corpus scan). `all` is true only when no selector was given. export type ResolvedSelection = { channels: ChannelRef[]; matchedGroups: ChannelGroup[]; unknownChannels: string[]; unknownGroups: string[]; all: boolean; }; function selectorTokens(one: string | undefined, many: string[] | undefined): string[] { return [...(one ? [one] : []), ...(many ?? [])] .map((t) => (typeof t === "string" ? t.trim() : "")) .filter((t) => t !== ""); } // Resolve a ChannelSelector to a concrete set of channels. Channel tokens match // slug|key|name (case-insensitive); group tokens match a group by id or name // (case-insensitive; blank group names are unselectable) and expand to the // channels whose resolved groupId is that group — mirroring the browser's // group→slugs mapping (resolveChannelGroupId folds unknown/absent groupIds onto // the default). The result is the deduped union; unmatched tokens are reported. export async function resolveSelectedChannels( source: ShardSource, sel: ChannelSelector, ): Promise { const allChannels = await source.listChannels(); const channelTokens = selectorTokens(sel.channel, sel.channels); const groupTokens = selectorTokens(sel.group, sel.groups); // No selector → all channels (unchanged primitive behavior). if (channelTokens.length === 0 && groupTokens.length === 0) { return { channels: allChannels, matchedGroups: [], unknownChannels: [], unknownGroups: [], all: true, }; } const selected = new Map(); // key -> ref (dedupe) const unknownChannels: string[] = []; const unknownGroups: string[] = []; const matchedGroups: ChannelGroup[] = []; for (const token of channelTokens) { const want = token.toLowerCase(); const matches = allChannels.filter( (c) => c.slug.toLowerCase() === want || c.key.toLowerCase() === want || c.name.toLowerCase() === want, ); if (matches.length === 0) { unknownChannels.push(token); continue; } for (const c of matches) selected.set(c.key, c); } if (groupTokens.length > 0) { const { groups, defaultGroupId } = await source.loadGroups(); for (const token of groupTokens) { const want = token.toLowerCase(); const group = groups.find( (g) => g.id.toLowerCase() === want || (g.name.trim() !== "" && g.name.toLowerCase() === want), ); if (!group) { unknownGroups.push(token); continue; } if (!matchedGroups.some((g) => g.id === group.id)) matchedGroups.push(group); for (const c of allChannels) { if (resolveChannelGroupId(c.groupId, groups, defaultGroupId) === group.id) { selected.set(c.key, c); } } } } return { channels: [...selected.values()], matchedGroups, unknownChannels, unknownGroups, all: false, }; } // Which corpora a search covers. Defaults to BOTH, so existing agent flows // pick up the social-post corpus automatically. export type ContentType = "video" | "post"; export const ALL_CONTENT_TYPES: ReadonlyArray = ["video", "post"]; export function parseContentTypes(raw: unknown): ContentType[] { if (!Array.isArray(raw) || raw.length === 0) return [...ALL_CONTENT_TYPES]; const out = raw.filter( (t): t is ContentType => t === "video" || t === "post", ); return out.length > 0 ? out : [...ALL_CONTENT_TYPES]; } export type SearchHit = { videoId: string; // The video's slug (`/`) — the archilyzer viewer's `?v=` // value, used to build a moment deep link (momentUrl). slug: string; channelSlug: string; channelName: string; siteTitle?: string; // The public origin of the site that owns this video (a member url in hub // mode, the site base for a single remote). Used with `slug` for a viewer // link; absent for a local source (→ platform-link fallback). siteUrl?: string; platform?: Platform; title: string; uploadDate: string; webpageUrl?: string; matches: number; snippets: Snippet[]; // Set on a hit from the social-post corpus. A post has no timeline, so its // snippets carry seconds 0 and its citation takes no `@ mm:ss`. contentType?: ContentType; author?: string; createdAt?: string; // Other copies of THIS SAME RECORDING that also matched and were collapsed // into this row (cross-platform mirrors, per duplicates.json). Present only // when collapsing actually removed something — the mirrors are named rather // than dropped, so the count is honest without the evidence disappearing. mirrors?: { videoId: string; channelName: string; slug: string }[]; }; export type SearchResult = { hits: SearchHit[]; // 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; offset: number; limit: number; hasMore: boolean; // The curated aliases that fired for this query (so the tool can report the // expansion it applied). Empty for a plain or explicit-regex search. firedAliases: SearchAlias[]; scanned: { channels: number; pages: number }; // The posts pass, counted separately. Without this a posts-only search over a // corpus where no channel ships a posts index reports "scanned 0 page(s) // across 0 channel(s)" — indistinguishable from "searched everything, found // nothing". `channels` counts the channels that actually HAVE a posts index, // so 0 with `requested` true means the post corpus is empty here. // `skippedForTagFilter` is the one way `requested: false` is worth saying out // loud: a post carries no curated tags, so a tag filter drops the whole posts // corpus rather than reporting it as searched-and-unmatched. The export UI // makes the same call (SearchSessionContext's globalScopeSlugs). postsScanned: { requested: boolean; channels: number; pages: number; skippedForTagFilter: boolean; }; // What duplicate collapsing did to the count. `available: false` means this // corpus ships no duplicates.json, so no claim about mirrors can be made // 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 // (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 // misleading in a specific way: the sample is biased toward whichever // channels happen to sort first. Naming the channels that were fully scanned // and the ones never reached turns that from alarming into actionable — you // can re-run scoped to the remainder. coverage: { // A filter let us plan an exact page set instead of scanning everything. pruned: boolean; pagesPlanned: number; pagesTotal: number; // Videos in a channel manifest that the summaries index didn't know about; // their pages were read unconditionally. unknownVideos: number; channelsCompleted: string[]; channelsNotReached: string[]; channelStopped?: { channel: string; page: number; pages: number }; }; // How the scope selector resolved (for the tool's scope note): whether it was // whole-corpus, the channels/groups it matched, and any tokens that matched // nothing (a typo'd channel/group is surfaced, not silently a full scan). selection: { all: boolean; channelCount: number; matchedGroups: ChannelGroup[]; unknownChannels: string[]; unknownGroups: string[]; }; }; // ─── Windowed page reading ─── // // Pages were read one `await` at a time, which on a local corpus leaves the // disk idle for the whole ~390 ms JSON.parse of each 8 MB page, and over HTTP // leaves the connection idle for a whole round trip. This reads a WINDOW of // pages concurrently and then folds them in index order. // // Windowed rather than one big fan-out over every page, for three reasons that // are all load-bearing: // 1. hit order is unchanged — results are folded in page order, which a dozen // ordering assertions in search.test.ts pin; // 2. `max_pages` is honored exactly — the window is clamped to the remaining // budget, so max_pages:1 still reads precisely one page; // 3. overshoot past a cap is bounded by the window width rather than by the // size of the corpus. async function readPageWindow( pages: readonly number[], concurrency: number, read: (page: number) => Promise, ): Promise<(T | null)[]> { // A page that fails to read becomes null and is skipped by the caller — // matching the per-page try/catch this replaces. return mapConcurrent(pages, concurrency, (p) => read(p).then( (v) => v, () => null, ), ); } function concurrencyOf(source: ShardSource): number { return source.pageConcurrency ?? DEFAULT_PAGE_CONCURRENCY; } // ─── Filter-first scan planning ─── // // A filtered query's exact page set is computable BEFORE any transcript is // read: the summaries index says which videos pass the filter, and each channel // manifest's `slugToPage` says which page each video lives on. Measured on the // live corpus (170 pages): the deleted/unlisted question touches 8 pages, a // single channel 8, one upload year 28, and "not livestreams" 162 — so this // pays in proportion to how selective the filter is, and degrades to today's // full scan when it isn't. // // The load-bearing invariant: the index only ever PRUNES PAGES; `passesFilters` // on the real record still DECIDES every hit. A page is skipped only when every // video the manifest places on it is known to be excluded. So a stale or // incomplete index can cost time, never correctness — and the two can't drift, // because pruning and deciding call the same predicate. export type ScanPlan = { // Per channel key: the page indices to read, ascending. Channels whose // manifest could not be read are absent (matching the old `continue`). perChannel: Map; pruned: boolean; pagesPlanned: number; pagesTotal: number; // Videos the manifest lists but the index has never heard of. Their pages are // included unconditionally; a non-zero count here is why a "pruned" scan may // still read more than the filter suggests. unknownVideos: number; }; // 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 // no summaries/). export async function buildScanPlan( source: ShardSource, channels: readonly ChannelRef[], filters: SearchFilters | null | undefined, ): Promise { const perChannel = new Map(); let pagesTotal = 0; let pagesPlanned = 0; let unknownVideos = 0; const wantPrune = filterIsSelective(filters) && typeof source.videoIndex === "function"; const index = wantPrune ? await source.videoIndex!() : null; for (const ch of channels) { let manifest; try { manifest = await source.transcriptsManifest(ch); } catch { continue; // unreachable/missing channel — skip, as before } const pageCount = manifest.pageCount; pagesTotal += pageCount; if (!index || !filters) { const pages = Array.from({ length: pageCount }, (_, i) => i); perChannel.set(ch.key, { ch, pages }); pagesPlanned += pages.length; continue; } const keep = new Set(); for (const [videoId, page] of Object.entries(manifest.slugToPage)) { if (keep.has(page)) continue; // one surviving video is enough to read it const rec = index.get(`${ch.slug}/${videoId}`); if (!rec) { // Unknown to the index — summaries may be older or narrower than the // transcripts. Read the page; the record predicate will decide. unknownVideos++; keep.add(page); continue; } if (passesFilters(rec, filters, rec)) keep.add(page); } const pages = [...keep].sort((a, b) => a - b); perChannel.set(ch.key, { ch, pages }); pagesPlanned += pages.length; } return { perChannel, pruned: Boolean(index), pagesPlanned, pagesTotal, unknownVideos, }; } // 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[], enabled: boolean, ): Promise<{ collapsed: number; clusters: number; available: boolean }> { if (!enabled || typeof source.duplicateIndex !== "function") { return { collapsed: 0, clusters: 0, available: false }; } let index; try { index = await source.duplicateIndex(); } catch { return { collapsed: 0, clusters: 0, available: false }; } if (index.size === 0) return { collapsed: 0, clusters: 0, available: true }; const { collapsed, clusters, kept } = collapseHits(all, index); if (collapsed > 0) { all.length = 0; all.push(...kept); } return { collapsed, clusters, available: true }; } // Scan a source's transcript shards for `query` (alias-aware by default), // 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 policy.hardVideoCap. export async function searchTranscripts( source: ShardSource, opts: { query: string; channel?: string; channels?: string[]; group?: string; groups?: string[]; regex?: boolean; limit?: number; offset?: number; maxPages?: number; includeSnippets?: boolean; useAliases?: boolean; snippetsPerVideo?: number; aliases?: SearchAlias[]; // Which corpora to search. Defaults to both video transcripts and posts. contentTypes?: ContentType[]; // The share-link filter set (availability state / upload-date range / // media type / audience). When selective, it also drives filter-first page // planning, so a narrow question reads a fraction of the corpus. filters?: SearchFilters | null; // Video-level NOT: a video matching any of these is dropped even if the // query matched it. Evaluated over the same text the query is (cues, plus // title), which is what makes `"cup"` minus `"world cup"` mean what a // person means by it. exclude?: string[]; // Which layers of a video to match against. Omitted → today's behavior // exactly: spoken captions plus the title. Naming scopes replaces that // default outright, so ["description"] searches descriptions and NOT // captions. scopes?: LayerScope[]; // Count a recording mirrored across platforms ONCE (default true). The // 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 { const policy = opts.policy ?? MCP_POLICY; const limit = opts.limit ?? 20; const offset = Math.max(0, opts.offset ?? 0); const maxPages = opts.maxPages ?? policy.maxPages; const includeSnippets = opts.includeSnippets !== false; const snippetsPerVideo = opts.snippetsPerVideo ?? 4; const aliases = opts.aliases ?? (opts.useAliases !== false && !opts.regex ? await source.loadAliases() : []); const { match, firedAliases } = buildMatcher({ query: opts.query, regex: opts.regex, useAliases: opts.useAliases, aliases, }); // Video-level NOT. Each term is compiled the same way the query is (so a // regex search excludes by regex too) but WITHOUT alias expansion: an // exclusion is a thing the user named precisely, and quietly widening it via // a curated alias would drop videos they never asked to drop. const excludeMatchers = (opts.exclude ?? []) .map((t) => (typeof t === "string" ? t.trim() : "")) .filter((t) => t !== "") .map( (t) => buildMatcher({ query: t, regex: opts.regex, useAliases: false }).match, ); const isExcluded = (title: string, texts: readonly { text: string }[]): boolean => excludeMatchers.length > 0 && excludeMatchers.some((m) => m(title) || texts.some((c) => m(c.text))); // Which layers to match. The default is exactly what this scanner has always // done — spoken captions plus the title — so an existing caller's results are // byte-identical. Naming scopes replaces that default rather than adding to // it, which is the only reading under which ["description"] is honest. const scopes = opts.scopes; const wantCues = !scopes || scopes.includes("transcripts"); const wantMetadata = Boolean(scopes?.includes("metadata")); const wantTitle = !scopes || wantMetadata; const wantDescription = Boolean(scopes?.includes("description")); const wantTags = Boolean(scopes?.includes("tags")); const wantChat = Boolean(scopes?.includes("chat")); // Live chat is 1.2 GB of subs/ shards, so it is fetched lazily and only for // the videos a chat-scoped query actually reaches. const chatCuesFor = wantChat ? makeChatFetcher(source) : null; const selection = await resolveSelectedChannels(source, { channel: opts.channel, channels: opts.channels, group: opts.group, groups: opts.groups, }); const channels = selection.channels; const all: SearchHit[] = []; let pagesScanned = 0; let channelsScanned = 0; let truncated = false; const contentTypes = opts.contentTypes ?? [...ALL_CONTENT_TYPES]; // `scopes` and `content_types` are both narrowing, so they intersect: asking // for scopes:["description"] must not still scan the post corpus, and // content_types:["video"] must not be widened by a posts scope. const wantVideos = contentTypes.includes("video") && (!scopes || scopes.some((s) => s !== "posts")); const postsAsked = contentTypes.includes("post") && (!scopes || scopes.includes("posts")); // A tag filter takes the posts corpus out of the search entirely: curated // tags live on VIDEO records (lib/curatedTags.ts), so every post would fail // the filter, and scanning them to drop them all is a manifest probe and a // shard read per posting channel spent on a foregone conclusion. Counting // them as hits would be worse: the export viewer hit exactly that bug (see // tag-chips.spec.ts) — rows dropped by the display fold that the header had // already counted. const postsSkippedForTagFilter = postsAsked && (opts.filters?.curatedTags?.length ?? 0) > 0; const wantPosts = postsAsked && !postsSkippedForTagFilter; // Counted apart from the video pass so "no posts index anywhere in scope" is // distinguishable from "searched the posts and found nothing". let postChannelsScanned = 0; let postPagesScanned = 0; const filters = opts.filters ?? null; const availability = wantVideos && needsAvailability(filters) ? await source.availabilityMap() : null; const plan = wantVideos ? await buildScanPlan(source, channels, filters) : null; const concurrency = concurrencyOf(source); const channelsCompleted: string[] = []; let channelStopped: SearchResult["coverage"]["channelStopped"]; outer: if (plan) { for (const entry of plan.perChannel.values()) { const ch = entry.ch; channelsScanned++; for (let i = 0; i < entry.pages.length; ) { const budget = maxPages - pagesScanned; if (budget <= 0) { truncated = true; channelStopped = { channel: ch.name, page: i, pages: entry.pages.length }; break outer; } const window = entry.pages.slice(i, i + Math.min(concurrency, budget)); i += window.length; const loaded = await readPageWindow(window, concurrency, (p) => source.transcriptPage(ch, p), ); for (let w = 0; w < loaded.length; w++) { const records = loaded[w]; if (records === null) continue; // failed page — skipped, as before pagesScanned++; for (const rec of records) { if ( filters && !passesFilters(rec, filters, availability?.get(rec.slug)) ) { continue; } const snippets: Snippet[] = []; const push = (s: Snippet): void => { if (includeSnippets && snippets.length < snippetsPerVideo) { snippets.push(s); } }; let matches = 0; let otherHit = false; if (wantCues) { // Every English track of the record (lib/captionTracks.ts): a // match only an alternate holds names that alternate. for (const cue of hitsAcrossTracks(rec, (cues) => cues.filter((c) => match(c.text)), )) { matches++; push({ clock: clock(cue.start), seconds: cue.start, text: truncate(cue.text, policy.snippetChars), ...(cue.track ? { track: cue.track } : {}), }); } } // Title is matched under the default (no `scopes`) and under an // explicit `metadata` scope; `metadata` additionally matches the // channel name, mirroring the viewer's metadata layer. const titleHit = wantTitle && match(rec.title ?? ""); if (titleHit) otherHit = true; if (wantMetadata && !titleHit && match(rec.channel ?? ch.name)) { otherHit = true; push({ clock: clock(0), seconds: 0, text: `Channel: ${rec.channel ?? ch.name}`, scope: "metadata", }); } if (wantDescription && rec.description && match(rec.description)) { otherHit = true; push({ clock: clock(0), seconds: 0, text: truncate(rec.description, policy.snippetChars), scope: "description", }); } if (wantTags) { const tags = (rec.tags ?? []).join(", "); if (tags && match(tags)) { otherHit = true; push({ clock: clock(0), seconds: 0, text: truncate(tags, policy.snippetChars), scope: "tags", }); } } if (wantChat && chatCuesFor) { for (const cue of await chatCuesFor(ch, rec)) { if (!match(cue.text)) continue; matches++; push({ clock: clock(cue.start), seconds: cue.start, text: truncate(cue.text, policy.snippetChars), scope: "chat", }); } } if (matches === 0 && !otherHit) continue; if (isExcluded(rec.title ?? "", rec.cues ?? [])) continue; all.push({ videoId: rec.id, slug: rec.slug, channelSlug: ch.slug, channelName: ch.name, ...(ch.siteTitle ? { siteTitle: ch.siteTitle } : {}), ...(ch.siteUrl ? { siteUrl: ch.siteUrl } : {}), ...(rec.platform ? { platform: rec.platform } : {}), title: rec.title, uploadDate: rec.uploadDate, webpageUrl: rec.webpageUrl, matches: matches || 1, snippets, }); if (all.length >= policy.hardVideoCap) { truncated = true; channelStopped = { channel: ch.name, page: i - loaded.length + w, pages: entry.pages.length, }; break outer; } } } } channelsCompleted.push(ch.name); } } // ── the social-post corpus ── // A parallel pass over the same selected channels: only social channels ship // a posts manifest, so a video-only site costs one 404 per channel and adds // nothing to the result set. if (wantPosts) { postsOuter: for (const ch of channels) { let pm; try { pm = await source.postsManifest(ch); } catch { continue; } if (!pm) continue; // not a social channel postChannelsScanned++; for (let page = 0; page < pm.pageCount; page++) { if (pagesScanned >= maxPages) { truncated = true; break postsOuter; } let posts: Post[]; try { posts = await source.postsPage(ch, page); } catch { continue; } pagesScanned++; postPagesScanned++; for (const post of posts) { if (!match(post.text)) continue; if (isExcluded("", [{ text: post.text }])) continue; all.push({ videoId: post.id, slug: post.slug, channelSlug: ch.slug, channelName: ch.name, ...(ch.siteTitle ? { siteTitle: ch.siteTitle } : {}), ...(ch.siteUrl ? { siteUrl: ch.siteUrl } : {}), platform: post.platform, // A post has no title; its body stands in so a worklist row is // still readable without pulling snippets. title: truncate(post.text, 120), uploadDate: post.uploadDate, webpageUrl: post.url, matches: 1, contentType: "post", author: post.authorName || post.author, createdAt: post.createdAt, snippets: includeSnippets ? [{ clock: "", seconds: 0, text: truncate(post.text, 480) }] : [], }); if (all.length >= policy.hardVideoCap) { truncated = true; break postsOuter; } } } } } // ── collapse cross-platform duplicates ── // Before slicing, so `total` is the count of distinct RECORDINGS rather than // of uploads. A mirrored video counted twice is not a rounding error inside a // narrow result set, and "N videos said X" is the sentence a sweep actually // writes. const duplicates = await collapseDuplicates( source, all, opts.collapseDuplicates !== false, ); const total = all.length; const hits = all.slice(offset, offset + limit); const planned = [...(plan?.perChannel.values() ?? [])].map((e) => e.ch.name); const reached = new Set([ ...channelsCompleted, ...(channelStopped ? [channelStopped.channel] : []), ]); return { hits, total, offset, limit, hasMore: offset + limit < total, firedAliases, scanned: { channels: channelsScanned, pages: pagesScanned }, postsScanned: { requested: wantPosts, channels: postChannelsScanned, pages: postPagesScanned, skippedForTagFilter: postsSkippedForTagFilter, }, duplicates, truncated, coverage: { pruned: plan?.pruned ?? false, pagesPlanned: plan?.pagesPlanned ?? 0, pagesTotal: plan?.pagesTotal ?? 0, unknownVideos: plan?.unknownVideos ?? 0, channelsCompleted, channelsNotReached: planned.filter((n) => !reached.has(n)), ...(channelStopped ? { channelStopped } : {}), }, selection: { all: selection.all, channelCount: selection.channels.length, matchedGroups: selection.matchedGroups, unknownChannels: selection.unknownChannels, unknownGroups: selection.unknownGroups, }, }; } // Render a single 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), and // return timestamped excerpt lines plus the total match count. Bounded and // high-signal — the batch read the sweep prompt drives. // Locate one archived post by id, mirroring findVideo. Only social channels // ship a posts manifest, so a video-only corpus resolves this cheaply. export async function findPost( source: ShardSource, postId: string, channelHint?: string | string[], ): Promise<{ ch: ChannelRef; post: Post } | null> { let channels = await source.listChannels(); const hints = ( Array.isArray(channelHint) ? channelHint : channelHint ? [channelHint] : [] ) .map((h) => (typeof h === "string" ? h.trim().toLowerCase() : "")) .filter((h) => h !== ""); if (hints.length > 0) { const want = new Set(hints); const filtered = channels.filter( (c) => want.has(c.slug.toLowerCase()) || want.has(c.key.toLowerCase()) || want.has(c.name.toLowerCase()), ); if (filtered.length > 0) channels = filtered; } for (const ch of channels) { let manifest; try { manifest = await source.postsManifest(ch); } catch { continue; } if (!manifest) continue; const page = manifest.slugToPage[postId]; if (page === undefined) continue; let posts: Post[]; try { posts = await source.postsPage(ch, page); } catch { continue; } const post = posts.find((p) => p.id === postId); if (post) return { ch, post }; } return null; } // Every archived post in the same thread as `post`, oldest first. A thread can // straddle shard pages, so this walks the channel's whole (byte-capped) tree. // A FORUM post's "thread" is its conversation — the posts it quotes and the // posts quoting it (postConversation) — not the whole forum thread, which is // the channel itself and may run to thousands of posts. export async function getThread( source: ShardSource, ch: ChannelRef, post: Post, ): Promise { const threadId = post.threadId || post.id; const manifest = await source.postsManifest(ch); if (!manifest) return [post]; const forum = post.platform === "xenforo"; const thread: Post[] = []; const all: Post[] = []; for (let page = 0; page < manifest.pageCount; page++) { let posts: Post[]; try { posts = await source.postsPage(ch, page); } catch { continue; } for (const p of posts) { if (forum) all.push(p); else if ((p.threadId || p.id) === threadId) thread.push(p); } } if (forum) return postConversation(post, all); rankThread(thread); return thread.length > 0 ? thread : [post]; } export function getWindowedTranscript( record: TranscriptDetail, 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 (default // MCP_POLICY.windowLineCap). The earliest lines are kept; `maxCues` above // stays the per-window bound. maxLines?: number; } = {}, ): { lines: string[]; matchCount: number } { return windowedTranscript(record.cues ?? [], matcher, { ...opts, maxLines: opts.maxLines ?? MCP_POLICY.windowLineCap, }); } // Locate a single video across the source's channels via each channel's // slugToPage map, returning the full record + its channel. `channelHint` // (slug/key/name, or a list of them) short-circuits the scan when the caller // knows the channel(s) — it's a lookup accelerator only, so an unmatched hint // silently falls back to a full scan. export async function findVideo( source: ShardSource, videoId: string, channelHint?: string | string[], ): Promise<{ ch: ChannelRef; record: TranscriptDetail } | null> { let channels = await source.listChannels(); const hints = (Array.isArray(channelHint) ? channelHint : channelHint ? [channelHint] : []) .map((h) => (typeof h === "string" ? h.trim().toLowerCase() : "")) .filter((h) => h !== ""); if (hints.length > 0) { const want = new Set(hints); const filtered = channels.filter( (c) => want.has(c.slug.toLowerCase()) || want.has(c.key.toLowerCase()) || want.has(c.name.toLowerCase()), ); // Prefer the hinted channel(s), but fall back to a full scan if not found. if (filtered.length > 0) channels = filtered; } for (const ch of channels) { let manifest; try { manifest = await source.transcriptsManifest(ch); } catch { continue; } const page = manifest.slugToPage[videoId]; if (page === undefined) continue; let records: TranscriptDetail[]; try { records = await source.transcriptPage(ch, page); } catch { continue; } const record = records.find((r) => r.id === videoId); if (record) return { ch, record }; } return null; } // ─── Full-fidelity spec engine: query tree + filter predicate ─── // // `runSearchSpec` is the per-record evaluator that backs the `open_link` / // `sweep link=` flows: it honors everything an archilyzer share link can carry // (a composite `qt=` query tree over any mix of scopes, plus the `fc/ft/fa/fav/ // fdf/fdt` filters), while keeping the same paging / total / truncation // contract as `searchTranscripts`. The plain `search_transcripts` path stays on // the simpler `searchTranscripts` scanner above (a trivial single-leaf query), // so today's callers/tests are unaffected. // // 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 // for a legacy `q`/`re` link). tree: GroupNode; // The decoded share filters, or null for an unfiltered search. filters?: SearchFilters | null; // Curated aliases to expand transcripts-scope leaves through (as in the plain // path). Empty/omitted → no expansion. aliases?: SearchAlias[]; // Expand transcripts leaves via aliases (default true; skipped per-leaf for a // regex leaf). useAliases?: boolean; }; export type SpecHit = { videoId: string; slug: string; channelSlug: string; channelName: string; siteTitle?: string; siteUrl?: string; platform?: Platform; title: string; uploadDate: string; webpageUrl?: string; matches: number; snippets: ScopedSnippet[]; }; export type SpecResult = { hits: SpecHit[]; total: number; offset: number; limit: number; hasMore: boolean; firedAliases: SearchAlias[]; scanned: { channels: number; pages: number }; truncated: boolean; }; // 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) { const manifests = new Map>>>(); const pages = new Map>>>(); return async function chatCues(ch: ChannelRef, rec: TranscriptDetail): Promise { let mp = manifests.get(ch.key); if (!mp) { mp = source.subsManifest(ch).catch(() => null); manifests.set(ch.key, mp); } const manifest = await mp; if (!manifest) return []; const page = manifest.slugToPage[rec.id]; if (page === undefined) return []; const pk = `${ch.key}:${page}`; let pp = pages.get(pk); if (!pp) { pp = source.subsPage(ch, page).catch(() => []); pages.set(pk, pp); } const records = await pp; const found = records.find((r) => r.slug === rec.slug || r.id === rec.id); const lc = found?.tracks?.live_chat; return Array.isArray(lc) ? lc : []; }; } // Run a full-fidelity spec over the given (already-resolved) channel set. Same // paging/total/truncation semantics as searchTranscripts. export async function runSearchSpec( source: ShardSource, channels: ChannelRef[], spec: SearchSpec, opts: { offset?: number; limit?: number; includeSnippets?: boolean; maxPages?: number; snippetsPerVideo?: number; policy?: SearchPolicy; } = {}, ): Promise { const policy = opts.policy ?? MCP_POLICY; const limit = opts.limit ?? 20; const offset = Math.max(0, opts.offset ?? 0); const maxPages = opts.maxPages ?? policy.maxPages; const includeSnippets = opts.includeSnippets !== false; const snippetsPerVideo = opts.snippetsPerVideo ?? 4; const filters = spec.filters ?? null; const { matchers, fired } = buildLeafMatchers( spec.tree, spec.aliases ?? [], spec.useAliases !== false, ); const wantsChat = [...matchers.values()].some((m) => m.scope === "chat"); const chatCuesFor = wantsChat ? makeChatFetcher(source) : null; const availability = needsAvailability(filters) ? await source.availabilityMap() : null; const all: SpecHit[] = []; let pagesScanned = 0; let channelsScanned = 0; let truncated = false; // Filter-first: a link carrying an availability or date filter reads only the // pages that can hold a match, which is the whole point of `open_link` being // the one caller with real filters today. const plan = await buildScanPlan(source, channels, filters); const concurrency = concurrencyOf(source); outer: for (const entry of plan.perChannel.values()) { const ch = entry.ch; channelsScanned++; for (let i = 0; i < entry.pages.length; ) { const budget = maxPages - pagesScanned; if (budget <= 0) { truncated = true; break outer; } const window = entry.pages.slice(i, i + Math.min(concurrency, budget)); i += window.length; const loaded = await readPageWindow(window, concurrency, (p) => source.transcriptPage(ch, p), ); for (const records of loaded) { if (records === null) continue; pagesScanned++; for (const rec of records) { if (filters && !passesFilters(rec, filters, availability?.get(rec.slug))) { continue; } const ctx: RecordCtx = { title: rec.title ?? "", channel: rec.channel ?? ch.name, description: rec.description ?? "", tags: (rec.tags ?? []).join(", "), cues: rec.cues ?? [], ...(rec.altTracks ? { altTracks: rec.altTracks } : {}), chatCues: chatCuesFor ? await chatCuesFor(ch, rec) : [], snippetsPerVideo, includeSnippets, snippetChars: policy.snippetChars, }; const r = evalNode(spec.tree, matchers, ctx); if (!r.match) continue; all.push({ videoId: rec.id, slug: rec.slug, channelSlug: ch.slug, channelName: ch.name, ...(ch.siteTitle ? { siteTitle: ch.siteTitle } : {}), ...(ch.siteUrl ? { siteUrl: ch.siteUrl } : {}), ...(rec.platform ? { platform: rec.platform } : {}), title: rec.title, uploadDate: rec.uploadDate, webpageUrl: rec.webpageUrl, matches: r.count || 1, snippets: r.hits, }); if (all.length >= policy.hardVideoCap) { truncated = true; break outer; } } } } } // ── posts pass ── // Only run when the tree actually has a posts leaf: a video-only spec must // not pay a manifest probe per channel. Post records reuse the same evalNode // with `postText` set and no cues, so AND/OR/negate semantics are identical. // …and the same exclusion as the plain path: a posts leaf under a tag filter // can only ever match nothing, because a post carries no curated tags. const wantsPosts = [...matchers.values()].some((m) => m.scope === "posts") && (filters?.curatedTags?.length ?? 0) === 0; if (wantsPosts) { postsOuter: for (const ch of channels) { let pm; try { pm = await source.postsManifest(ch); } catch { continue; } if (!pm) continue; for (let page = 0; page < pm.pageCount; page++) { if (pagesScanned >= maxPages) { truncated = true; break postsOuter; } let posts: Post[]; try { posts = await source.postsPage(ch, page); } catch { continue; } pagesScanned++; for (const post of posts) { const ctx: RecordCtx = { title: post.text, channel: post.authorName || post.author, description: "", tags: "", cues: [], chatCues: [], postText: post.text, snippetsPerVideo, includeSnippets, snippetChars: policy.snippetChars, }; const r = evalNode(spec.tree, matchers, ctx); if (!r.match) continue; all.push({ videoId: post.id, slug: post.slug, channelSlug: ch.slug, channelName: ch.name, ...(ch.siteTitle ? { siteTitle: ch.siteTitle } : {}), ...(ch.siteUrl ? { siteUrl: ch.siteUrl } : {}), platform: post.platform, title: truncate(post.text, 120), uploadDate: post.uploadDate, webpageUrl: post.url, matches: r.count || 1, snippets: r.hits, }); if (all.length >= policy.hardVideoCap) { truncated = true; break postsOuter; } } } } } const total = all.length; return { hits: all.slice(offset, offset + limit), total, offset, limit, hasMore: offset + limit < total, firedAliases: fired, scanned: { channels: channelsScanned, pages: pagesScanned }, truncated, }; }